Files
vigilcare-clinical/VigilCareClinicalAPI/BackgroundServices/Notifications/NotificationPublisherService.cs
T
voltsrage 2a3ef62a7d
CI / frontend (push) Failing after 57s
CI / backend (push) Failing after 6m27s
Add deployment
2026-08-05 00:26:20 +08:00

169 lines
6.3 KiB
C#

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<RabbitMqOptions> _rabbitOpts;
private readonly KafkaOptions _kafkaOptions;
private readonly ILogger<NotificationPublisherService> _logger;
private readonly ClinicalSyncOptions _syncOptions;
public NotificationPublisherService(
IOptions<RabbitMqOptions> rabbitOpts,
IOptions<KafkaOptions> kafkaOptions,
ILogger<NotificationPublisherService> logger,
IOptions<ClinicalSyncOptions> 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<AutoOffsetReset>(
_kafkaOptions.NotificationPublisherAutoOffsetReset, ignoreCase: true),
EnableAutoCommit = false,
}.ApplySecurity(_kafkaOptions);
using var consumer = new ConsumerBuilder<string, string>(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<string, string>? 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);
}
}