234 lines
8.1 KiB
C#
234 lines
8.1 KiB
C#
using System.Text.Json;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Npgsql;
|
|
using Prometheus;
|
|
using VigilCare.ClinicalContracts.Sync;
|
|
|
|
public class ClinicalSyncBatchProcessor
|
|
{
|
|
private readonly AppDbContext _db;
|
|
private readonly IObservationService _observations;
|
|
private readonly IAlertService _alerts;
|
|
private readonly ClinicalMetrics _metrics;
|
|
private readonly ILogger<ClinicalSyncBatchProcessor> _logger;
|
|
|
|
public ClinicalSyncBatchProcessor(
|
|
AppDbContext db,
|
|
IObservationService observations,
|
|
IAlertService alerts,
|
|
ClinicalMetrics metrics,
|
|
ILogger<ClinicalSyncBatchProcessor> logger)
|
|
{
|
|
_db = db;
|
|
_observations = observations;
|
|
_alerts = alerts;
|
|
_metrics = metrics;
|
|
_logger = logger;
|
|
}
|
|
|
|
public async Task ProcessBatchAsync(Guid batchId, CancellationToken ct)
|
|
{
|
|
using var timer = _metrics.ClinicalSyncBatchDuration.NewTimer();
|
|
|
|
// Phase 1 — lock batch row
|
|
ClinicalSyncBatch? batch;
|
|
try
|
|
{
|
|
await using var lockTx = await _db.Database.BeginTransactionAsync(ct);
|
|
batch = await _db.ClinicalSyncBatches
|
|
.FromSqlInterpolated($"""
|
|
SELECT * FROM clinical_sync_batches
|
|
WHERE id = {batchId}
|
|
FOR UPDATE NOWAIT
|
|
""")
|
|
.FirstOrDefaultAsync(ct);
|
|
|
|
if (batch is null || batch.Status != ClinicalSyncBatchStatus.Received)
|
|
return;
|
|
|
|
batch.MarkProcessing();
|
|
await _db.SaveChangesAsync(ct);
|
|
await lockTx.CommitAsync(ct);
|
|
}
|
|
catch (PostgresException ex) when (ex.SqlState == "55P03") // lock_not_available
|
|
{
|
|
_logger.LogInformation("Batch {BatchId} already locked by another consumer", batchId);
|
|
return;
|
|
}
|
|
|
|
// Phase 2 — deserialize payload
|
|
var request = JsonSerializer.Deserialize<ClinicalSyncBatchRequest>(batch!.Payload)!;
|
|
var hasConflict = false;
|
|
|
|
// Phase 3 — replay order: observations → alerts → acks → resolutions
|
|
foreach (var obs in request.Observations.OrderBy(o => o.RecordedAt))
|
|
{
|
|
try
|
|
{
|
|
await ApplyObservationAsync(obs, batch, ct);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_db.ChangeTracker.Clear();
|
|
await RecordConflictAsync(batch.Id, obs.ClientRef, "OBSERVATION", ex.Message, ct);
|
|
hasConflict = true;
|
|
}
|
|
}
|
|
|
|
foreach (var alert in request.AlertEvents.OrderBy(a => a.GeneratedAt))
|
|
{
|
|
try
|
|
{
|
|
await ApplyAlertEventAsync(alert, batch, ct);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_db.ChangeTracker.Clear();
|
|
await RecordConflictAsync(batch.Id, alert.ClientAlertId, "ALERT", ex.Message, ct);
|
|
hasConflict = true;
|
|
}
|
|
}
|
|
|
|
foreach (var ack in request.AlertAcknowledgments.OrderBy(a => a.AcknowledgedAt))
|
|
{
|
|
try
|
|
{
|
|
if (await ApplyAckAsync(ack, batch, ct))
|
|
hasConflict = true;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_db.ChangeTracker.Clear();
|
|
await RecordConflictAsync(batch.Id, ack.ClientRef, "ACK", ex.Message, ct);
|
|
hasConflict = true;
|
|
}
|
|
}
|
|
|
|
foreach (var resolve in request.AlertResolutions.OrderBy(r => r.ResolvedAt))
|
|
{
|
|
try
|
|
{
|
|
if (await ApplyResolveAsync(resolve, batch, ct))
|
|
hasConflict = true;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
_db.ChangeTracker.Clear();
|
|
await RecordConflictAsync(batch.Id, resolve.ClientRef, "RESOLVE", ex.Message, ct);
|
|
hasConflict = true;
|
|
}
|
|
}
|
|
|
|
// Phase 4 — finalize batch
|
|
batch = await _db.ClinicalSyncBatches.FindAsync([batchId], ct);
|
|
if (batch is null) return;
|
|
|
|
if (hasConflict) batch.MarkConflict();
|
|
else batch.MarkApplied();
|
|
|
|
var gateway = await _db.WardGateways.FindAsync([batch.GatewayId], ct);
|
|
gateway?.MarkSynced(DateTimeOffset.UtcNow);
|
|
|
|
await _db.SaveChangesAsync(ct);
|
|
|
|
_metrics.ClinicalSyncBatchesTotal.WithLabels(batch.Status.ToDbString()).Inc();
|
|
_logger.LogInformation("Batch {BatchId} finalized as {Status}", batchId, batch.Status);
|
|
}
|
|
|
|
private async Task ApplyObservationAsync(
|
|
SyncedObservation obs, ClinicalSyncBatch batch, CancellationToken ct)
|
|
{
|
|
if (await _db.Observations.AnyAsync(o => o.IdempotencyKey == obs.IdempotencyKey, ct))
|
|
return;
|
|
|
|
await _observations.ApplySyncedObservationAsync(obs, ct);
|
|
}
|
|
|
|
private async Task ApplyAlertEventAsync(
|
|
SyncedAlertEvent alert, ClinicalSyncBatch batch, CancellationToken ct)
|
|
{
|
|
if (await _db.ClinicalAlerts.AnyAsync(a => a.ClientAlertId == alert.ClientAlertId, ct))
|
|
return;
|
|
|
|
var encounter = await _db.Encounters.FindAsync([alert.EncounterId], ct)
|
|
?? throw new ValidationException("Encounter not found.", "ENCOUNTER_NOT_FOUND");
|
|
|
|
var clinicalAlert = new ClinicalAlert
|
|
{
|
|
Id = Guid.NewGuid(),
|
|
ClientAlertId = alert.ClientAlertId,
|
|
SyncedFromGateway = true,
|
|
EncounterId = alert.EncounterId,
|
|
PatientId = encounter.PatientId,
|
|
AlertType = AlertTypeExtensions.FromDbString(alert.AlertType),
|
|
Severity = AlertSeverityExtensions.FromDbString(alert.Severity),
|
|
Details = alert.Details,
|
|
Status = AlertStatus.Open,
|
|
TriggeredAt = alert.GeneratedAt
|
|
};
|
|
_db.ClinicalAlerts.Add(clinicalAlert);
|
|
|
|
_db.OutboxEvents.Add(new OutboxEvent
|
|
{
|
|
Id = Guid.NewGuid(),
|
|
Topic = "alert.generated",
|
|
Payload = JsonSerializer.Serialize(new
|
|
{
|
|
alertId = clinicalAlert.Id,
|
|
encounterId = alert.EncounterId,
|
|
patientId = encounter.PatientId,
|
|
alertType = alert.AlertType,
|
|
severity = alert.Severity,
|
|
details = alert.Details,
|
|
syncedFromGateway = true,
|
|
triggeredAt = alert.GeneratedAt,
|
|
partitionKey = alert.EncounterId.ToString()
|
|
}),
|
|
PartitionKey = alert.EncounterId.ToString(),
|
|
CreatedAt = DateTimeOffset.UtcNow
|
|
});
|
|
|
|
await _db.SaveChangesAsync(ct);
|
|
}
|
|
|
|
private async Task<bool> ApplyAckAsync(
|
|
SyncedAlertAcknowledgment ack, ClinicalSyncBatch batch, CancellationToken ct)
|
|
{
|
|
var alert = await _db.ClinicalAlerts
|
|
.FirstOrDefaultAsync(a => a.ClientAlertId == ack.ClientAlertId, ct);
|
|
if (alert is null)
|
|
{
|
|
await RecordConflictAsync(batch.Id, ack.ClientRef, "ACK", "ALERT_NOT_YET_SYNCED", ct);
|
|
return true;
|
|
}
|
|
if (alert.Status != AlertStatus.Open && alert.Status != AlertStatus.Escalated)
|
|
return false;
|
|
|
|
await _alerts.ApplySyncedAcknowledgmentAsync(alert.Id, ack, ct);
|
|
return false;
|
|
}
|
|
|
|
private async Task<bool> ApplyResolveAsync(
|
|
SyncedAlertResolution resolve, ClinicalSyncBatch batch, CancellationToken ct)
|
|
{
|
|
var alert = await _db.ClinicalAlerts
|
|
.FirstOrDefaultAsync(a => a.ClientAlertId == resolve.ClientAlertId, ct);
|
|
if (alert is null)
|
|
{
|
|
await RecordConflictAsync(batch.Id, resolve.ClientRef, "RESOLVE", "ALERT_NOT_YET_SYNCED", ct);
|
|
return true;
|
|
}
|
|
if (alert.Status == AlertStatus.Resolved)
|
|
return false;
|
|
|
|
await _alerts.ApplySyncedResolutionAsync(alert.Id, resolve, ct);
|
|
return false;
|
|
}
|
|
|
|
private async Task RecordConflictAsync(
|
|
Guid batchId, Guid clientRef, string itemType, string reason, CancellationToken ct)
|
|
{
|
|
_db.ClinicalSyncConflicts.Add(new ClinicalSyncConflict(batchId, clientRef, itemType, reason));
|
|
await _db.SaveChangesAsync(ct);
|
|
}
|
|
} |