diff --git a/CONTEXT.md b/CONTEXT.md index fa3e696..a5e961f 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -226,6 +226,14 @@ _Avoid_: Release (that word belongs to the Concurrency-Limit slot), return the L One execution try of a job, numbered and visible to the handler. A lease expiry counts as an attempt, the same as a thrown exception. _Avoid_: Retry (retry is attempts after the first; counting "retries" invites off-by-one ambiguity) +**Retry Cause**: +Why a job last went back to Scheduled because an Attempt went wrong: the handler failed and the retry policy scheduled another Attempt (**handler failed**), or the Lease lapsed before the worker reported an outcome (**lease expired**). Recorded on the job row by the store, atomically with the reschedule. Sticky: a claim, a Relinquish, a Cancel and a terminal outcome leave it as it is (a terminal job keeps it as a record of its last retry); only an operator requeue clears it, because the requeued job starts over. A new job has none. It names the kind of problem only; the message is the Failure Detail on the Transition Log. +_Avoid_: Retry reason, last error (it is a kind, not a message), Attempt > 0 (a relinquished or requeued job has Attempts but no problem) + +**Retrying**: +A Scheduled job that carries a Retry Cause: still live, waiting for another Attempt because the last one went wrong. Not a state of its own but a view of Scheduled, so an operator sees trouble before it ends as Dead-Lettered. A new job, a requeued job and a job a stopping worker handed back are Scheduled but not Retrying, so a deploy raises no false alarm. Leaves the view when the job is claimed again. +_Avoid_: Failed (the job is not terminal), Retried (it has not run again yet), Retry state (there is no such state) + **At-Least-Once Execution**: BackWave's delivery contract: a job's handler body may run more than once, and idempotency is the handler author's responsibility. The field standard — Hangfire, Sidekiq, River, Celery (acks-late), RabbitMQ-with-acks, and Temporal *activities* all land here. Exactly-once *body* execution is not offered, because the only two roads to it are both rejected: at-most-once (accept job loss on crash) or durable execution. What BackWave still guarantees exactly once is the Effect-Once property. _Avoid_: Exactly-once execution, at-most-once, deliver-once diff --git a/src/BackWave.Conformance/ConformanceSuite.cs b/src/BackWave.Conformance/ConformanceSuite.cs index daaa74f..2f23f14 100644 --- a/src/BackWave.Conformance/ConformanceSuite.cs +++ b/src/BackWave.Conformance/ConformanceSuite.cs @@ -1878,6 +1878,196 @@ public async Task Clause_5_9_CountMatchingJobs_EmptyStore_IsZero() new JobQuery { State = JobState.Scheduled, TagPredicates = [JobTagPredicate.HasLabel("urgent")] })); } + // ── §5.9 Retrying: the jobs an attempt went wrong for ─────────────────────── + + private static async Task> ListRetryingAsync(IJobStore store, string? queue = null) + => await store.ListJobsAsync(new JobQuery { Retrying = true, Queue = queue, SortDirection = JobSortDirection.OldestFirst }); + + /// + /// Certifies that a handler failure the retry policy reschedules records the HandlerFailed cause, so + /// the Scheduled job lists as Retrying, and that the Retrying filter ANDs with the other filters. + /// + [Fact] + public async Task Clause_5_9_Retrying_AFailureRetry_IsRetrying_WithTheHandlerFailedCause() + { + var store = await CreateStoreAsync(); + await store.EnqueueAsync(Job(), now: T0); + var claimed = Assert.Single(await ClaimAsync(store, T0)); + + var retryAt = T0.AddMinutes(5); + await store.ReportOutcomeAsync(claimed.JobId, "w1", claimed.Attempt, new JobOutcome.Failure(retryAt, "transient"), T0); + + var job = await store.GetJobAsync(claimed.JobId); + Assert.Equal(JobState.Scheduled, job!.State); + Assert.Equal(RetryCause.HandlerFailed, job.RetryCause); + var listed = Assert.Single(await ListRetryingAsync(store)); + Assert.Equal(claimed.JobId, listed.JobId); + Assert.Equal(RetryCause.HandlerFailed, listed.RetryCause); + Assert.Equal(retryAt, listed.DueTime); + Assert.Equal(1, listed.Attempt); + Assert.Empty(await ListRetryingAsync(store, queue: "other")); + } + + /// + /// Certifies that the batched outcome path records the same HandlerFailed cause as a single report + /// for a failure it reschedules, and records none for a success or a dead-letter in the same batch. + /// + [Fact] + public async Task Clause_5_9_Retrying_ABatchedFailureRetry_IsRetrying_LikeASingleReport() + { + var store = await CreateStoreAsync(); + var retried = Job(); + var succeeded = Job(); + var dead = Job(); + await store.EnqueueAsync(retried, now: T0); + await store.EnqueueAsync(succeeded, now: T0); + await store.EnqueueAsync(dead, now: T0); + var claimed = await ClaimAsync(store, T0); + OutcomeReport Report(NewJob job, JobOutcome outcome) + => new(job.JobId, "w1", claimed.Single(j => j.JobId == job.JobId).Attempt, outcome); + + await store.ReportOutcomesAsync( + [ + Report(retried, new JobOutcome.Failure(T0.AddMinutes(5), "transient")), + Report(succeeded, new JobOutcome.Success()), + Report(dead, new JobOutcome.Failure(null, "fatal")), + ], T0); + + Assert.Equal(RetryCause.HandlerFailed, (await store.GetJobAsync(retried.JobId))!.RetryCause); + Assert.Null((await store.GetJobAsync(succeeded.JobId))!.RetryCause); + Assert.Null((await store.GetJobAsync(dead.JobId))!.RetryCause); + Assert.Equal(retried.JobId, Assert.Single(await ListRetryingAsync(store)).JobId); + } + + /// + /// Certifies that a lapsed lease the expiry sweep reschedules records the LeaseExpired cause, so the + /// Scheduled job lists as Retrying. + /// + [Fact] + public async Task Clause_5_9_Retrying_ALeaseExpiry_IsRetrying_WithTheLeaseExpiredCause() + { + var store = await CreateStoreAsync(); + await store.EnqueueAsync(Job(), now: T0); + var claimed = Assert.Single(await ClaimAsync(store, T0)); + + var afterExpiry = T0 + Lease + TimeSpan.FromSeconds(1); + Assert.Equal(1, await store.ExpireLeasesAsync(afterExpiry, maxJobs: 32, DefaultQueues, TwoAttempts)); + + var job = await store.GetJobAsync(claimed.JobId); + Assert.Equal(JobState.Scheduled, job!.State); + Assert.Equal(RetryCause.LeaseExpired, job.RetryCause); + var listed = Assert.Single(await ListRetryingAsync(store)); + Assert.Equal(claimed.JobId, listed.JobId); + Assert.Equal(RetryCause.LeaseExpired, listed.RetryCause); + } + + /// + /// Certifies that a clean-stop hand-back is not a problem: a job that never went wrong comes back + /// Scheduled without a cause and does not list as Retrying, while a job that was already Retrying + /// keeps its cause through the hand-back, so a deploy neither raises nor hides an alarm. + /// + [Fact] + public async Task Clause_5_9_Retrying_ARelinquish_IsNotRetrying_AndKeepsAnEarlierCause() + { + var store = await CreateStoreAsync(); + var healthy = Job(); + var failing = Job(); + await store.EnqueueAsync(healthy, now: T0); + await store.EnqueueAsync(failing, now: T0); + var first = await ClaimAsync(store, T0); + var failed = first.Single(j => j.JobId == failing.JobId); + await store.ReportOutcomeAsync(failed.JobId, "w1", failed.Attempt, new JobOutcome.Failure(T0, "transient"), T0); + Assert.Single(await ClaimAsync(store, T0)); // the failing job's second attempt; both are now leased + + // Three attempts, so the failing job's second attempt is below the ceiling and is handed back + // rather than dead-lettered. + var threeAttempts = new RetryPolicy { MaxAttempts = 3, Backoff = _ => TimeSpan.FromMinutes(1) }.ToDisposition(); + var handBack = T0.AddSeconds(5); + if (!Declares(ConformanceCapabilities.LeaseRelinquish)) + { + await AssertRelinquishIsANoOpAsync(store, handBack, threeAttempts, healthy.JobId, failing.JobId); + return; + } + Assert.Equal(2, await store.RelinquishLeasesAsync("w1", handBack, threeAttempts)); + + var handedBack = await store.GetJobAsync(healthy.JobId); + Assert.Equal(JobState.Scheduled, handedBack!.State); + Assert.Null(handedBack.RetryCause); + Assert.Equal(RetryCause.HandlerFailed, (await store.GetJobAsync(failing.JobId))!.RetryCause); + Assert.Equal(failing.JobId, Assert.Single(await ListRetryingAsync(store)).JobId); + } + + /// + /// Certifies that a terminal outcome keeps the last retry cause as a record, and that an operator + /// requeue clears it: the requeued job starts over and does not list as Retrying. + /// + [Fact] + public async Task Clause_5_9_Retrying_ARequeue_IsNotRetrying_AndClearsTheCause() + { + var store = await CreateStoreAsync(); + await store.EnqueueAsync(Job(), now: T0); + var first = Assert.Single(await ClaimAsync(store, T0)); + await store.ReportOutcomeAsync(first.JobId, "w1", first.Attempt, new JobOutcome.Failure(T0, "transient"), T0); + var second = Assert.Single(await ClaimAsync(store, T0)); + await store.ReportOutcomeAsync(second.JobId, "w1", second.Attempt, new JobOutcome.Failure(null, "fatal"), T0); + + var dead = await store.GetJobAsync(first.JobId); + Assert.Equal(JobState.DeadLettered, dead!.State); + Assert.Equal(RetryCause.HandlerFailed, dead.RetryCause); // kept as the record of the last retry + Assert.Empty(await ListRetryingAsync(store)); // terminal, so not Retrying + + var requeueTime = T0.AddMinutes(1); + Assert.Equal(RequeueResult.Requeued, await store.RequeueAsync(first.JobId, "alice", requeueTime)); + + var requeued = await store.GetJobAsync(first.JobId); + Assert.Equal(JobState.Scheduled, requeued!.State); + Assert.Null(requeued.RetryCause); + Assert.Empty(await ListRetryingAsync(store)); + } + + /// + /// Certifies that new work is never Retrying: a job enqueued due now, one enqueued for later, and one + /// awaiting a parent all carry no cause. + /// + [Fact] + public async Task Clause_5_9_Retrying_AFreshEnqueue_IsNotRetrying() + { + var store = await CreateStoreAsync(); + var dueNow = Job(); + var later = Job(dueTime: T0.AddHours(1)); + await store.EnqueueAsync(dueNow, now: T0); + await store.EnqueueAsync(later, now: T0); + var child = Job() with { Parents = [dueNow.JobId] }; + await store.EnqueueAsync(child, now: T0); + + foreach (var id in (Guid[])[dueNow.JobId, later.JobId, child.JobId]) + { + Assert.Null((await store.GetJobAsync(id))!.RetryCause); + } + Assert.Empty(await ListRetryingAsync(store)); + } + + /// + /// Certifies that a Retrying job leaves the Retrying list once it is claimed again, and stays gone + /// when that attempt succeeds. + /// + [Fact] + public async Task Clause_5_9_Retrying_ClaimedAgainAndSucceeded_LeavesTheRetryingList() + { + var store = await CreateStoreAsync(); + await store.EnqueueAsync(Job(), now: T0); + var first = Assert.Single(await ClaimAsync(store, T0)); + await store.ReportOutcomeAsync(first.JobId, "w1", first.Attempt, new JobOutcome.Failure(T0, "transient"), T0); + Assert.Single(await ListRetryingAsync(store)); + + var second = Assert.Single(await ClaimAsync(store, T0)); + Assert.Empty(await ListRetryingAsync(store)); // running, not waiting + + await store.ReportOutcomeAsync(second.JobId, "w1", second.Attempt, new JobOutcome.Success(), T0); + Assert.Equal(JobState.Succeeded, (await store.GetJobAsync(first.JobId))!.State); + Assert.Empty(await ListRetryingAsync(store)); + } + // ── Mutation teeth: boundary/misc contract facts (issue 0235) ──────────────── /// diff --git a/src/BackWave.Dashboard/Components/JobDetailPanel.razor b/src/BackWave.Dashboard/Components/JobDetailPanel.razor index 1077c89..8c2c167 100644 --- a/src/BackWave.Dashboard/Components/JobDetailPanel.razor +++ b/src/BackWave.Dashboard/Components/JobDetailPanel.razor @@ -22,6 +22,10 @@ { Terminal cause@cause } + @if (Job.RetryCause is { } retryCause) + { + Retry cause@DashboardGlossary.RetryCauseName(retryCause) + } @if (Job.ScheduleId is { } scheduleId) { diff --git a/src/BackWave.Dashboard/Components/JobTable.razor b/src/BackWave.Dashboard/Components/JobTable.razor index 436c132..43b2055 100644 --- a/src/BackWave.Dashboard/Components/JobTable.razor +++ b/src/BackWave.Dashboard/Components/JobTable.razor @@ -15,8 +15,16 @@ else Queue State Attempt - Due Time - Terminal At + @if (ShowRetryCause) + { + Next Attempt + Retry Cause + } + else + { + Due Time + Terminal At + } @if (ShowTags) { Tags @@ -42,8 +50,16 @@ else @job.Attempt - @DashboardGlossary.Instant(job.DueTime) - @(job.TerminalAt is { } terminalAt ? DashboardGlossary.Instant(terminalAt) : "—") + @if (ShowRetryCause) + { + @DashboardGlossary.Instant(job.DueTime) + @(job.RetryCause is { } retryCause ? DashboardGlossary.RetryCauseName(retryCause) : "-") + } + else + { + @DashboardGlossary.Instant(job.DueTime) + @(job.TerminalAt is { } terminalAt ? DashboardGlossary.Instant(terminalAt) : "—") + } @if (TagHref is { } tagHref) { @* Tag pills (ADR 0022, issue 0113): an empty tag set renders NOTHING — no empty @@ -95,5 +111,11 @@ else /// [Parameter] public Func? TagHref { get; set; } + /// + /// Show the Retrying columns: the due time reads as the next attempt, and the retry cause (handler + /// failed or lease expired) takes the place of the terminal instant, which a live job never has. + /// + [Parameter] public bool ShowRetryCause { get; set; } + private bool ShowTags => TagHref is not null; } diff --git a/src/BackWave.Dashboard/Components/Pages/Failures.razor b/src/BackWave.Dashboard/Components/Pages/Failures.razor index c5d5238..74e403b 100644 --- a/src/BackWave.Dashboard/Components/Pages/Failures.razor +++ b/src/BackWave.Dashboard/Components/Pages/Failures.razor @@ -1,7 +1,9 @@ @* Failures: Dead-Lettered and Quarantined as distinct categories (invariant I5, never collapsed) — ran-and-kept-failing vs could-not-be-routed/decoded. Two tabs so a long Dead-Lettered list never buries the Quarantined one; the active tab rides the URL - (?tab=quarantine) so it survives each live SSE tick. Via the Monitor API. *@ + (?tab=quarantine) so it survives each live SSE tick. A third tab (?tab=retrying) lists the + Retrying jobs: still live, but waiting for another attempt because the handler failed or the + lease expired. Via the Monitor API. *@
@@ -11,17 +13,22 @@
- @Tab(JobState.DeadLettered, "Dead-Lettered", $"{BasePath}/failures", DeadLettered.Count) - @Tab(JobState.Quarantined, "Quarantined", $"{BasePath}/failures?tab=quarantine", Quarantined.Count) + @Tab(Category.DeadLettered, "Dead-Lettered", $"{BasePath}/failures", DeadLettered.Count) + @Tab(Category.Quarantined, "Quarantined", $"{BasePath}/failures?tab=quarantine", Quarantined.Count) + @Tab(Category.Retrying, "Retrying", $"{BasePath}/failures?tab=retrying", Retrying.Count)
- @if (ShowQuarantined) + @switch (Active) { - @Section(JobState.Quarantined, "Quarantined", Quarantined) - } - else - { - @Section(JobState.DeadLettered, "Dead-Lettered", DeadLettered) + case Category.Quarantined: + @Section(JobState.Quarantined, "Quarantined", Quarantined) + break; + case Category.Retrying: + @RetryingSection + break; + default: + @Section(JobState.DeadLettered, "Dead-Lettered", DeadLettered) + break; }
@@ -30,11 +37,15 @@ [Parameter, EditorRequired] public string BasePath { get; set; } = ""; [Parameter, EditorRequired] public IReadOnlyList DeadLettered { get; set; } = []; [Parameter, EditorRequired] public IReadOnlyList Quarantined { get; set; } = []; + + /// The Retrying jobs: Scheduled again because the handler failed or the lease expired. + [Parameter, EditorRequired] public IReadOnlyList Retrying { get; set; } = []; [Parameter, EditorRequired] public int PageSize { get; set; } [Parameter, EditorRequired] public DashboardActions Actions { get; set; } = DashboardActions.None; - /// Which category tab is open — "quarantine" shows Quarantined, anything else (the default) - /// shows Dead-Lettered. Carried in the URL so it survives every live SSE re-render. + /// Which category tab is open - "quarantine" shows Quarantined, "retrying" shows Retrying, + /// anything else (the default) shows Dead-Lettered. Carried in the URL so it survives every live SSE + /// re-render. [Parameter] public string ActiveTab { get; set; } = ""; /// Render only the live region (SSE fragment) rather than the full document. @@ -43,7 +54,12 @@ /// Wrap in the #bw-live region and inline the SSE client. [Parameter] public bool Live { get; set; } - private bool ShowQuarantined => string.Equals(ActiveTab, "quarantine", StringComparison.OrdinalIgnoreCase); + private enum Category { DeadLettered, Quarantined, Retrying } + + private Category Active => + string.Equals(ActiveTab, "quarantine", StringComparison.OrdinalIgnoreCase) ? Category.Quarantined + : string.Equals(ActiveTab, "retrying", StringComparison.OrdinalIgnoreCase) ? Category.Retrying + : Category.DeadLettered; // Requeue: a Dead-Lettered/Quarantined job back to Scheduled (Attempt reset). Rendered // only where AuthorizeRequeue passed. @@ -60,10 +76,10 @@ // many" — shown as "50+" rather than an exact count that would under-report. private string TabCount(int count) => count >= PageSize ? $"{PageSize}+" : count.ToString(); - private RenderFragment Tab(JobState state, string label, string href, int count) =>@ + private RenderFragment Tab(Category category, string label, string href, int count) =>@ @label@TabCount(count) - @if (ShowQuarantined == (state == JobState.Quarantined)) + @if (Active == category) { } @@ -76,4 +92,15 @@

