using System.Security.Cryptography; using System.Text; using CareFix.Api.Infrastructure; using CareFix.Core; using Dapper; namespace CareFix.Api.Hospitals; /// Server side of the agent protocol: key issue, authentication, long-poll and results. public sealed class AgentHub(ControlDb db, AuditService audit) { private static readonly TimeSpan PollWindow = TimeSpan.FromSeconds(25); private static byte[] Hash(string key) => SHA256.HashData(Encoding.UTF8.GetBytes(key)); public async Task IssueKeyAsync(int adminId, int hospitalId, CancellationToken ct) { var key = "cfa_" + Convert.ToBase64String(RandomNumberGenerator.GetBytes(32)).Replace('+', '-').Replace('/', '_').TrimEnd('='); await using var c = await db.OpenAsync(ct); if (await c.ExecuteScalarAsync("SELECT COUNT(*) FROM CF_HOSPITAL WHERE HospitalId = @hospitalId", new { hospitalId }) == 0) throw AppException.NotFound("Hospital not found."); await c.ExecuteAsync(""" MERGE CF_AGENT AS t USING (SELECT @hospitalId AS HospitalId) s ON t.HospitalId = s.HospitalId WHEN MATCHED THEN UPDATE SET KeyHash = @h, KeyIssuedAt = SYSUTCDATETIME(), KeyIssuedBy = @adminId WHEN NOT MATCHED THEN INSERT (HospitalId, KeyHash, KeyIssuedBy) VALUES (@hospitalId, @h, @adminId); """, new { hospitalId, h = Hash(key), adminId }); await audit.LogAsync(adminId, null, "AgentKeyIssued", new { hospitalId }, ct); return key; } /// Returns the hospital id for a valid agent, or throws 401. public async Task AuthenticateAsync(HttpRequest req, CancellationToken ct) { var code = req.Headers[AgentProtocol.HospitalHeader].ToString(); var key = req.Headers[AgentProtocol.KeyHeader].ToString(); if (string.IsNullOrEmpty(code) || string.IsNullOrEmpty(key)) throw new AppException(401, "Agent credentials missing."); await using var c = await db.OpenAsync(ct); var row = await c.QuerySingleOrDefaultAsync<(int HospitalId, byte[] KeyHash)>(""" SELECT a.HospitalId, a.KeyHash FROM CF_AGENT a JOIN CF_HOSPITAL h ON h.HospitalId = a.HospitalId WHERE h.Code = @code AND h.Channel = 'Agent' AND h.Status = 'Active' """, new { code }); if (row.KeyHash is null || !CryptographicOperations.FixedTimeEquals(row.KeyHash, Hash(key))) throw new AppException(401, "Agent key not valid for this hospital."); return row.HospitalId; } public async Task PollAsync(int hospitalId, AgentPoll poll, string? ip, CancellationToken ct) { // "CareFix.Agent.exe test" only checks the key; it must never claim a real job. if (poll.Version == "test") return null; await using var c = await db.OpenAsync(ct); await c.ExecuteAsync(""" UPDATE CF_AGENT SET LastSeenAt = SYSUTCDATETIME(), Version = @Version, MachineName = @MachineName, LastIp = @ip WHERE HospitalId = @hospitalId """, new { hospitalId, poll.Version, poll.MachineName, ip }); var until = DateTime.UtcNow + PollWindow; while (!ct.IsCancellationRequested && DateTime.UtcNow < until) { var job = await c.QuerySingleOrDefaultAsync<(long JobId, string Kind, string PayloadJson)>(""" WITH next AS ( SELECT TOP (1) * FROM CF_AGENT_JOB WITH (ROWLOCK, READPAST, UPDLOCK) WHERE HospitalId = @hospitalId AND State = 'Pending' AND ExpiresAt > SYSUTCDATETIME() ORDER BY JobId) UPDATE next SET State = 'Taken', TakenAt = SYSUTCDATETIME() OUTPUT INSERTED.JobId, INSERTED.Kind, INSERTED.PayloadJson; """, new { hospitalId }); if (job.Kind is not null) { using var doc = System.Text.Json.JsonDocument.Parse(job.PayloadJson); return new AgentJob(job.JobId, job.Kind, doc.RootElement.Clone()); } // Check often for the first few seconds (a waiting engineer), then ease off: with hundreds of // agents connected this is the difference between ~1,000 and ~250 queries a second. var quick = DateTime.UtcNow < until - PollWindow + TimeSpan.FromSeconds(5); try { await Task.Delay(quick ? 500 : 2000, ct); } catch (OperationCanceledException) { break; } } return null; } public async Task CompleteAsync(int hospitalId, long jobId, AgentResult result, CancellationToken ct) { await using var c = await db.OpenAsync(ct); var n = await c.ExecuteAsync(""" UPDATE CF_AGENT_JOB SET State = @state, ResultJson = @json, ErrorStatus = @ErrorStatus, Error = @Error, DoneAt = SYSUTCDATETIME() WHERE JobId = @jobId AND HospitalId = @hospitalId AND State = 'Taken' """, new { jobId, hospitalId, state = result.Ok ? "Done" : "Failed", json = result.Ok && result.Result is { } r ? r.GetRawText() : null, result.ErrorStatus, Error = result.Error is { Length: > 2000 } e ? e[..2000] : result.Error, }); if (n == 0) throw AppException.Conflict("Job not found or already finished."); } public async Task> StatusAsync(CancellationToken ct) { await using var c = await db.OpenAsync(ct); return await c.QueryAsync("SELECT HospitalId, LastSeenAt, Version, MachineName, KeyIssuedAt FROM CF_AGENT"); } } /// Removes finished agent jobs after 7 days and clears any leftover results. public sealed class AgentJobJanitor(ControlDb db, ILogger log) : BackgroundService { protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { await using var c = await db.OpenAsync(stoppingToken); await c.ExecuteAsync(""" UPDATE CF_AGENT_JOB SET State = 'Expired' WHERE State = 'Pending' AND ExpiresAt < SYSUTCDATETIME(); UPDATE CF_AGENT_JOB SET ResultJson = NULL WHERE ResultJson IS NOT NULL AND DoneAt < DATEADD(MINUTE, -10, SYSUTCDATETIME()); DELETE CF_AGENT_JOB WHERE CreatedAt < DATEADD(DAY, -7, SYSUTCDATETIME()); """); } catch (Exception ex) when (ex is not OperationCanceledException) { log.LogWarning(ex, "Agent job cleanup failed"); } try { await Task.Delay(TimeSpan.FromMinutes(5), stoppingToken); } catch (OperationCanceledException) { } } } }