using CareFix.Api.Hospitals; using CareFix.Api.Options; using Dapper; using Microsoft.Extensions.Options; namespace CareFix.Api.Infrastructure; /// /// Nightly housekeeping: enforces the retention periods promised in the hospital addendum and /// recaptures schema snapshots that have gone stale (usually after a HIS upgrade). /// Runs once a day at the configured UTC hour; safe to run more than once. /// public sealed class MaintenanceWorker( ControlDb db, HospitalAdminService hospitals, AuditService audit, IOptions options, ILogger log) : BackgroundService { private DateOnly _lastRun = DateOnly.MinValue; protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // Give the service a minute to finish starting before any heavy work. try { await Task.Delay(TimeSpan.FromMinutes(1), stoppingToken); } catch (OperationCanceledException) { return; } while (!stoppingToken.IsCancellationRequested) { var now = DateTime.UtcNow; var today = DateOnly.FromDateTime(now); if (now.Hour == Math.Clamp(options.Value.Maintenance.HourUtc, 0, 23) && _lastRun != today) { _lastRun = today; try { await RunAsync(stoppingToken); } catch (Exception ex) when (ex is not OperationCanceledException) { log.LogError(ex, "Nightly maintenance failed"); } } try { await Task.Delay(TimeSpan.FromMinutes(10), stoppingToken); } catch (OperationCanceledException) { return; } } } /// Also callable on demand from the admin API. public async Task RunAsync(CancellationToken ct) { var purged = await PurgeAsync(ct); var schema = await RecaptureSchemasAsync(ct); var summary = new { purged, schema }; await audit.LogAsync(null, null, "MaintenanceRun", summary, ct); log.LogInformation("Maintenance done: {Summary}", System.Text.Json.JsonSerializer.Serialize(summary)); return summary; } private async Task PurgeAsync(CancellationToken ct) { var r = options.Value.Retention; await using var c = await db.OpenAsync(ct); // Query results and AI transcripts hold masked patient data; clear them first and keep the metadata for reporting. var resultsCleared = await ChunkedAsync(c, "UPDATE TOP (5000) CF_QUERY_LOG SET ResultJson = NULL WHERE ResultJson IS NOT NULL AND CreatedAt < DATEADD(DAY, -@days, SYSUTCDATETIME())", r.QueryResultDays, ct); var transcripts = await ChunkedAsync(c, """ DELETE TOP (2000) t FROM CF_AI_TRANSCRIPT t JOIN CF_TICKET k ON k.TicketId = t.TicketId WHERE k.State = 'Closed' AND t.UpdatedAt < DATEADD(DAY, -@days, SYSUTCDATETIME()) """, r.TranscriptDays, ct); var queryRows = await ChunkedAsync(c, "DELETE TOP (5000) FROM CF_QUERY_LOG WHERE CreatedAt < DATEADD(DAY, -@days, SYSUTCDATETIME())", r.QueryLogDays, ct); // Snapshots are what a rollback restores from, so they outlive the query log. var snapshots = await ChunkedAsync(c, """ DELETE TOP (2000) s FROM CF_ROW_SNAPSHOT s JOIN CF_EXECUTION e ON e.ExecutionId = s.ExecutionId WHERE e.ExecutedAt < DATEADD(DAY, -@days, SYSUTCDATETIME()) """, r.SnapshotDays, ct); var mined = await ChunkedAsync(c, """ DELETE TOP (5000) m FROM CF_MINE_TICKET m JOIN CF_MINE_BATCH b ON b.BatchId = m.BatchId WHERE b.CreatedAt < DATEADD(DAY, -@days, SYSUTCDATETIME()) """, r.LearningDays, ct); var webhooks = await ChunkedAsync(c, "DELETE TOP (5000) FROM CF_WEBHOOK_OUTBOX WHERE State <> 'Pending' AND CreatedAt < DATEADD(DAY, -@days, SYSUTCDATETIME())", 30, ct); // CF_AUDIT is append-only by trigger and is never purged here. Archiving it is a deliberate DBA task. return new { resultsCleared, transcripts, queryRows, snapshots, minedTickets = mined, webhooks }; } private static async Task ChunkedAsync(Microsoft.Data.SqlClient.SqlConnection c, string sql, int days, CancellationToken ct) { var total = 0; while (!ct.IsCancellationRequested) { var n = await c.ExecuteAsync(new CommandDefinition(sql, new { days }, commandTimeout: 120, cancellationToken: ct)); total += n; if (n == 0) break; await Task.Delay(200, ct); // let the live workload breathe } return total; } private async Task RecaptureSchemasAsync(CancellationToken ct) { var days = options.Value.Maintenance.SchemaRecaptureDays; if (days <= 0) return new { attempted = 0, done = 0, failed = 0 }; List due; await using (var c = await db.OpenAsync(ct)) due = (await c.QueryAsync(""" SELECT TOP 25 h.HospitalId FROM CF_HOSPITAL h WHERE h.Status = 'Active' AND (h.SchemaCapturedAt IS NULL OR h.SchemaCapturedAt < DATEADD(DAY, -@days, SYSUTCDATETIME())) AND (EXISTS (SELECT 1 FROM CF_CONNECTION cn WHERE cn.HospitalId = h.HospitalId) OR EXISTS (SELECT 1 FROM CF_AGENT a WHERE a.HospitalId = h.HospitalId AND a.LastSeenAt > DATEADD(MINUTE, -2, SYSUTCDATETIME()))) ORDER BY h.SchemaCapturedAt """, new { days })).ToList(); int done = 0, failed = 0; foreach (var id in due) { try { await hospitals.CaptureSchemaAsync(0, id, ct); done++; } catch (Exception ex) when (ex is not OperationCanceledException) { failed++; log.LogWarning("Schema recapture failed for hospital {HospitalId}: {Message}", id, ex.Message); } } return new { attempted = due.Count, done, failed }; } }