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 _logger; public ClinicalSyncBatchProcessor( AppDbContext db, IObservationService observations, IAlertService alerts, ClinicalMetrics metrics, ILogger 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(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 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 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); } }