using System.Text.Json; using Microsoft.Extensions.Logging; namespace VFXReviewWorker; /// /// Durable lifecycle reports (§7.6): fail/complete reports are written to the /// spool directory before the first send attempt and replayed strictly in /// order with exponential backoff (1 s → 60 s). Rendering continues during /// server downtime; a report is deleted only once the server accepts it (or /// permanently rejects it as contradictory). /// public sealed class ReportSpooler { private readonly ApiClient _api; private readonly ILogger _log; private readonly string _dir; private readonly SemaphoreSlim _wake = new(0); private long _seq; public ReportSpooler(ApiClient api, ILogger log, string? dir = null) { _api = api; _log = log; _dir = dir ?? WorkerOptions.SpoolDir; Directory.CreateDirectory(_dir); } public void Enqueue(string kind, string jobId, Dictionary payload) { var report = new SpooledReport { Kind = kind, JobId = jobId, Payload = payload }; var name = $"{DateTimeOffset.UtcNow.UtcTicks:D20}_{Interlocked.Increment(ref _seq):D6}.json"; var tmp = Path.Combine(_dir, name + ".tmp"); var final = Path.Combine(_dir, name); File.WriteAllText(tmp, JsonSerializer.Serialize(report, WorkerOptions.JsonOpts)); File.Move(tmp, final); _wake.Release(); } public bool HasPending => Directory.EnumerateFiles(_dir, "*.json").Any(); /// Long-running replay loop; started once by the worker service. public async Task RunAsync(CancellationToken ct) { var backoff = TimeSpan.FromSeconds(1); while (!ct.IsCancellationRequested) { var next = Directory.EnumerateFiles(_dir, "*.json").OrderBy(f => f, StringComparer.Ordinal).FirstOrDefault(); if (next is null) { try { await _wake.WaitAsync(TimeSpan.FromSeconds(5), ct); } catch (OperationCanceledException) { break; } continue; } bool done; try { var report = JsonSerializer.Deserialize(File.ReadAllText(next), WorkerOptions.JsonOpts); if (report is null || string.IsNullOrEmpty(report.JobId)) { _log.LogWarning("Discarding unreadable spool file {File}", next); done = true; } else { done = await _api.SendLifecycleAsync(report.Kind, report.JobId, report.Payload, ct); } } catch (OperationCanceledException) { break; } catch (Exception ex) { _log.LogWarning(ex, "Spool replay attempt failed for {File}", next); done = false; } if (done) { try { File.Delete(next); } catch { /* re-sent next pass; server side is idempotent */ } backoff = TimeSpan.FromSeconds(1); } else { try { await Task.Delay(backoff, ct); } catch (OperationCanceledException) { break; } backoff = TimeSpan.FromSeconds(Math.Min(backoff.TotalSeconds * 2, 60)); } } } }