namespace ChangeT50.Agent; public sealed class AgentRunner { private readonly AgentApiClient _api; private readonly IChangeT50Provider _provider; private readonly LocalJournal _journal; public AgentRunner(AgentApiClient api, IChangeT50Provider provider, LocalJournal journal) { _api = api; _provider = provider; _journal = journal; } public async Task RunOnceAsync(CancellationToken cancellationToken) { var health = await _provider.GetHealthAsync(_journal.PendingUploadCount(), cancellationToken); await _api.PostHeartbeatAsync(health, cancellationToken); var summary = new AgentRunSummary(); var tasks = await _api.GetTasksAsync(cancellationToken: cancellationToken); foreach (var task in tasks.Tasks) { summary.Pulled++; await ProcessNewTaskAsync(task, summary, cancellationToken); } var recovery = await _api.GetRecoveryTasksAsync(cancellationToken: cancellationToken); foreach (var task in recovery.Tasks) { summary.RecoveryPulled++; await RecoverTaskAsync(task, summary, cancellationToken); } return summary; } private async Task ProcessNewTaskAsync(AgentTask task, AgentRunSummary summary, CancellationToken cancellationToken) { if (!string.Equals(task.TaskType, "consume", StringComparison.OrdinalIgnoreCase)) { return; } var acceptance = await _api.AcceptTaskAsync( task.ClientOrderNo, $"{task.ClientOrderNo}-{Guid.NewGuid():N}", cancellationToken); summary.Accepted++; var journal = _journal.Get(task.ClientOrderNo) ?? new JournalEntry { ClientOrderNo = task.ClientOrderNo, LeaseToken = acceptance.LeaseToken, State = "accepted" }; journal.LeaseToken = acceptance.LeaseToken; _journal.Save(journal); // If a previous process already wrote the card, the local journal is the // source of truth. Never call DebitCard a second time for this order. var consume = journal.Consume; if (consume is null) { consume = await _provider.ConsumeAsync(task, cancellationToken); journal.Consume = consume; journal.State = consume.Result == "success" ? "card_debited" : "card_write_failed"; _journal.Save(journal); } await _api.PostConsumeResultAsync(consume, journal.LeaseToken, cancellationToken); summary.ConsumeResults++; if (consume.Result != "success") { return; } await UploadExistingAsync(task, journal, summary, cancellationToken); } private async Task RecoverTaskAsync(AgentTask task, AgentRunSummary summary, CancellationToken cancellationToken) { var journal = _journal.Get(task.ClientOrderNo); if (journal?.Consume is { Result: "success" }) { if (task.RecoveryAction == "inspect_local_journal" && !string.IsNullOrEmpty(journal.LeaseToken)) { await _api.PostConsumeResultAsync(journal.Consume, journal.LeaseToken, cancellationToken); summary.ConsumeResults++; } await UploadExistingAsync(task, journal, summary, cancellationToken); return; } // There is no safe local proof that the card was or was not debited. // Record an unknown reconciliation result; never retry the debit. await _api.PostReconciliationAsync( $"unknown-{task.ClientOrderNo}-{DateTimeOffset.UtcNow.ToUnixTimeSeconds()}", task.ClientOrderNo, "unknown", "local_journal_missing", task.AmountCents, task.AmountCents, cancellationToken); summary.ManualReview++; } private async Task UploadExistingAsync(AgentTask task, JournalEntry journal, AgentRunSummary summary, CancellationToken cancellationToken) { var retryCount = (journal.Upload?.RetryCount ?? task.RetryCount) + 1; var upload = await _provider.UploadAsync(task, journal.Consume!, retryCount, cancellationToken); journal.Upload = upload; journal.State = upload.Result == "success" ? "settled" : "upload_retrying"; _journal.Save(journal); await _api.PostUploadResultAsync(task.ClientOrderNo, upload, cancellationToken); summary.UploadResults++; if (upload.Result == "success") { summary.Settled++; } } } public sealed class AgentRunSummary { public int Pulled { get; set; } public int Accepted { get; set; } public int ConsumeResults { get; set; } public int UploadResults { get; set; } public int RecoveryPulled { get; set; } public int Settled { get; set; } public int ManualReview { get; set; } }