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