Files
vigilcare-clinical/VigilCareClinicalAPI/Services/ClinicalSyncBatchProcessor.cs
T

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);
}
}