using System.Net.Http.Headers; using System.Net.Http.Json; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; using StackExchange.Redis; public sealed class EncounterReplicaSyncService : BackgroundService { private const int EncounterPageSize = 100; private readonly IServiceScopeFactory _scopes; private readonly IHttpClientFactory _http; private readonly CentralReachabilityService _reachability; private readonly GatewayOptions _gateway; private readonly ILogger _logger; public EncounterReplicaSyncService( IServiceScopeFactory scopes, IHttpClientFactory http, CentralReachabilityService reachability, IOptions gateway, ILogger logger) { _scopes = scopes; _http = http; _reachability = reachability; _gateway = gateway.Value; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken ct) { await SyncOnceAsync(ct); using var timer = new PeriodicTimer( TimeSpan.FromMinutes(_gateway.EncounterSyncIntervalMinutes)); while (await timer.WaitForNextTickAsync(ct)) { if (_reachability.IsCentralReachable) await SyncOnceAsync(ct); } } internal async Task SyncOnceAsync(CancellationToken ct) { await using var scope = _scopes.CreateAsyncScope(); var db = scope.ServiceProvider.GetRequiredService(); var redis = scope.ServiceProvider.GetRequiredService(); var centralOpts = scope.ServiceProvider.GetRequiredService>().Value; var client = _http.CreateClient("central"); client.BaseAddress = new Uri(centralOpts.BaseUrl); ConfigureAuth(client); if (client.DefaultRequestHeaders.Authorization is null) { _logger.LogWarning( "GATEWAY_SYNC_JWT not set — skipping encounter replica sync (seed local replica manually for dev)"); return; } try { var syncedCount = await SyncEncountersAsync(client, db, ct); await SyncThresholdsAsync(client, db, redis, ct); var state = await db.SyncState.FirstOrDefaultAsync(ct) ?? db.SyncState.Add(new GatewaySyncState { Id = Guid.NewGuid() }).Entity; state.LastEncounterSyncAt = DateTimeOffset.UtcNow; state.InitialSyncCompleted = true; await db.SaveChangesAsync(ct); _logger.LogInformation("Encounter replica sync completed — {Count} encounters", syncedCount); } catch (HttpRequestException ex) { _logger.LogWarning(ex, "Encounter replica sync failed — central API unreachable"); } } private async Task SyncEncountersAsync( HttpClient client, GatewayDbContext db, CancellationToken ct) { var dept = Uri.EscapeDataString(_gateway.Department); var summaries = new List(); var page = 1; while (true) { var listResp = await client.GetFromJsonAsync>( $"/api/v1/encounters?status=ACTIVE&department={dept}&page={page}&pageSize={EncounterPageSize}", CentralApiJson.Options, ct); if (listResp?.Data?.Items is null || listResp.Data.Items.Count == 0) break; summaries.AddRange(listResp.Data.Items); if (page >= listResp.Data.TotalPages) break; page++; } var synced = 0; foreach (var summary in summaries) { var detailResp = await client.GetFromJsonAsync>( $"/api/v1/encounters/{summary.EncounterId}", CentralApiJson.Options, ct); if (detailResp?.Data is null) { _logger.LogWarning("Central returned no detail for encounter {EncounterId}", summary.EncounterId); continue; } UpsertReplica(db, detailResp.Data); synced++; } return synced; } private static void UpsertReplica(GatewayDbContext db, CentralEncounterDetail detail) { var now = DateTimeOffset.UtcNow; var patient = detail.Patient; var replicaPatient = db.Patients.Find(patient.Id); if (replicaPatient is null) { db.Patients.Add(new ReplicaPatient { Id = patient.Id, Mrn = patient.Mrn, FirstName = patient.FirstName, LastName = patient.LastName, DateOfBirth = patient.DateOfBirth, Gender = patient.Gender, SyncedAt = now }); } else { replicaPatient.Mrn = patient.Mrn; replicaPatient.FirstName = patient.FirstName; replicaPatient.LastName = patient.LastName; replicaPatient.DateOfBirth = patient.DateOfBirth; replicaPatient.Gender = patient.Gender; replicaPatient.SyncedAt = now; } var enc = db.Encounters.Find(detail.Id); if (enc is null) { db.Encounters.Add(new ReplicaEncounter { Id = detail.Id, PatientId = detail.PatientId, EncounterType = detail.EncounterType, Status = detail.Status, Department = detail.Department, AttendingPhysician = detail.AttendingPhysician, RoomBed = detail.RoomBed, AdmittedAt = detail.AdmittedAt, SyncedAt = now }); } else { enc.PatientId = detail.PatientId; enc.EncounterType = detail.EncounterType; enc.Status = detail.Status; enc.Department = detail.Department; enc.AttendingPhysician = detail.AttendingPhysician; enc.RoomBed = detail.RoomBed; enc.AdmittedAt = detail.AdmittedAt; enc.SyncedAt = now; } } private async Task SyncThresholdsAsync( HttpClient client, GatewayDbContext db, IConnectionMultiplexer redis, CancellationToken ct) { var thrResp = await client.GetFromJsonAsync>>( "/api/v1/alert-thresholds", CentralApiJson.Options, ct); if (thrResp?.Data is null) return; await db.Database.ExecuteSqlRawAsync("DELETE FROM alert_thresholds", ct); var syncedAt = DateTimeOffset.UtcNow; db.AlertThresholds.AddRange(thrResp.Data.Select(t => new ReplicaAlertThreshold { Id = t.Id, ObservationCode = t.ObservationCode, DisplayName = t.DisplayName, Unit = t.Unit, CriticalLow = t.CriticalLow, WarningLow = t.WarningLow, WarningHigh = t.WarningHigh, CriticalHigh = t.CriticalHigh, SuppressionWindowMinutes = t.SuppressionWindowMinutes, SyncedAt = syncedAt })); await db.SaveChangesAsync(ct); await ThresholdCacheLoader.RefreshAsync(db, redis, _logger, ct); } private static void ConfigureAuth(HttpClient client) { client.DefaultRequestHeaders.Authorization = null; var jwt = Environment.GetEnvironmentVariable("GATEWAY_SYNC_JWT"); if (!string.IsNullOrEmpty(jwt)) client.DefaultRequestHeaders.Authorization = new AuthenticationHeaderValue("Bearer", jwt); } }