using System.Text; using System.Text.Json; using Microsoft.Extensions.Options; using RabbitMQ.Client; using RabbitMQ.Client.Events; public sealed class LocalEscalationWorkerService : BackgroundService { private readonly IOptions _opts; private readonly IServiceScopeFactory _scopes; private readonly ILogger _logger; public LocalEscalationWorkerService( IOptions opts, IServiceScopeFactory scopes, ILogger logger) { _opts = opts; _scopes = scopes; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); while (!stoppingToken.IsCancellationRequested) { try { await RunConsumerAsync(stoppingToken); return; } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (Exception ex) { _logger.LogError(ex, "LocalEscalationWorkerService disconnected — retrying in 5s"); await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken); } } } private async Task RunConsumerAsync(CancellationToken stoppingToken) { var o = _opts.Value; var factory = RabbitMqConnectionFactory.Create(o, dispatchConsumersAsync: true); using var connection = factory.CreateConnection("gateway-escalation-worker"); using var channel = connection.CreateModel(); channel.BasicQos(0, prefetchCount: 5, global: false); var consumer = new AsyncEventingBasicConsumer(channel); consumer.Received += async (_, ea) => { await HandleEscalationAsync(channel, ea, stoppingToken); }; channel.BasicConsume("alerts.escalation.queue", autoAck: false, consumer); _logger.LogInformation("LocalEscalationWorkerService consuming alerts.escalation.queue"); await Task.Delay(Timeout.Infinite, stoppingToken); } private async Task HandleEscalationAsync(IModel channel, BasicDeliverEventArgs ea, CancellationToken ct) { var payload = Encoding.UTF8.GetString(ea.Body.Span); var doc = JsonDocument.Parse(payload); var alertId = Guid.Parse(doc.RootElement.GetProperty("alertId").GetString()!); _logger.LogCritical("[ESCALATION] Paging on-call backup for alert {AlertId}", alertId); try { var escalated = await UpdateAlertStatusEscalatedAsync(alertId, ct); channel.BasicAck(ea.DeliveryTag, multiple: false); if (escalated) { _logger.LogWarning( "[ESCALATION-DONE] Alert {AlertId} status → Escalated in gateway DB", alertId); } } catch (Exception ex) { _logger.LogError(ex, "Local escalation worker failed for alert {AlertId}", alertId); channel.BasicNack(ea.DeliveryTag, multiple: false, requeue: true); } } private async Task UpdateAlertStatusEscalatedAsync(Guid alertId, CancellationToken ct) { await using var scope = _scopes.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); var alert = await db.ClinicalAlerts.FindAsync([alertId], ct); if (alert is null) return false; if (alert.Status != AlertStatus.Open) return false; alert.Status = AlertStatus.Escalated; await db.SaveChangesAsync(ct); return true; } }