using System.Text; using System.Text.Json; using Confluent.Kafka; using Microsoft.Extensions.Options; using RabbitMQ.Client; public sealed class NotificationPublisherService : BackgroundService { private readonly IOptions _rabbitOpts; private readonly KafkaOptions _kafkaOptions; private readonly ILogger _logger; private readonly ClinicalSyncOptions _syncOptions; public NotificationPublisherService( IOptions rabbitOpts, IOptions kafkaOptions, ILogger logger, IOptions syncOptions) { _rabbitOpts = rabbitOpts; _kafkaOptions = kafkaOptions.Value; _logger = logger; _syncOptions = syncOptions.Value; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // Short delay so topology provisioner finishes before first publish. await Task.Delay(TimeSpan.FromSeconds(3), stoppingToken); var factory = BuildRabbitFactory(); using var conn = factory.CreateConnection("notification-publisher"); using var chan = conn.CreateModel(); var props = chan.CreateBasicProperties(); props.Persistent = true; // messages survive broker restart var consumerConfig = new ConsumerConfig { BootstrapServers = _kafkaOptions.BootstrapServers, GroupId = _kafkaOptions.NotificationPublisherGroupId, AutoOffsetReset = Enum.Parse( _kafkaOptions.NotificationPublisherAutoOffsetReset, ignoreCase: true), EnableAutoCommit = false, }.ApplySecurity(_kafkaOptions); using var consumer = new ConsumerBuilder(consumerConfig).Build(); consumer.Subscribe(new[] { _kafkaOptions.Topics.AlertGenerated, _kafkaOptions.Topics.EncounterStatusChanged, }); _logger.LogInformation("NotificationPublisherService started — consuming {Topics}", string.Join(", ", _kafkaOptions.Topics.AlertGenerated, _kafkaOptions.Topics.EncounterStatusChanged)); var guard = new PoisonPillGuard("notification-publisher", _kafkaOptions.MaxPoisonRetries, _logger); try { while (!stoppingToken.IsCancellationRequested) { ConsumeResult? result; try { result = consumer.Consume(TimeSpan.FromMilliseconds(500)); } catch (ConsumeException ex) { _logger.LogError(ex, "NotificationPublisher consume error"); continue; } if (result is null) continue; try { if (result.Topic == _kafkaOptions.Topics.AlertGenerated) await HandleAlertGeneratedAsync(chan, props, result.Message.Value, stoppingToken); else if (result.Topic == _kafkaOptions.Topics.EncounterStatusChanged) await HandleEncounterStatusChangedAsync(chan, props, result.Message.Value, stoppingToken); consumer.Commit(result); guard.OnSuccess(); } catch (Exception ex) { if (guard.ShouldSkip(result, ex)) { consumer.Commit(result); continue; } _logger.LogError(ex, "NotificationPublisher failed to process message from {Topic} — will retry", result.Topic); } } } finally { consumer.Close(); } } private Task HandleAlertGeneratedAsync( IModel chan, IBasicProperties props, string payload, CancellationToken ct) { var doc = JsonDocument.Parse(payload); var severity = doc.RootElement.GetProperty("severity").GetString(); if (!string.Equals(severity, "Critical", StringComparison.OrdinalIgnoreCase)) { // Warning-level alerts are not paged — they appear in the clinician dashboard only. return Task.CompletedTask; } var alertId = doc.RootElement.GetProperty("alertId").GetString(); // Skip paging for gateway-synced alerts when configured if (_syncOptions.SuppressPagingForSyncedAlerts && doc.RootElement.TryGetProperty("syncedFromGateway", out var synced) && synced.GetBoolean()) { _logger.LogInformation( "Skipping central paging for gateway-synced alert {AlertId} — ward already paged locally", alertId); return Task.CompletedTask; } var body = Encoding.UTF8.GetBytes(payload); chan.BasicPublish( exchange: RabbitMqTopologyProvisioner.Exchange, routingKey: RabbitMqTopologyProvisioner.PagingKey, basicProperties: props, body: body); _logger.LogInformation("Published paging job to alerts.paging.queue for alert {AlertId}", alertId); return Task.CompletedTask; } private Task HandleEncounterStatusChangedAsync( IModel chan, IBasicProperties props, string payload, CancellationToken ct) { var doc = JsonDocument.Parse(payload); var newStatus = doc.RootElement.GetProperty("newStatus").GetString(); if (!string.Equals(newStatus, "Discharged", StringComparison.OrdinalIgnoreCase)) return Task.CompletedTask; var body = Encoding.UTF8.GetBytes(payload); chan.BasicPublish( exchange: RabbitMqTopologyProvisioner.Exchange, routingKey: RabbitMqTopologyProvisioner.DischargeKey, basicProperties: props, body: body); var encounterId = doc.RootElement.GetProperty("encounterId").GetString(); _logger.LogInformation("Published discharge summary job for encounter {EncounterId}", encounterId); return Task.CompletedTask; } private IConnectionFactory BuildRabbitFactory() { var o = _rabbitOpts.Value; return RabbitMqConnectionFactory.Create(o, dispatchConsumersAsync: true); } }