Search all @title jobs →

} ; + + // A Retrying job is live, so it has nothing to requeue; the row shows when it runs next and why it + // came back instead. A tag pill deep-links to the Jobs list filtered to Retrying AND the clicked Tag. + private RenderFragment RetryingSection =>@
+ + @if (Retrying.Count >= PageSize) + { +

Search all Retrying jobs →

+ } +
; } diff --git a/src/BackWave.Dashboard/Components/Pages/Jobs.razor b/src/BackWave.Dashboard/Components/Pages/Jobs.razor index d4518bb..d284df2 100644 --- a/src/BackWave.Dashboard/Components/Pages/Jobs.razor +++ b/src/BackWave.Dashboard/Components/Pages/Jobs.razor @@ -1,4 +1,4 @@ -@* Jobs list: filters (state, Queue, Wire Name, Recurring Schedule) and §5.9 cursor +@* Jobs list: filters (state or Retrying, Queue, Wire Name, Recurring Schedule) and §5.9 cursor pagination (?after={sequence}), ported verbatim from the prior hand-written view onto the render spine. All data arrives via the Monitor API. *@ @@ -20,6 +20,12 @@ @foreach (var state in DashboardGlossary.StateOrder) { + @if (state == JobState.Scheduled) + { + @* Retrying narrows Scheduled to the jobs waiting for another + attempt because the handler failed or the lease expired. *@ + + } } @Chevron @@ -158,7 +164,7 @@ }
- +
@if (NextHref is { } next) @@ -243,6 +249,10 @@ { query.Add($"state={state}"); } + if (Filter.Retrying) + { + query.Add($"state={DashboardGlossary.RetryingFilterValue}"); + } if (Filter.Queue is { } queue) { query.Add($"queue={Uri.EscapeDataString(queue)}"); @@ -274,6 +284,10 @@ { parts.Add($"state={state}"); } + if (Filter.Retrying) + { + parts.Add($"state={DashboardGlossary.RetryingFilterValue}"); + } if (Filter.Queue is { } queue) { parts.Add($"queue={Uri.EscapeDataString(queue)}"); diff --git a/src/BackWave.Dashboard/DashboardGlossary.cs b/src/BackWave.Dashboard/DashboardGlossary.cs index 07c3a16..654962b 100644 --- a/src/BackWave.Dashboard/DashboardGlossary.cs +++ b/src/BackWave.Dashboard/DashboardGlossary.cs @@ -118,6 +118,20 @@ static bool TryFixed(string field, int min, int max, out int value) /// public static string TagLabel(JobTag tag) => tag.IsLabel ? tag.Value : $"{tag.Key}:{tag.Value}"; + /// + /// The State-filter value that selects Retrying jobs. Retrying is not a state: it is a Scheduled job + /// that carries a retry cause, so it is waiting for another attempt because something went wrong. + /// + public const string RetryingFilterValue = "Retrying"; + + /// A retry cause in plain words, for the Retrying list and the job detail. + public static string RetryCauseName(RetryCause cause) => cause switch + { + RetryCause.HandlerFailed => "Handler failed", + RetryCause.LeaseExpired => "Lease expired", + _ => cause.ToString(), + }; + /// Terminal states are settled; only non-terminal jobs can be Cancelled. public static bool IsTerminal(JobState state) => state is JobState.Succeeded or JobState.Cancelled or JobState.DeadLettered or JobState.Quarantined; diff --git a/src/BackWave.Dashboard/DashboardRequestHandler.cs b/src/BackWave.Dashboard/DashboardRequestHandler.cs index 1e79557..f732e31 100644 --- a/src/BackWave.Dashboard/DashboardRequestHandler.cs +++ b/src/BackWave.Dashboard/DashboardRequestHandler.cs @@ -11,6 +11,7 @@ using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Http.Features; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; namespace BackWave.Dashboard; @@ -446,11 +447,17 @@ private static async Task RenderOnceAsync(HttpContext context, LiveView view, Di /// Holds the response open and pushes the re-rendered #bw-live fragment as Server-Sent /// Events, every , but only when the markup changed since the last /// push (a comment heartbeat keeps the connection warm otherwise). Ends when the browser - /// disconnects. + /// disconnects or the application starts to stop. private static async Task StreamAsync( HttpContext context, LiveView view, TimeSpan interval, Dictionary? seed) { - var ct = context.RequestAborted; + // The server waits for open requests before the hosted services stop, so a stream that only + // ends on disconnect makes an open dashboard tab spend the host's whole shutdown window. The + // worker groups then have no time left to give their leases back on a clean stop. + var stopping = context.RequestServices.GetService()?.ApplicationStopping + ?? CancellationToken.None; + using var streamEnd = CancellationTokenSource.CreateLinkedTokenSource(context.RequestAborted, stopping); + var ct = streamEnd.Token; context.Response.ContentType = "text/event-stream"; context.Response.Headers.CacheControl = "no-cache"; context.Response.Headers["X-Accel-Buffering"] = "no"; // don't let a reverse proxy buffer the stream @@ -496,7 +503,7 @@ private static async Task StreamAsync( } catch (OperationCanceledException) { - // The browser navigated away or closed the tab — a normal end to the stream. + // The browser navigated away or closed the tab, or the application is stopping - a normal end to the stream. } } @@ -538,15 +545,25 @@ private static async Task JobsAsync(HttpContext context, BackWaveMonitor monitor ? parsedSize : PageSize; + // Retrying is not a state but a narrower view of Scheduled (a job waiting for another attempt + // because the handler failed or the lease expired), offered in the same State filter. JobState? state = null; + var retrying = false; if (query["state"] is [{ Length: > 0 } rawState]) { - if (!Enum.TryParse(rawState, ignoreCase: true, out var parsed)) + if (string.Equals(rawState, DashboardGlossary.RetryingFilterValue, StringComparison.OrdinalIgnoreCase)) + { + retrying = true; + } + else if (!Enum.TryParse(rawState, ignoreCase: true, out var parsed)) { await BadRequestAsync(context, $"Unknown state '{rawState}'.").ConfigureAwait(false); return; } - state = parsed; + else + { + state = parsed; + } } long? after = null; if (query["after"] is [{ Length: > 0 } rawAfter]) @@ -570,6 +587,7 @@ private static async Task JobsAsync(HttpContext context, BackWaveMonitor monitor Queue = NonEmpty(query["queue"]), WireName = NonEmpty(query["wire"]), ScheduleId = NonEmpty(query["schedule"]), + Retrying = retrying, TagPredicates = tagPredicates, AfterSequence = after, SortDirection = JobSortDirection.NewestFirst, // historical table: most recent jobs first @@ -709,10 +727,15 @@ private static async Task JobDetailAsync( // Glossary distinction, never collapsed (invariant I5): Dead-Lettered jobs ran and // kept failing; Quarantined jobs could not be routed or decoded. Both lists load every // tick — the inactive tab still shows a live count badge — but only the active tab's - // table renders, so a long Dead-Lettered list never buries the Quarantined one. + // table renders, so a long Dead-Lettered list never buries the Quarantined one. Retrying + // jobs are the third list: still live, waiting for another attempt because the handler + // failed or the lease expired, so trouble shows before it ends in a dead letter. Oldest + // job first, so a job stuck in a retry loop does not sink below newer ones. async () => new Dictionary { ["BasePath"] = basePath, + ["Retrying"] = await monitor.ListJobsAsync( + new JobQuery { Retrying = true, MaxResults = PageSize }).ConfigureAwait(false), ["DeadLettered"] = await monitor.ListJobsAsync( new JobQuery { State = JobState.DeadLettered, SortDirection = JobSortDirection.NewestFirst, MaxResults = PageSize }).ConfigureAwait(false), ["Quarantined"] = await monitor.ListJobsAsync( diff --git a/src/BackWave.Oracle/OracleJobStore.cs b/src/BackWave.Oracle/OracleJobStore.cs index c76afc4..d0955c3 100644 --- a/src/BackWave.Oracle/OracleJobStore.cs +++ b/src/BackWave.Oracle/OracleJobStore.cs @@ -872,7 +872,7 @@ private async ValueTask ReportOutcomeUntracedAsync( } }), 3), JobOutcome.Failure { NextDueTime: { } retryAt } => - ("state = 0, due_time = :retryAt, lease_owner = NULL, lease_expiry = NULL", + ("state = 0, due_time = :retryAt, lease_owner = NULL, lease_expiry = NULL, retry_cause = 1", command => command.Parameters.Add(Tstz("retryAt", retryAt)), 0), JobOutcome.Failure failure => ("state = 5, lease_owner = NULL, lease_expiry = NULL, terminal_at = :now, terminal_cause = :cause", @@ -1071,6 +1071,7 @@ FOR UPDATE // authorizes every write - the verdict above only decides what the caller is told. due_time // moves only for a retry row (COALESCE keeps it otherwise); cancel_requested clears only for a // Cancelled row (CASE); terminal_at and terminal_cause carry per row and are null for a retry. + // retry_cause records a handler-failure retry and is left alone on every terminal row. // Both instants travel as ISO text under an explicit format. A JSON_TABLE column declared // TIMESTAMP WITH TIME ZONE takes second precision 6 and rounds away the seventh digit, which is // a digit this store hands back; TO_TIMESTAMP_TZ over the text keeps all of them. @@ -1098,7 +1099,8 @@ WHEN MATCHED THEN UPDATE SET j.terminal_at = d.terminal_at, j.terminal_cause = d.cause, j.due_time = COALESCE(d.due, j.due_time), - j.cancel_requested = CASE WHEN d.state = 4 THEN 0 ELSE j.cancel_requested END + j.cancel_requested = CASE WHEN d.state = 4 THEN 0 ELSE j.cancel_requested END, + j.retry_cause = CASE WHEN d.state = 0 THEN 1 ELSE j.retry_cause END WHERE j.state = 2 AND j.lease_owner = d.worker AND j.attempt = d.attempt AND j.lease_expiry > :now """, @@ -1540,7 +1542,7 @@ job_hex VARCHAR2(32) PATH '$.JobHex', due VARCHAR2(40) PATH '$.Due')) d) d ON (j.job_id = d.job_id) WHEN MATCHED THEN UPDATE SET - j.state = 0, j.due_time = d.due, j.lease_owner = NULL, j.lease_expiry = NULL + j.state = 0, j.due_time = d.due, j.lease_owner = NULL, j.lease_expiry = NULL, j.retry_cause = 2 """, connection, transaction); reschedule.Parameters.Add(Clob("payload", JsonSerializer.Serialize( @@ -1856,7 +1858,7 @@ public async ValueTask RequeueAsync( """ UPDATE backwave.jobs SET state = 0, attempt = 0, due_time = :now, lease_owner = NULL, lease_expiry = NULL, - cancel_requested = 0, terminal_at = NULL, terminal_cause = NULL + cancel_requested = 0, terminal_at = NULL, terminal_cause = NULL, retry_cause = NULL WHERE job_id = :id AND state IN (5, 6) """, connection, transaction); @@ -2629,6 +2631,10 @@ private static void AppendScopeConditions(JobQuery query, List condition conditions.Add("schedule_id = :scheduleId"); command.Parameters.Add(Str("scheduleId", scheduleId)); } + if (query.Retrying) + { + conditions.Add("state = 0 AND retry_cause IS NOT NULL"); + } for (var i = 0; i < query.TagPredicates.Count; i++) { var predicate = query.TagPredicates[i]; @@ -3855,7 +3861,7 @@ public async ValueTask> ListObserverDead private const string JobColumns = "job_id, wire_name, payload, queue, state, due_time, attempt, lease_owner, lease_expiry, " + "cancel_requested, terminal_at, terminal_cause, schedule_id, parents_remaining, job_mode, trace_context, " + - "sequence, workflow_id"; + "sequence, workflow_id, retry_cause"; private static JobRecord ReadJob(OracleDataReader reader) { @@ -3893,9 +3899,28 @@ private static JobRecord ReadJob(OracleDataReader reader) TraceContext = reader.IsDBNull(15) ? null : reader.GetString(15), Sequence = reader.GetInt64(16), WorkflowId = reader.IsDBNull(17) ? null : ReadGuid(reader, 17), + RetryCause = ReadRetryCause(reader, 18), }; } + // The retry cause is nullable (no cause is NULL, never a number), and a stored number outside the + // enum surfaces as the named violation, the same as an undefined state. + private static RetryCause? ReadRetryCause(OracleDataReader reader, int ordinal) + { + if (reader.IsDBNull(ordinal)) + { + return null; + } + var storedCause = reader.GetInt32(ordinal); + if (!Enum.IsDefined((RetryCause)storedCause)) + { + throw Invariant.Halt( + InvariantTrigger.UndefinedEnumValueStored, + $"Job {ReadGuid(reader, 0)} stores retry cause {storedCause}, which is not a defined RetryCause."); + } + return (RetryCause)storedCause; + } + // ── Job Tags ────────────────────────────────────────────────────────────────── /// diff --git a/src/BackWave.Oracle/OracleMigrator.cs b/src/BackWave.Oracle/OracleMigrator.cs index 92b270f..7aa206a 100644 --- a/src/BackWave.Oracle/OracleMigrator.cs +++ b/src/BackWave.Oracle/OracleMigrator.cs @@ -17,7 +17,7 @@ namespace BackWave.Oracle; public static class OracleMigrator { /// The schema version this build of the adapter requires the database to be at. - public const int ExpectedSchemaVersion = 1; + public const int ExpectedSchemaVersion = 2; // Transient connection faults a cold-booting fleet can hit that the bounded retry should ride out // rather than surface: the shared listener/handshake-storm and connection-lost connectivity set. @@ -100,8 +100,19 @@ public static async Task MigrateAsync( private static async Task ApplyScriptsAsync( string connectionString, SchemaRewriter rewriter, CancellationToken cancellationToken) { - await using var connection = new OracleConnection(connectionString); + // A later script alters a table that other nodes can still be creating indexes and constraints on + // during a cold boot, or that live workers write to during an upgrade. DDL therefore waits for the + // table lock instead of failing at once with ORA-00054. The session does not go back to the pool, + // so the wait never reaches a store connection. + var unpooled = new OracleConnectionStringBuilder(connectionString) { Pooling = false }; + await using var connection = new OracleConnection(unpooled.ConnectionString); await connection.OpenAsync(cancellationToken).ConfigureAwait(false); + await using (var wait = connection.CreateCommand()) + { + wait.CommandText = "ALTER SESSION SET ddl_lock_timeout = 30"; + // uncounted round trip: part of the one-time migration, like the scripts below. + await wait.ExecuteNonQueryAsync(cancellationToken).ConfigureAwait(false); + } var assembly = typeof(OracleMigrator).Assembly; var scripts = assembly.GetManifestResourceNames() diff --git a/src/BackWave.Oracle/Schema/0002_retry_cause.sql b/src/BackWave.Oracle/Schema/0002_retry_cause.sql new file mode 100644 index 0000000..64fa536 --- /dev/null +++ b/src/BackWave.Oracle/Schema/0002_retry_cause.sql @@ -0,0 +1,64 @@ +-- BackWave schema v2 (Oracle dialect). Idempotent: safe to run on every deploy. +-- v1 -> v2: record why a job went back to Scheduled, so a Retrying job is told apart from a new one. +-- +-- retry_cause: why the job most recently went back to Scheduled after an attempt went wrong. +-- 1 HandlerFailed (the handler failed and the retry policy scheduled another attempt), +-- 2 LeaseExpired (the lease lapsed before the worker reported an outcome), NULL none. +-- A Scheduled job with a cause is Retrying. A requeue clears it; a clean-stop hand-back of the lease +-- leaves it alone, so a deploy never makes a healthy job look like a failing one. +-- +-- Nullable with no default, so the ADD is a dictionary-only change that rewrites no rows, and an N-1 +-- node that neither reads nor writes the column keeps working. Existing rows read NULL: a job already +-- retrying when the fleet upgrades shows as Retrying from its next failed attempt or expired lease. +-- +-- One anonymous PL/SQL block. Unlike v1 it alters a table that live workers write to, so it differs in +-- two ways: +-- * It returns at once when the schema is already at v2 or later. A boot against a current schema +-- then takes no DDL lock on the jobs table, and an N-1 node that boots after the upgrade never +-- stamps the version back down. +-- * ddl() also retries ORA-00054 and ORA-14411 (another session runs DDL on the same table). The +-- migrator session already waits for the table lock, so the ALTER queues behind the in-flight claims +-- and outcome reports of a running fleet. When a fleet cold-boots, other nodes can still be creating +-- v1's indexes on the jobs table. After the winner commits, every other node gets "already exists" +-- and no-ops. + +DECLARE + deployed NUMBER; + + -- Runs one DDL statement and ignores the "already exists" family, so a re-run converges. + PROCEDURE ddl(statement IN VARCHAR2) IS + BEGIN + FOR attempt IN 1 .. 600 LOOP + BEGIN + EXECUTE IMMEDIATE statement; + RETURN; + EXCEPTION + WHEN OTHERS THEN + -- -955 name already used, -1430 column exists, -1408 index column list already indexed. + IF SQLCODE IN (-955, -1430, -1408) THEN + RETURN; + -- -54 resource busy, -14411 concurrent DDL on the same object: wait, then try again. + ELSIF SQLCODE IN (-54, -14411) AND attempt < 600 THEN + DBMS_SESSION.SLEEP(0.1); + ELSE + RAISE; + END IF; + END; + END LOOP; + END; +BEGIN + SELECT MAX(version) INTO deployed FROM backwave.schema_version; + IF deployed >= 2 THEN + RETURN; + END IF; + + ddl(q'{ALTER TABLE backwave.jobs ADD (retry_cause NUMBER(10) NULL)}'); + + -- The Retrying listing. A single-column index omits NULL keys, so only a row with a cause enters it, + -- and a claim (which changes state, not retry_cause) never touches it. + ddl(q'{CREATE INDEX backwave.ix_bw_jobs_retrying ON backwave.jobs (retry_cause)}'); + + -- v1 seeds the row with an INSERT guarded by WHERE NOT EXISTS, so every later version stamps by + -- UPDATE. The guard keeps a later version's stamp in place. + EXECUTE IMMEDIATE q'{UPDATE backwave.schema_version SET version = 2 WHERE version < 2}'; +END; diff --git a/src/BackWave.Postgres/PostgresJobStore.cs b/src/BackWave.Postgres/PostgresJobStore.cs index 474e30a..e2ef8d1 100644 --- a/src/BackWave.Postgres/PostgresJobStore.cs +++ b/src/BackWave.Postgres/PostgresJobStore.cs @@ -476,7 +476,7 @@ FROM candidates c RETURNING j.job_id, j.wire_name, j.payload, j.queue, j.state, j.due_time, j.attempt, j.lease_owner, j.lease_expiry, j.cancel_requested, j.terminal_at, j.terminal_cause, j.schedule_id, j.parents_remaining, j.mode, j.trace_context, - j.sequence, j.workflow_id, + j.sequence, j.workflow_id, j.retry_cause, -- Job Tags (ADR 0022) ride back with the claim as a correlated aggregate in -- THIS round-trip — never a second SELECT — so the no-tags hot path pays only a -- PK-indexed empty lookup (NULL) and is never N+1. The empty-string-key => Label @@ -502,10 +502,10 @@ FROM candidates c InvariantTrigger.ClaimedRowNotLeasedToWorker, $"Claim returned job {job.JobId} in state {job.State} leased to '{job.LeaseOwner}'; the same statement had just set Leased to '{request.WorkerId}'."); } - // Column 18 is the correlated tag aggregate (json or NULL when the job has none). - if (!reader.IsDBNull(18)) + // Column 19 is the correlated tag aggregate (json or NULL when the job has none). + if (!reader.IsDBNull(19)) { - job = job with { Tags = ParseTagsJson(reader.GetString(18)) }; + job = job with { Tags = ParseTagsJson(reader.GetString(19)) }; } queueClaims.Add(job); } @@ -675,7 +675,7 @@ private async ValueTask ReportOutcomeUntracedAsync( } })), JobOutcome.Failure { NextDueTime: { } retryAt } => - ("state = 0, due_time = @retryAt, lease_owner = NULL, lease_expiry = NULL", + ("state = 0, due_time = @retryAt, lease_owner = NULL, lease_expiry = NULL, retry_cause = 1", command => command.Parameters.AddWithValue("retryAt", retryAt.ToUniversalTime())), JobOutcome.Failure failure => ("state = 5, lease_owner = NULL, lease_expiry = NULL, terminal_at = @now, terminal_cause = @cause", @@ -828,6 +828,7 @@ private async ValueTask> ReportOutcomesUntrac // nothing (StaleLease); a matched row applies and is returned via RETURNING, keyed by job id. // due_time moves only for a retry row (COALESCE keeps it for everyone else); cancel_requested // clears only for a Cancelled row. terminal_at/terminal_cause carry per-row (null for retry). + // retry_cause records a handler-failure retry and is left alone on every terminal row. var matched = new Dictionary(); await using (var update = Cmd( """ @@ -838,7 +839,8 @@ UPDATE backwave.jobs j terminal_at = d.terminal_at, terminal_cause = d.cause, due_time = COALESCE(d.due, j.due_time), - cancel_requested = CASE WHEN d.state = 4 THEN false ELSE j.cancel_requested END + cancel_requested = CASE WHEN d.state = 4 THEN false ELSE j.cancel_requested END, + retry_cause = CASE WHEN d.state = 0 THEN 1 ELSE j.retry_cause END FROM unnest(@ids::uuid[], @workers::text[], @attempts::int[], @states::int[], @causes::text[], @dues::timestamptz[], @terminalAts::timestamptz[]) AS d(job_id, worker, attempt, state, cause, due, terminal_at) @@ -1219,7 +1221,7 @@ await RecordTransitionsBatchAsync(connection, transaction, transitions, now, can await using var reschedule = Cmd( """ UPDATE backwave.jobs j - SET state = 0, due_time = d.due, lease_owner = NULL, lease_expiry = NULL + SET state = 0, due_time = d.due, lease_owner = NULL, lease_expiry = NULL, retry_cause = 2 FROM unnest(@ids::uuid[], @dues::timestamptz[]) AS d(job_id, due) WHERE j.job_id = d.job_id """, @@ -1519,7 +1521,7 @@ public async ValueTask RequeueAsync( """ UPDATE backwave.jobs SET state = 0, attempt = 0, due_time = @now, lease_owner = NULL, lease_expiry = NULL, - cancel_requested = false, terminal_at = NULL, terminal_cause = NULL + cancel_requested = false, terminal_at = NULL, terminal_cause = NULL, retry_cause = NULL WHERE job_id = @id AND state IN (5, 6) RETURNING job_id """, @@ -2217,6 +2219,10 @@ private static void AppendScopeConditions(JobQuery query, List condition conditions.Add("schedule_id = @scheduleId"); command.Parameters.AddWithValue("scheduleId", scheduleId); } + if (query.Retrying) + { + conditions.Add("state = 0 AND retry_cause IS NOT NULL"); + } for (var i = 0; i < query.TagPredicates.Count; i++) { var predicate = query.TagPredicates[i]; @@ -3491,7 +3497,7 @@ public async ValueTask DisposeAsync() private const string JobColumns = "job_id, wire_name, payload, queue, state, due_time, attempt, lease_owner, lease_expiry, " + "cancel_requested, terminal_at, terminal_cause, schedule_id, parents_remaining, mode, trace_context, " + - "sequence, workflow_id"; + "sequence, workflow_id, retry_cause"; private static JobRecord ReadJob(NpgsqlDataReader reader) { @@ -3529,6 +3535,7 @@ private static JobRecord ReadJob(NpgsqlDataReader reader) TraceContext = reader.IsDBNull(15) ? null : reader.GetString(15), Sequence = reader.GetInt64(16), WorkflowId = reader.IsDBNull(17) ? null : reader.GetGuid(17), + RetryCause = ReadRetryCause(reader, 18), }; } @@ -3546,6 +3553,24 @@ private static JobState ReadState(NpgsqlDataReader reader, int ordinal) return (JobState)storedState; } + // The retry cause is nullable (no cause is NULL, never a number), and a stored number outside the + // enum surfaces as the named violation, the same as an undefined state. + private static RetryCause? ReadRetryCause(NpgsqlDataReader reader, int ordinal) + { + if (reader.IsDBNull(ordinal)) + { + return null; + } + var storedCause = reader.GetInt32(ordinal); + if (!Enum.IsDefined((RetryCause)storedCause)) + { + throw Invariant.Halt( + InvariantTrigger.UndefinedEnumValueStored, + $"Job {reader.GetGuid(0)} stores retry cause {storedCause}, which is not a defined RetryCause."); + } + return (RetryCause)storedCause; + } + // ── Job Tags (ADR 0022) ───────────────────────────────────────────────────── /// diff --git a/src/BackWave.Postgres/PostgresMigrator.cs b/src/BackWave.Postgres/PostgresMigrator.cs index e191190..cdef439 100644 --- a/src/BackWave.Postgres/PostgresMigrator.cs +++ b/src/BackWave.Postgres/PostgresMigrator.cs @@ -17,7 +17,7 @@ namespace BackWave.Postgres; public static class PostgresMigrator { /// The schema version this build of the adapter requires. - public const int ExpectedSchemaVersion = 1; + public const int ExpectedSchemaVersion = 2; // Reserved advisory-lock classid for migration coordination (ADR 0046). pg_advisory_xact_lock has // a two-int32 key space that is DISJOINT from the single-bigint per-queue config lock (issue 0193), diff --git a/src/BackWave.Postgres/Schema/0002_retry_cause.sql b/src/BackWave.Postgres/Schema/0002_retry_cause.sql new file mode 100644 index 0000000..410a4a5 --- /dev/null +++ b/src/BackWave.Postgres/Schema/0002_retry_cause.sql @@ -0,0 +1,21 @@ +-- BackWave schema v2 (Postgres dialect). Idempotent: safe to run on every deploy. +-- v1 -> v2: record why a job went back to Scheduled, so a Retrying job is told apart from a new one. + +-- retry_cause: why the job most recently went back to Scheduled after an attempt went wrong. +-- 1 HandlerFailed (the handler failed and the retry policy scheduled another attempt), +-- 2 LeaseExpired (the lease lapsed before the worker reported an outcome), NULL none. +-- A Scheduled job with a cause is Retrying. A requeue clears it; a clean-stop hand-back of the lease +-- leaves it alone, so a deploy never makes a healthy job look like a failing one. +-- +-- Nullable with no default, so the ADD is a catalog-only change that rewrites no rows, and an N-1 node +-- that neither reads nor writes the column keeps working. Existing rows read NULL: a job already +-- retrying when the fleet upgrades shows as Retrying from its next failed attempt or expired lease. +ALTER TABLE backwave.jobs ADD COLUMN IF NOT EXISTS retry_cause int NULL; + +-- The Retrying listing, in sequence order. Only a Scheduled row with a cause enters the index, so a +-- claim (which moves the row out of Scheduled) and every healthy job leave it untouched. +CREATE INDEX IF NOT EXISTS ix_backwave_jobs_retrying + ON backwave.jobs (sequence) WHERE state = 0 AND retry_cause IS NOT NULL; + +-- v1 seeds the row with an INSERT guarded by WHERE NOT EXISTS, so every later version stamps by UPDATE. +UPDATE backwave.schema_version SET version = 2; diff --git a/src/BackWave.Pro.Mcp/Tools/JobTools.cs b/src/BackWave.Pro.Mcp/Tools/JobTools.cs index 22ef782..bd10fce 100644 --- a/src/BackWave.Pro.Mcp/Tools/JobTools.cs +++ b/src/BackWave.Pro.Mcp/Tools/JobTools.cs @@ -37,6 +37,8 @@ internal sealed class JobTools( public async Task SearchJobsAsync( [Description("Only jobs in this state: Scheduled, AwaitingParent, Leased, Succeeded, Cancelled, DeadLettered, or Quarantined. Omit to match any state.")] string? state = null, + [Description("When true, only Retrying jobs: Scheduled jobs waiting for another attempt because the handler failed or the lease expired. A new job, a requeued job, and a job a worker handed back on a clean stop are not Retrying. Combine with state only as Scheduled; any other state matches nothing. Omit or false to match jobs whether Retrying or not.")] + bool? retrying = null, [Description("Only jobs on this queue. Omit to match any queue.")] string? queue = null, [Description("Only jobs of this wire name (the job type's stable string identity; list_wire_names enumerates them). Omit to match any type.")] @@ -82,6 +84,7 @@ _ when sort.Equals("oldest_first", StringComparison.OrdinalIgnoreCase) => JobSor var query = new JobQuery { State = ParseState(state), + Retrying = retrying ?? false, Queue = queue, WireName = wire_name, ScheduleId = schedule_id, @@ -463,6 +466,10 @@ internal sealed record JobRow [Description("A short reason for the terminal outcome (for example why it was dead-lettered); null while still active.")] public string? TerminalCause { get; init; } + /// Why the job last went back to Scheduled after an attempt went wrong; null when none has. + [Description("Why the job last went back to Scheduled after an attempt went wrong: HandlerFailed or LeaseExpired. Null when no attempt has gone wrong since the job was enqueued or last requeued. A Scheduled job with a retry cause is Retrying.")] + public string? RetryCause { get; init; } + /// The recurring schedule that minted this instance, when any. [Description("The recurring schedule that minted this instance; null for a directly enqueued job.")] public string? ScheduleId { get; init; } @@ -492,6 +499,7 @@ internal sealed record JobRow CancelRequested = snapshot.CancelRequested, TerminalAt = snapshot.TerminalAt, TerminalCause = snapshot.TerminalCause, + RetryCause = snapshot.RetryCause?.ToString(), ScheduleId = snapshot.ScheduleId, Sequence = snapshot.Sequence, WorkflowId = snapshot.WorkflowId, diff --git a/src/BackWave.SqlServer/Schema/0003_retry_cause.sql b/src/BackWave.SqlServer/Schema/0003_retry_cause.sql new file mode 100644 index 0000000..365bd7a --- /dev/null +++ b/src/BackWave.SqlServer/Schema/0003_retry_cause.sql @@ -0,0 +1,25 @@ +-- BackWave schema v3 (SQL Server dialect). Idempotent: safe to run on every deploy. +-- v2 -> v3: record why a job went back to Scheduled, so a Retrying job is told apart from a new one. + +-- retry_cause: why the job most recently went back to Scheduled after an attempt went wrong. +-- 1 HandlerFailed (the handler failed and the retry policy scheduled another attempt), +-- 2 LeaseExpired (the lease lapsed before the worker reported an outcome), NULL none. +-- A Scheduled job with a cause is Retrying. A requeue clears it; a clean-stop hand-back of the lease +-- leaves it alone, so a deploy never makes a healthy job look like a failing one. +-- +-- Nullable with no default, so the ADD is a metadata-only change that rewrites no rows, and an N-1 node +-- that neither reads nor writes the column keeps working. Existing rows read NULL: a job already +-- retrying when the fleet upgrades shows as Retrying from its next failed attempt or expired lease. +IF COL_LENGTH('backwave.jobs', 'retry_cause') IS NULL + ALTER TABLE backwave.jobs ADD retry_cause int NULL; + +-- The Retrying listing, in sequence order. Only a Scheduled row with a cause enters the index, so a +-- claim (which moves the row out of Scheduled) and every healthy job leave it untouched. Run through +-- EXEC because the script is one batch, and a statement that names a column added earlier in the same +-- batch does not compile. +IF NOT EXISTS (SELECT 1 FROM sys.indexes WHERE name = 'ix_backwave_jobs_retrying') + EXEC('CREATE INDEX ix_backwave_jobs_retrying + ON backwave.jobs (sequence) WHERE state = 0 AND retry_cause IS NOT NULL'); + +-- v1 seeds the row with an INSERT guarded by WHERE NOT EXISTS, so every later version stamps by UPDATE. +UPDATE backwave.schema_version SET version = 3; diff --git a/src/BackWave.SqlServer/SqlServerJobStore.cs b/src/BackWave.SqlServer/SqlServerJobStore.cs index 9ad73fa..f120d27 100644 --- a/src/BackWave.SqlServer/SqlServerJobStore.cs +++ b/src/BackWave.SqlServer/SqlServerJobStore.cs @@ -570,7 +570,8 @@ UPDATE j inserted.state, inserted.due_time, inserted.attempt, inserted.lease_owner, inserted.lease_expiry, inserted.cancel_requested, inserted.terminal_at, inserted.terminal_cause, inserted.schedule_id, inserted.parents_remaining, - inserted.mode, inserted.trace_context, inserted.[sequence], inserted.workflow_id + inserted.mode, inserted.trace_context, inserted.[sequence], inserted.workflow_id, + inserted.retry_cause FROM backwave.jobs j INNER JOIN candidates c ON j.job_id = c.job_id """, @@ -806,7 +807,7 @@ private async ValueTask ReportOutcomeUntracedAsync( } })), JobOutcome.Failure { NextDueTime: { } retryAt } => - ("state = 0, due_time = @retryAt, lease_owner = NULL, lease_expiry = NULL", + ("state = 0, due_time = @retryAt, lease_owner = NULL, lease_expiry = NULL, retry_cause = 1", command => command.Parameters.AddWithValue("retryAt", retryAt)), JobOutcome.Failure failure => ("state = 5, lease_owner = NULL, lease_expiry = NULL, terminal_at = @now, terminal_cause = @cause", @@ -956,7 +957,8 @@ private async ValueTask> ReportOutcomesUntrac // no longer live simply fails to join and changes nothing (StaleLease); a matched row applies // and is returned via OUTPUT, keyed by job id. due_time moves only for a retry row (COALESCE // keeps it for everyone else); cancel_requested clears only for a Cancelled row (CASE). - // terminal_at/terminal_cause carry per-row (null for a retry). + // terminal_at/terminal_cause carry per-row (null for a retry). retry_cause records a + // handler-failure retry and is left alone on every terminal row. var matched = new Dictionary(); // The payload leads the join and INNER LOOP JOIN pins the shape, so this seeks the clustered // PK once per row and locks only the batch's own jobs. Left to itself the optimizer reads no @@ -972,7 +974,8 @@ UPDATE j terminal_at = d.terminal_at, terminal_cause = d.cause, due_time = COALESCE(d.due, j.due_time), - cancel_requested = CASE WHEN d.state = 4 THEN 0 ELSE j.cancel_requested END + cancel_requested = CASE WHEN d.state = 4 THEN 0 ELSE j.cancel_requested END, + retry_cause = CASE WHEN d.state = 0 THEN 1 ELSE j.retry_cause END OUTPUT inserted.job_id, inserted.state FROM OPENJSON(@payload) WITH (job_id uniqueidentifier '$.JobId', worker nvarchar(450) '$.WorkerId', @@ -1369,7 +1372,7 @@ await RecordTransitionsBatchAsync(connection, transaction, transitions, now, can var rows = string.Join(", ", retries.Select((_, i) => $"(@rid{i}, @rdue{i})")); await using var reschedule = Cmd( $""" - UPDATE j SET state = 0, due_time = d.due, lease_owner = NULL, lease_expiry = NULL + UPDATE j SET state = 0, due_time = d.due, lease_owner = NULL, lease_expiry = NULL, retry_cause = 2 FROM (VALUES {rows}) AS d(job_id, due) INNER LOOP JOIN backwave.jobs j ON j.job_id = d.job_id """, @@ -1675,7 +1678,7 @@ public async ValueTask RequeueAsync( """ UPDATE backwave.jobs SET state = 0, attempt = 0, due_time = @now, lease_owner = NULL, lease_expiry = NULL, - cancel_requested = 0, terminal_at = NULL, terminal_cause = NULL + cancel_requested = 0, terminal_at = NULL, terminal_cause = NULL, retry_cause = NULL OUTPUT inserted.job_id WHERE job_id = @id AND state IN (5, 6) """, @@ -2395,6 +2398,10 @@ private static void AppendScopeConditions(JobQuery query, List condition conditions.Add("schedule_id = @scheduleId"); command.Parameters.Add("scheduleId", SqlDbType.NVarChar, 450).Value = scheduleId; } + if (query.Retrying) + { + conditions.Add("state = 0 AND retry_cause IS NOT NULL"); + } for (var i = 0; i < query.TagPredicates.Count; i++) { var predicate = query.TagPredicates[i]; @@ -3631,7 +3638,7 @@ public async ValueTask> ListObserverDead private const string JobColumns = "job_id, wire_name, payload, queue, state, due_time, attempt, lease_owner, lease_expiry, " + "cancel_requested, terminal_at, terminal_cause, schedule_id, parents_remaining, mode, trace_context, " + - "[sequence], workflow_id"; + "[sequence], workflow_id, retry_cause"; private static JobRecord ReadJob(SqlDataReader reader) { @@ -3669,9 +3676,28 @@ private static JobRecord ReadJob(SqlDataReader reader) TraceContext = reader.IsDBNull(15) ? null : reader.GetString(15), Sequence = reader.GetInt64(16), WorkflowId = reader.IsDBNull(17) ? null : reader.GetGuid(17), + RetryCause = ReadRetryCause(reader, 18), }; } + // The retry cause is nullable (no cause is NULL, never a number), and a stored number outside the + // enum surfaces as the named violation, the same as an undefined state. + private static RetryCause? ReadRetryCause(SqlDataReader reader, int ordinal) + { + if (reader.IsDBNull(ordinal)) + { + return null; + } + var storedCause = reader.GetInt32(ordinal); + if (!Enum.IsDefined((RetryCause)storedCause)) + { + throw Invariant.Halt( + InvariantTrigger.UndefinedEnumValueStored, + $"Job {reader.GetGuid(0)} stores retry cause {storedCause}, which is not a defined RetryCause."); + } + return (RetryCause)storedCause; + } + // Every read of a state column an out-of-band write can reach goes through here: a value outside the // enum surfaces as the named violation, never as a cast that hands the caller an undefined JobState. private static JobState ReadState(SqlDataReader reader, int ordinal) diff --git a/src/BackWave.SqlServer/SqlServerMigrator.cs b/src/BackWave.SqlServer/SqlServerMigrator.cs index 8cd41aa..f51248c 100644 --- a/src/BackWave.SqlServer/SqlServerMigrator.cs +++ b/src/BackWave.SqlServer/SqlServerMigrator.cs @@ -16,7 +16,7 @@ namespace BackWave.SqlServer; public static class SqlServerMigrator { /// The schema version this build of the adapter requires the database to be at. - public const int ExpectedSchemaVersion = 2; + public const int ExpectedSchemaVersion = 3; /// /// Runs every schema script in version order, bringing the database up to the version this diff --git a/src/BackWave.Sqlite/Schema/0003_retry_cause.sql b/src/BackWave.Sqlite/Schema/0003_retry_cause.sql new file mode 100644 index 0000000..e74cdfd --- /dev/null +++ b/src/BackWave.Sqlite/Schema/0003_retry_cause.sql @@ -0,0 +1,23 @@ +-- BackWave SQLite schema v3. Runs once: the migrator skips a step the file already carries. +-- v2 -> v3: record why a job went back to Scheduled, so a Retrying job is told apart from a new one. + +-- retry_cause: why the job most recently went back to Scheduled after an attempt went wrong. +-- 1 HandlerFailed (the handler failed and the retry policy scheduled another attempt), +-- 2 LeaseExpired (the lease lapsed before the worker reported an outcome), NULL none. +-- A Scheduled job with a cause is Retrying. A requeue clears it; a clean-stop hand-back of the lease +-- leaves it alone, so a deploy never makes a healthy job look like a failing one. +-- +-- SQLite has no ADD COLUMN IF NOT EXISTS, so this step is not safe to run twice; the migrator runs only +-- the steps above the version the file is stamped at. Nullable with no default, so the ADD rewrites no +-- rows, and an N-1 node that neither reads nor writes the column keeps working. Existing rows read NULL: +-- a job already retrying when the file upgrades shows as Retrying from its next failed attempt or +-- expired lease. +ALTER TABLE backwave_jobs ADD COLUMN retry_cause INTEGER NULL; + +-- The Retrying listing, in sequence order. Only a Scheduled row with a cause enters the index, so a +-- claim (which moves the row out of Scheduled) and every healthy job leave it untouched. +CREATE INDEX IF NOT EXISTS ix_backwave_jobs_retrying + ON backwave_jobs (sequence) WHERE state = 0 AND retry_cause IS NOT NULL; + +-- v1 seeds the row with an INSERT guarded by WHERE NOT EXISTS, so every later version stamps by UPDATE. +UPDATE backwave_schema_version SET version = 3; diff --git a/src/BackWave.Sqlite/SqliteJobStore.cs b/src/BackWave.Sqlite/SqliteJobStore.cs index 6f32ba4..66a3776 100644 --- a/src/BackWave.Sqlite/SqliteJobStore.cs +++ b/src/BackWave.Sqlite/SqliteJobStore.cs @@ -685,7 +685,7 @@ private static (string Sql, Action Configure) BuildOutcomeUpdate( } })), JobOutcome.Failure { NextDueTime: { } retryAt } => - ($"state = {(int)JobState.Scheduled}, due_time = $retryAt, lease_owner = NULL, lease_expiry = NULL", + ($"state = {(int)JobState.Scheduled}, due_time = $retryAt, lease_owner = NULL, lease_expiry = NULL, retry_cause = {(int)RetryCause.HandlerFailed}", command => command.Parameters.AddWithValue("$retryAt", SqliteValueCodec.ToTicks(retryAt))), JobOutcome.Failure failure => ($"state = {(int)JobState.DeadLettered}, lease_owner = NULL, lease_expiry = NULL, terminal_at = $now, terminal_cause = $cause", @@ -1069,7 +1069,8 @@ ORDER BY lease_expiry await using var reschedule = Cmd( $""" UPDATE backwave_jobs - SET state = {(int)JobState.Scheduled}, due_time = $due, lease_owner = NULL, lease_expiry = NULL + SET state = {(int)JobState.Scheduled}, due_time = $due, lease_owner = NULL, lease_expiry = NULL, + retry_cause = {(int)RetryCause.LeaseExpired} WHERE job_id = $id """, connection, transaction); @@ -1359,7 +1360,7 @@ public async ValueTask RequeueAsync( $""" UPDATE backwave_jobs SET state = {(int)JobState.Scheduled}, attempt = 0, due_time = $now, lease_owner = NULL, lease_expiry = NULL, - cancel_requested = 0, terminal_at = NULL, terminal_cause = NULL + cancel_requested = 0, terminal_at = NULL, terminal_cause = NULL, retry_cause = NULL WHERE job_id = $id AND state IN ({(int)JobState.DeadLettered}, {(int)JobState.Quarantined}) RETURNING job_id """, @@ -2082,6 +2083,10 @@ private static void AppendScopeConditions(JobQuery query, List condition conditions.Add("schedule_id = $scheduleId"); command.Parameters.AddWithValue("$scheduleId", scheduleId); } + if (query.Retrying) + { + conditions.Add($"state = {(int)JobState.Scheduled} AND retry_cause IS NOT NULL"); + } for (var i = 0; i < query.TagPredicates.Count; i++) { var predicate = query.TagPredicates[i]; @@ -3265,7 +3270,7 @@ private sealed class NoopSubscription : IAsyncDisposable private const string JobColumns = "sequence, job_id, wire_name, payload, trace_context, queue, state, due_time, attempt, " + "lease_owner, lease_expiry, cancel_requested, terminal_at, terminal_cause, schedule_id, " + - "parents_remaining, mode, workflow_id"; + "parents_remaining, mode, workflow_id, retry_cause"; private static JobRecord ReadJob(SqliteDataReader reader) => new() { @@ -3287,6 +3292,7 @@ private sealed class NoopSubscription : IAsyncDisposable ParentsRemaining = reader.GetInt32(15), Mode = SqliteValueCodec.ToEnum(reader.GetInt64(16)), WorkflowId = reader.IsDBNull(17) ? null : SqliteValueCodec.ToGuid(reader.GetString(17)), + RetryCause = reader.IsDBNull(18) ? null : SqliteValueCodec.ToEnum(reader.GetInt64(18)), }; // ── Job Tags (ADR 0022) ───────────────────────────────────────────────────── diff --git a/src/BackWave.Sqlite/SqliteMigrator.cs b/src/BackWave.Sqlite/SqliteMigrator.cs index 1ce8dd1..df67182 100644 --- a/src/BackWave.Sqlite/SqliteMigrator.cs +++ b/src/BackWave.Sqlite/SqliteMigrator.cs @@ -18,7 +18,7 @@ namespace BackWave.Sqlite; public static class SqliteMigrator { /// The schema version this build of the adapter requires the database to be at. - public const int ExpectedSchemaVersion = 2; + public const int ExpectedSchemaVersion = 3; // 3.35 is the floor that ships UPDATE … RETURNING, which the claim path relies on (ADR 0019). internal static readonly Version MinimumEngineVersion = new(3, 35, 0); @@ -131,15 +131,21 @@ public static async Task MigrateAsync( // Runs every embedded schema script in version order on the given connection, optionally inside a // transaction. Shared by the coordinated (in-transaction) and opt-out (autocommit) paths. The WAL - // pragma is intentionally NOT here — it runs once, before, outside any transaction. + // pragma is intentionally NOT here - it runs once, before, outside any transaction. A script is + // step N of the schema (its position in version order), and a step at or below the version the file + // is already stamped at is skipped: SQLite has no ADD COLUMN IF NOT EXISTS, so a step that adds a + // column cannot be made safe to run twice in SQL alone. private static async Task ApplyScriptsAsync( SqliteConnection connection, SqliteTransaction? transaction, SchemaRewriter rewriter, CancellationToken cancellationToken) { + var deployed = await ReadDeployedVersionAsync(connection, transaction, rewriter, cancellationToken) + .ConfigureAwait(false); var assembly = typeof(SqliteMigrator).Assembly; var scripts = assembly.GetManifestResourceNames() .Where(name => name.EndsWith(".sql", StringComparison.Ordinal)) - .OrderBy(name => name, StringComparer.Ordinal); + .OrderBy(name => name, StringComparer.Ordinal) + .Skip((int)Math.Max(deployed, 0)); foreach (var script in scripts) { @@ -166,6 +172,13 @@ private static async Task ApplyScriptsAsync( private static async Task IsSchemaCurrentAsync( SqliteConnection connection, SqliteTransaction? transaction, SchemaRewriter rewriter, CancellationToken cancellationToken) + => await ReadDeployedVersionAsync(connection, transaction, rewriter, cancellationToken).ConfigureAwait(false) + >= ExpectedSchemaVersion; + + // The version the file is stamped at, or 0 when it carries no BackWave schema yet. + private static async Task ReadDeployedVersionAsync( + SqliteConnection connection, SqliteTransaction? transaction, SchemaRewriter rewriter, + CancellationToken cancellationToken) { await using (var probe = connection.CreateCommand()) { @@ -177,7 +190,7 @@ private static async Task IsSchemaCurrentAsync( var exists = (long)(await probe.ExecuteScalarAsync(cancellationToken).ConfigureAwait(false))!; if (exists == 0) { - return false; + return 0; } } @@ -188,7 +201,7 @@ private static async Task IsSchemaCurrentAsync( // connection. It decides whether the scripts still need running, which is boot work, not store // work. var version = await command.ExecuteScalarAsync(cancellationToken).ConfigureAwait(false); - return version is long deployed && deployed >= ExpectedSchemaVersion; + return version is long deployed ? deployed : 0; } /// diff --git a/src/BackWave/Monitor/BackWaveMonitor.cs b/src/BackWave/Monitor/BackWaveMonitor.cs index a6c079d..05e033d 100644 --- a/src/BackWave/Monitor/BackWaveMonitor.cs +++ b/src/BackWave/Monitor/BackWaveMonitor.cs @@ -402,6 +402,7 @@ public ValueTask> ListObserverDeadLetter CancelRequested = record.CancelRequested, TerminalAt = record.TerminalAt, TerminalCause = record.TerminalCause, + RetryCause = record.RetryCause, ScheduleId = record.ScheduleId, Sequence = record.Sequence, Tags = record.Tags, diff --git a/src/BackWave/Monitor/JobSnapshot.cs b/src/BackWave/Monitor/JobSnapshot.cs index 2aba8d9..20a40f1 100644 --- a/src/BackWave/Monitor/JobSnapshot.cs +++ b/src/BackWave/Monitor/JobSnapshot.cs @@ -41,6 +41,14 @@ public sealed record JobSnapshot /// A short reason for the terminal outcome (for example why it was dead-lettered); null while still active. public string? TerminalCause { get; init; } + /// + /// Why the job most recently went back to Scheduled after an attempt went wrong: the handler failed, + /// or the lease expired. Null when no attempt has gone wrong since the job was enqueued or last + /// requeued. A job handed back by a worker on a clean stop (for example during a deploy) keeps + /// whatever it had, so a deploy alone never sets it. A Scheduled job with a cause is Retrying. + /// + public RetryCause? RetryCause { get; init; } + /// The recurring schedule that minted this instance; null for a directly enqueued job. public string? ScheduleId { get; init; } diff --git a/src/BackWave/Storage/IJobStore.cs b/src/BackWave/Storage/IJobStore.cs index 54227ff..222c4af 100644 --- a/src/BackWave/Storage/IJobStore.cs +++ b/src/BackWave/Storage/IJobStore.cs @@ -814,6 +814,15 @@ public sealed record JobQuery /// Match only jobs minted by this recurring schedule; null matches jobs from any source. public string? ScheduleId { get; init; } + /// + /// When true, match only Retrying jobs: jobs that are Scheduled and carry a + /// , so they are waiting for another attempt because the handler + /// failed or the lease expired. A new job and a requeued job do not match, and a clean-stop hand-back + /// of the lease does not make a job match. False (the default) adds no constraint. Like every filter it is AND-ed with the + /// others, so a other than Scheduled together with this flag matches nothing. + /// + public bool Retrying { get; init; } + /// /// Tag predicates AND-ed together and AND-composed with the scalar filters above: a job matches /// only when it satisfies EVERY predicate. An empty list adds no constraint (matches everything). diff --git a/src/BackWave/Storage/InMemory/InMemoryJobStore.cs b/src/BackWave/Storage/InMemory/InMemoryJobStore.cs index 05c898b..0e3bb63 100644 --- a/src/BackWave/Storage/InMemory/InMemoryJobStore.cs +++ b/src/BackWave/Storage/InMemory/InMemoryJobStore.cs @@ -783,6 +783,7 @@ public ValueTask ReportOutcomeAsync( DueTime = retryAt, LeaseOwner = null, LeaseExpiry = null, + RetryCause = RetryCause.HandlerFailed, }, JobOutcome.Failure failure => job with { @@ -939,6 +940,7 @@ public ValueTask ExpireLeasesAsync( DueTime = dueTime, LeaseOwner = null, LeaseExpiry = null, + RetryCause = RetryCause.LeaseExpired, } : job with { @@ -1069,6 +1071,7 @@ public ValueTask RequeueAsync( CancelRequested = false, TerminalAt = null, TerminalCause = null, + RetryCause = null, }; RecordTransition(jobId, JobState.Scheduled, 0, now); // Attempt budget reset (§3) AppendAudit(actor, OperatorAction.Requeue, jobId.ToString(), now); @@ -1731,6 +1734,7 @@ private static bool MatchesScope(JobRecord j, JobQuery query) && (query.Queue is null || j.Queue == query.Queue) && (query.WireName is null || j.WireName == query.WireName) && (query.ScheduleId is null || j.ScheduleId == query.ScheduleId) + && (!query.Retrying || (j.State == JobState.Scheduled && j.RetryCause is not null)) // Tag predicates are AND-ed (ADR 0022): a job must satisfy EVERY predicate. // An empty list adds no constraint (All over empty is true). OR is out of scope. && query.TagPredicates.All(p => p.Matches(j.Tags)); diff --git a/src/BackWave/Storage/JobRecord.cs b/src/BackWave/Storage/JobRecord.cs index 8a8de88..9c1ddc3 100644 --- a/src/BackWave/Storage/JobRecord.cs +++ b/src/BackWave/Storage/JobRecord.cs @@ -42,6 +42,14 @@ public sealed record JobRecord /// A short human-readable reason for the terminal state (the failure error, cancel actor, or unroutable reason), or null while live. public string? TerminalCause { get; init; } + /// + /// Why the job most recently went back to Scheduled after an attempt went wrong: the handler failed, + /// or the lease expired. Null when no attempt has gone wrong since the job was enqueued or last + /// requeued. A clean-stop hand-back of the lease does not change it, and a terminal outcome keeps it + /// as a record of the last retry. A Scheduled job with a cause is Retrying. + /// + public RetryCause? RetryCause { get; init; } + /// The id of the recurring schedule that minted this instance, or null for a directly enqueued job. public string? ScheduleId { get; init; } diff --git a/src/BackWave/Storage/RetryCause.cs b/src/BackWave/Storage/RetryCause.cs new file mode 100644 index 0000000..301fbb2 --- /dev/null +++ b/src/BackWave/Storage/RetryCause.cs @@ -0,0 +1,20 @@ +namespace BackWave.Storage; + +/// +/// Why an attempt went wrong and sent its job back to Scheduled for another attempt. A job that carries +/// a cause and is Scheduled is Retrying: it is waiting to run again because something failed, not because +/// it is new or because a worker handed it back on a clean stop. +/// +/// Members are a stable wire identity (persisted by number): every adapter writes (int) of this +/// enum into a nullable int column, and no cause is stored as null, never as a number. The assigned +/// value, not the member's position, is the storage contract - give a new member the next free number and +/// never reuse a retired one. +/// +public enum RetryCause +{ + /// The handler failed (it threw or reported a failure) and the retry policy scheduled another attempt. + HandlerFailed = 1, + + /// The lease lapsed before the worker reported an outcome (for example the worker crashed or stalled), and the store scheduled another attempt. + LeaseExpired = 2, +} diff --git a/tests/BackWave.Dashboard.Tests/DashboardTests.cs b/tests/BackWave.Dashboard.Tests/DashboardTests.cs index 1490b66..7f22d6b 100644 --- a/tests/BackWave.Dashboard.Tests/DashboardTests.cs +++ b/tests/BackWave.Dashboard.Tests/DashboardTests.cs @@ -304,6 +304,31 @@ public async Task LiveView_StreamsTheRenderedFragmentOverSse_WhenAskedWithLiveFl } } + [Fact] + public async Task LiveView_EndsTheSseStream_WhenTheApplicationStartsToStop() + { + // The server waits for open requests before it stops the hosted services. A stream that + // outlives ApplicationStopping makes one open dashboard tab spend the whole shutdown window, + // and the worker groups then cannot give their leases back on a clean stop. + var (app, _, http) = await StartAsync(new BackWaveDashboardOptions { LiveRefreshInterval = TimeSpan.FromMilliseconds(50) }); + await using (app) + { + using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(15)); + using var response = await http.GetAsync( + "/backwave/?live=1", HttpCompletionOption.ResponseHeadersRead, cts.Token); + await using var stream = await response.Content.ReadAsStreamAsync(cts.Token); + using var reader = new StreamReader(stream); + Assert.Equal("event: update", await reader.ReadLineAsync(cts.Token)); + + app.Lifetime.StopApplication(); + + // Without the stop link the stream pings every interval until the test times out. + while (await reader.ReadLineAsync(cts.Token) is not null) + { + } + } + } + [Fact] public async Task LiveView_ClosesTheSseStreamOnNavigation_SoItNeverStrandsAConnection() { @@ -474,6 +499,98 @@ public async Task Failures_ShowDeadLetteredAndQuarantined_Separately() } } + /// + /// Seeds one job per way a job can be Scheduled: a handler failure the policy retries, a lapsed + /// lease the sweep reschedules, a clean-stop hand-back, and a fresh enqueue. Each lives on its own + /// Queue so each claim takes only its own job. + /// + private static async Task SeedRetryCasesAsync(InMemoryJobStore store) + { + var disposition = new RetryPolicy { MaxAttempts = 5, Backoff = _ => TimeSpan.FromMinutes(1) }.ToDisposition(); + + await store.EnqueueAsync(Job(wireName: "flaky-charge", queue: "q-failed"), now: T0); + var failed = Assert.Single(await store.ClaimAsync(new ClaimRequest("w1", ["q-failed"], 32, Lease, T0))); + await store.ReportOutcomeAsync(failed.JobId, "w1", failed.Attempt, new JobOutcome.Failure(T0.AddMinutes(5), "card declined"), T0); + + await store.EnqueueAsync(Job(wireName: "stalled-export", queue: "q-expired"), now: T0); + Assert.Single(await store.ClaimAsync(new ClaimRequest("w2", ["q-expired"], 32, Lease, T0))); + await store.ExpireLeasesAsync(T0 + Lease + TimeSpan.FromSeconds(1), 32, ["q-expired"], disposition); + + await store.EnqueueAsync(Job(wireName: "handed-back", queue: "q-relinquished"), now: T0); + Assert.Single(await store.ClaimAsync(new ClaimRequest("w3", ["q-relinquished"], 32, Lease, T0))); + await store.RelinquishLeasesAsync("w3", T0.AddSeconds(5), disposition); + + await store.EnqueueAsync(Job(wireName: "brand-new", queue: "q-fresh"), now: T0); + } + + [Fact] + public async Task Failures_RetryingTab_ListsOnlyJobsAnAttemptWentWrongFor() + { + var (app, store, http) = await StartAsync(); + await using (app) + { + await SeedRetryCasesAsync(store); + + // The default tab still opens on Dead-Lettered, with Retrying as a counted third tab. + var html = await http.GetStringAsync("/backwave/failures"); + Assert.Contains("/backwave/failures?tab=retrying", html); + Assert.DoesNotContain("flaky-charge", html); + + var retrying = await http.GetStringAsync("/backwave/failures?tab=retrying"); + Assert.Contains("flaky-charge", retrying); + Assert.Contains("stalled-export", retrying); + Assert.DoesNotContain("handed-back", retrying); // a clean stop is not a problem + Assert.DoesNotContain("brand-new", retrying); + // The row says when the next attempt runs and why the job came back. + Assert.Contains("data-label=\"Next Attempt\"", retrying); + Assert.Contains("data-label=\"Retry Cause\"", retrying); + Assert.Contains("Handler failed", retrying); + Assert.Contains("Lease expired", retrying); + Assert.Contains("2", retrying); + } + } + + [Fact] + public async Task JobSearch_RetryingFilter_NarrowsScheduledToTheRetryingJobs() + { + var (app, store, http) = await StartAsync(); + await using (app) + { + await SeedRetryCasesAsync(store); + + var html = await http.GetStringAsync("/backwave/jobs?state=Retrying"); + Assert.Contains("""""", html); + Assert.Contains("flaky-charge", html); + Assert.Contains("stalled-export", html); + Assert.DoesNotContain("handed-back", html); + Assert.DoesNotContain("brand-new", html); + Assert.Contains("data-label=\"Retry Cause\"", html); + + // Plain Scheduled still lists every Scheduled job, problem or not. + var scheduled = await http.GetStringAsync("/backwave/jobs?state=Scheduled"); + Assert.Contains("handed-back", scheduled); + Assert.Contains("brand-new", scheduled); + Assert.Contains("flaky-charge", scheduled); + } + } + + [Fact] + public async Task JobDetail_ShowsTheRetryCause_OfARetryingJob() + { + var (app, store, http) = await StartAsync(); + await using (app) + { + await SeedRetryCasesAsync(store); + var flaky = Assert.Single(await store.ListJobsAsync(new JobQuery { Queue = "q-failed" })); + var fresh = Assert.Single(await store.ListJobsAsync(new JobQuery { Queue = "q-fresh" })); + + var html = await http.GetStringAsync($"/backwave/jobs/{flaky.JobId}"); + Assert.Contains("Retry causeHandler failed", html); + + Assert.DoesNotContain("Retry cause", await http.GetStringAsync($"/backwave/jobs/{fresh.JobId}")); + } + } + [Fact] public async Task Failures_RenderTags_DeepLinkingToTheJobsListFilteredByStateAndTag() { diff --git a/tests/BackWave.Oracle.Tests/OracleRoundTripBudgetTests.cs b/tests/BackWave.Oracle.Tests/OracleRoundTripBudgetTests.cs index e8edf37..3af9860 100644 --- a/tests/BackWave.Oracle.Tests/OracleRoundTripBudgetTests.cs +++ b/tests/BackWave.Oracle.Tests/OracleRoundTripBudgetTests.cs @@ -43,7 +43,9 @@ private static NewJob Job(string queue = "budget") => new(Guid.NewGuid(), "budget-test", "{}"u8.ToArray(), queue, T0); // The recorded budgets: measured 2026-08-22 against Oracle Free 23 on ODP.NET 23.9.1, at schema - // version 1. Each is the cost of ONE call; the arithmetic behind each number is in its test. + // version 1. The fetch window was measured again 2026-10-06 at schema version 2, where the jobs row + // gained the 22-byte retry cause column. Each is the cost of ONE call; the arithmetic behind each + // number is in its test. private static readonly Budget Claim = new( "ClaimBatchAsync of 32 jobs (one queue, cold caches)", @@ -74,12 +76,12 @@ private static NewJob Job(string queue = "budget") => Statements: 3, LobReads: 0, FetchWindowBytes: 0); // The window a statement selecting the full jobs column set declares: 32 rows (one claim batch) of - // 140,447 bytes, which is the driver's own size for that row - both LOB columns at the 65,536 + // 140,469 bytes, which is the driver's own size for that row - both LOB columns at the 65,536 // payload prefetch, plus about 9 KB of scalars. Claim and job list select the same columns, so they // share it. The page size does NOT enter it: a 200-row page arrives in seven windows of this size // rather than one window seven times as wide, which is what keeps the monitor listing off the // memory ceiling. - private const long JobPageWindow = 4_494_304; + private const long JobPageWindow = 4_495_008; [Fact] public async Task Claim_of_a_full_batch_stays_within_its_round_trip_budget() diff --git a/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs b/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs index 0e80f74..f5de2f4 100644 --- a/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs +++ b/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs @@ -51,7 +51,7 @@ JsonElement OutputProperties(string name) => // The input contract is snake_case (the fixed tool shapes). var searchInputs = tools.Single(t => t.Name == "search_jobs").InputSchema!.Value.GetProperty("properties"); foreach (var parameter in new[] - { "state", "queue", "wire_name", "schedule_id", "tags", "after_cursor", "sort", "max_results" }) + { "state", "retrying", "queue", "wire_name", "schedule_id", "tags", "after_cursor", "sort", "max_results" }) { Assert.True(searchInputs.TryGetProperty(parameter, out _), $"search_jobs is missing input '{parameter}'"); } @@ -398,6 +398,34 @@ public async Task SearchJobs_InvalidStateAndSort_AreInvalidInputErrors() Assert.Contains("oldest_first", badSort.Text); } + [Fact] + public async Task SearchJobs_Retrying_ListsOnlyJobsAnAttemptWentWrongFor_WithTheirCause() + { + await using var server = await McpTestServer.StartAsync(); + var failing = await server.SeedJobAsync("critical"); + var claimed = Assert.Single(await server.Store.ClaimAsync( + new ClaimRequest("w1", ["critical"], 32, TimeSpan.FromMinutes(1), DateTimeOffset.UtcNow))); + await server.Store.ReportOutcomeAsync( + claimed.JobId, "w1", claimed.Attempt, new JobOutcome.Failure(DateTimeOffset.UtcNow.AddMinutes(5), "boom"), DateTimeOffset.UtcNow); + var fresh = await server.SeedJobAsync("critical"); + + var jobs = (await server.Client.CallToolAsync("search_jobs", new Dictionary + { + ["retrying"] = true, + })).StructuredContent!.Value.GetProperty("jobs").EnumerateArray().ToList(); + + var job = Assert.Single(jobs); + Assert.Equal(failing, job.GetProperty("jobId").GetGuid()); + Assert.Equal("Scheduled", job.GetProperty("state").GetString()); + Assert.Equal("HandlerFailed", job.GetProperty("retryCause").GetString()); + + // Without the filter the fresh job lists too, and carries no cause. + var all = (await server.Client.CallToolAsync("search_jobs")).StructuredContent!.Value + .GetProperty("jobs").EnumerateArray().ToList(); + var freshRow = Assert.Single(all, j => j.GetProperty("jobId").GetGuid() == fresh); + Assert.True(!freshRow.TryGetProperty("retryCause", out var cause) || cause.ValueKind == JsonValueKind.Null); + } + [Fact] public async Task SearchJobs_OldestFirst_ReversesTheOrder() { diff --git a/tests/BackWave.SchemaGate.Tests/SchemaGateTests.cs b/tests/BackWave.SchemaGate.Tests/SchemaGateTests.cs index afb595f..fef5479 100644 --- a/tests/BackWave.SchemaGate.Tests/SchemaGateTests.cs +++ b/tests/BackWave.SchemaGate.Tests/SchemaGateTests.cs @@ -14,8 +14,8 @@ namespace BackWave.SchemaGate.Tests; public sealed class SchemaGateTests { // Each adapter's assembly, reached through a type it ships, so the gate reads the SAME embedded - // scripts the migrator runs. SQLite is here too: its consolidated v1 script plus the v1 -> v2 - // step that adds the transition-position high-water mark, inspected with zero extra wiring. + // scripts the migrator runs. SQLite is here too: its consolidated v1 script plus its incremental + // steps, inspected with zero extra wiring. public static TheoryData Adapters() => new() { { "Postgres", typeof(PostgresMigrator).Assembly }, @@ -46,21 +46,26 @@ public void EveryShippedMigrationIsAdditive(string adapter, Assembly adapterAsse } [Fact] - public void Sqlite_ShipsTheTransitionPositionStepAsItsFirstIncrementalMigration() + public void Sqlite_ShipsOneIncrementalStepPerVersion_InVersionOrder() { - // SQLite's first real vN-1 -> vN step since its schema was consolidated into v1: 0002 adds the - // transition-position high-water mark. The script count is the version the adapter requires - // and the step stamps that same version, so the two cannot drift apart unnoticed; the gate - // above polices the step's DDL like any other adapter's. + // SQLite's real vN-1 -> vN steps since its schema was consolidated into v1: 0002 adds the + // transition-position high-water mark and 0003 adds the retry cause. The script count is the + // version the adapter requires and the last step stamps that same version, so the two cannot + // drift apart unnoticed; the gate above polices each step's DDL like any other adapter's. var scripts = SchemaScripts.Load(typeof(SqliteMigrator).Assembly); Assert.Equal(SqliteMigrator.ExpectedSchemaVersion, scripts.Count); - var step = scripts[^1]; - Assert.EndsWith("0002_transition_position.sql", step.ResourceName, StringComparison.Ordinal); - Assert.Contains("CREATE TABLE IF NOT EXISTS backwave_transition_position", step.Sql, StringComparison.Ordinal); + var transitionPosition = scripts[1]; + Assert.EndsWith("0002_transition_position.sql", transitionPosition.ResourceName, StringComparison.Ordinal); + Assert.Contains( + "CREATE TABLE IF NOT EXISTS backwave_transition_position", transitionPosition.Sql, StringComparison.Ordinal); + + var retryCause = scripts[^1]; + Assert.EndsWith("0003_retry_cause.sql", retryCause.ResourceName, StringComparison.Ordinal); + Assert.Contains("ADD COLUMN retry_cause INTEGER NULL", retryCause.Sql, StringComparison.Ordinal); Assert.Contains( $"UPDATE backwave_schema_version SET version = {SqliteMigrator.ExpectedSchemaVersion};", - step.Sql, StringComparison.Ordinal); + retryCause.Sql, StringComparison.Ordinal); } // ---- Sabotage self-tests: prove the gate turns RED on a synthetic non-additive migration. ---- diff --git a/tests/BackWave.SqlServer.Tests/SqlServerCoordinatedMigrationTests.cs b/tests/BackWave.SqlServer.Tests/SqlServerCoordinatedMigrationTests.cs index 1ccc838..294a43f 100644 --- a/tests/BackWave.SqlServer.Tests/SqlServerCoordinatedMigrationTests.cs +++ b/tests/BackWave.SqlServer.Tests/SqlServerCoordinatedMigrationTests.cs @@ -182,7 +182,7 @@ public async Task ConcurrentFirstBootAgainstFreshDatabase_EnablesRcsiAndMigrates // so a v2 script that inserted again would leave the row at 1 and every node would fail-stop on // skew. Pinning both the row count and the version is what catches that. [Fact] - public async Task AV1Database_UpgradesInPlaceToV2() + public async Task AV1Database_UpgradesInPlaceToTheCurrentVersion() { await DropSchemaAsync(); await ApplyScriptAsync("0001_initial.sql"); @@ -191,7 +191,6 @@ public async Task AV1Database_UpgradesInPlaceToV2() await SqlServerMigrator.MigrateAsync(SqlServerTestDatabase.ConnectionString, Schema); Assert.Equal(1, await SchemaVersionRowCountAsync()); - Assert.Equal(2, await DeployedVersionAsync()); Assert.Equal(SqlServerMigrator.ExpectedSchemaVersion, await DeployedVersionAsync()); Assert.Equal("lease_owner", await LeaseOwnerIndexKeyColumnAsync()); } diff --git a/tests/BackWave.Tests/CapturingLogger.cs b/tests/BackWave.Tests/CapturingLogger.cs index f023a91..ada9987 100644 --- a/tests/BackWave.Tests/CapturingLogger.cs +++ b/tests/BackWave.Tests/CapturingLogger.cs @@ -12,22 +12,39 @@ internal sealed record LogRecord( internal sealed class LogCapture { - public List Records { get; } = []; + private readonly List _records = []; + + // A handler the pump abandoned unwinds on a pool thread and can log while the test reads, so a read + // takes a copy under the same lock the logger writes under. + public IReadOnlyList Records + { + get + { + lock (_records) + { + return [.. _records]; + } + } + } + + public void Add(LogRecord record) + { + lock (_records) + { + _records.Add(record); + } + } public bool Enabled { get; set; } = true; } internal sealed class CapturingLogger(LogCapture capture) : ILogger { - // The pump drives one execution at a time and scopes open/close in order, so a simple stack is enough - // for these single-job tests. - private readonly List _scopes = []; + // Scopes follow the async flow that opened them, as in a real logging provider, so a handler that + // unwinds on a pool thread neither sees nor closes the scopes of the code that drives the pump. + private readonly LoggerExternalScopeProvider _scopes = new(); - public IDisposable BeginScope(TState state) where TState : notnull - { - _scopes.Add(state); - return new Pop(_scopes); - } + public IDisposable? BeginScope(TState state) where TState : notnull => _scopes.Push(state); public bool IsEnabled(LogLevel logLevel) => capture.Enabled; @@ -36,19 +53,16 @@ public void Log( Func formatter) { var scope = new List>(); - foreach (var open in _scopes) - { - if (open is IEnumerable> pairs) + _scopes.ForEachScope( + (open, collected) => { - scope.AddRange(pairs); - } - } - capture.Records.Add(new LogRecord(logLevel, eventId.Id, formatter(state, exception), scope)); - } - - private sealed class Pop(List scopes) : IDisposable - { - public void Dispose() => scopes.RemoveAt(scopes.Count - 1); + if (open is IEnumerable> pairs) + { + collected.AddRange(pairs); + } + }, + scope); + capture.Add(new LogRecord(logLevel, eventId.Id, formatter(state, exception), scope)); } } diff --git a/tests/BackWave.Tests/JobStateWireFormatTests.cs b/tests/BackWave.Tests/JobStateWireFormatTests.cs index 0aa9378..d00ea32 100644 --- a/tests/BackWave.Tests/JobStateWireFormatTests.cs +++ b/tests/BackWave.Tests/JobStateWireFormatTests.cs @@ -34,13 +34,14 @@ public class JobStateWireFormatTests ["ix_backwave_jobs_claim"] = JobState.Scheduled, ["ix_backwave_jobs_leased_queue"] = JobState.Leased, ["ix_backwave_jobs_lease_owner"] = JobState.Leased, + ["ix_backwave_jobs_retrying"] = JobState.Scheduled, }; - // Postgres, SQL Server, and SQLite each carry the claim and leased-queue predicates, SQL Server carries - // the lease-owner one as well, and Oracle carries none, because it has no partial index. Pinned so that + // Postgres, SQL Server, and SQLite each carry the claim, leased-queue, and retrying predicates, SQL Server + // carries the lease-owner one as well, and Oracle carries none, because it has no partial index. Pinned so that // dropping a predicate, or adding an adapter that needs one, is a deliberate edit here rather than a // silent loss of coverage. - private const int GuardedPredicateCount = 7; + private const int GuardedPredicateCount = 10; // The `-- States: 0 Scheduled, ...` gloss each schema carries above its jobs table, which is the one // comment that has to spell the numbers out: it is the only documentation a DBA reading the canonical diff --git a/tests/BackWave.Tests/RetryCauseWireFormatTests.cs b/tests/BackWave.Tests/RetryCauseWireFormatTests.cs new file mode 100644 index 0000000..2a0518c --- /dev/null +++ b/tests/BackWave.Tests/RetryCauseWireFormatTests.cs @@ -0,0 +1,37 @@ +using BackWave.Storage; + +namespace BackWave.Tests; + +// RetryCause's numbers are a storage wire format: every adapter writes (int)cause into a nullable int +// column. The conformance suite and the upgrade harness write and read with the same code, so both sides +// agree on a wrong number. This test is the asymmetric side: it holds the numbers that rows already in +// customer databases were written with. + +public class RetryCauseWireFormatTests +{ + // An entry only ever changes alongside a migration that rewrites the rows. + private static readonly Dictionary PersistedValues = new(StringComparer.Ordinal) + { + [nameof(RetryCause.HandlerFailed)] = 1, + [nameof(RetryCause.LeaseExpired)] = 2, + }; + + [Fact] + public void EveryMember_IsPinned_AndKeepsThePersistedNumberItsRowsWereWrittenWith() + { + var actual = Enum.GetValues().ToDictionary(cause => cause.ToString(), cause => (int)cause, StringComparer.Ordinal); + + Assert.True( + actual.Count == PersistedValues.Count && actual.All(pair => PersistedValues.TryGetValue(pair.Key, out var pinned) && pinned == pair.Value), + $""" + RetryCause no longer matches its pinned wire numbers. + + Now: {string.Join(", ", actual.Select(pair => $"{pair.Key}={pair.Value}"))} + Pinned: {string.Join(", ", PersistedValues.Select(pair => $"{pair.Key}={pair.Value}"))} + + A RetryCause number is persisted, so renumbering migrates nothing - every row already in a + customer database reads back as a different cause. To add a cause, give it the next free number + and pin it here; to retire one, leave its number reserved rather than reusing it. + """); + } +} diff --git a/tests/BackWave.Upgrade.Tests/UpgradeHarnessTests.cs b/tests/BackWave.Upgrade.Tests/UpgradeHarnessTests.cs index f43c1f3..b630546 100644 --- a/tests/BackWave.Upgrade.Tests/UpgradeHarnessTests.cs +++ b/tests/BackWave.Upgrade.Tests/UpgradeHarnessTests.cs @@ -16,10 +16,9 @@ namespace BackWave.Upgrade.Tests; public sealed class UpgradeHarnessTests { // Short workload per prior version keeps the shipped-prior-version sweep battery-friendly while still - // running a real concurrent workload across the freshly migrated schema. SQL Server ships v2, so its - // sweep carries one real step (v1 -> v2) that populates, migrates, works and audits. Postgres is still - // at the re-baselined consolidated v1, so its sweep (v1..v(current-1)) is legitimately empty and that - // clean fact passes vacuously; the sabotage fact below still exercises the oracle end to end. + // running a real concurrent workload across the freshly migrated schema. Each prior version in the + // sweep (v1..v(current-1)) populates, migrates to current, works and audits: SQL Server carries v1 and + // v2, Postgres carries v1. The sabotage fact below proves the oracle turns red on a broken upgrade. private static readonly TimeSpan BatteryWorkload = TimeSpan.FromSeconds(3); [Fact] @@ -47,8 +46,7 @@ public async Task SqlServer_EveryShippedPriorVersion_UpgradesInPlaceCleanly() [Fact] public async Task Sabotage_LosingAPopulatedJobDuringMigration_TurnsTheHarnessRed() { - // Hand-break the migration on the consolidated v1 (the only shipped version, so the empty sweep - // cannot exercise the oracle on its own): populate the base v1 fixture inventory, run the real + // Hand-break the migration from the consolidated v1: populate the base v1 fixture inventory, run the real // idempotent migrate-to-current, delete a populated fixture job, and prove the conservation oracle // goes RED. Proves the harness has teeth — a broken upgrade cannot pass green. var exit = await UpgradeRun.RunAsync(new UpgradeOptions