From 7390084b8733cc63765d0b90ca77f8b6b412df1e Mon Sep 17 00:00:00 2001 From: Paul DeVito Date: Tue, 6 Oct 2026 02:29:38 +0300 Subject: [PATCH 1/3] fix(dashboard): end the live SSE stream when the application stops The server waits for open requests before it stops the hosted services. An open dashboard tab held its live stream until the host shutdown timeout ended. The worker groups then had no time to give their leases back on a clean stop, and the leases expired later. The stream now also ends on ApplicationStopping. --- .../DashboardRequestHandler.cs | 13 +++++++--- .../DashboardTests.cs | 25 +++++++++++++++++++ 2 files changed, 35 insertions(+), 3 deletions(-) diff --git a/src/BackWave.Dashboard/DashboardRequestHandler.cs b/src/BackWave.Dashboard/DashboardRequestHandler.cs index 1e79557..2281919 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. } } diff --git a/tests/BackWave.Dashboard.Tests/DashboardTests.cs b/tests/BackWave.Dashboard.Tests/DashboardTests.cs index 1490b66..31e637e 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() { From 5e0ce55af0379a3f97ca1cb3f13b02602a2dca11 Mon Sep 17 00:00:00 2001 From: Paul DeVito Date: Tue, 6 Oct 2026 02:29:44 +0300 Subject: [PATCH 2/3] test: make the capturing logger safe for concurrent writers An abandoned handler can unwind on a thread-pool thread and log while the test reads the records. The read then failed with "Collection was modified". The scope stack was also shared across threads. The capture now locks its records and returns a copy. Scopes use LoggerExternalScopeProvider, which keeps one scope stack per async flow. LeaseReclaimLogTests failed 3 times in 25 runs before and 0 times in 40 runs after. --- tests/BackWave.Tests/CapturingLogger.cs | 56 +++++++++++++++---------- 1 file changed, 35 insertions(+), 21 deletions(-) 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)); } } From c561f652595799ac05c6d8186004cf49979a111f Mon Sep 17 00:00:00 2001 From: Paul DeVito Date: Tue, 6 Oct 2026 02:30:07 +0300 Subject: [PATCH 3/3] feat: record why a job went back to Scheduled and show Retrying jobs A Scheduled job did not tell an operator if an attempt went wrong. A job with Attempt > 0 is not a good signal, because a clean-stop relinquish also gives back a claimed job. A deploy then looks like a wave of problems. Each job now records a retry cause: - HandlerFailed: the handler failed and the retry policy rescheduled the job (single and batched outcome reports). - LeaseExpired: the lease expired and the sweep rescheduled the job. A requeue clears the cause. A claim, a relinquish, a terminal outcome, a cancel and a parent latch keep it. A new job has no cause. Retrying means Scheduled with a retry cause. It is available as: - JobQuery.Retrying and JobSnapshot.RetryCause on the Monitor API. - A Retrying tab on the Failures page, and a Retrying option in the Jobs state filter. Both show the next attempt and the cause. - A Retry cause row on the job detail page. - A "retrying" filter on the MCP search_jobs tool. Each SQL adapter gets one additive migration: a nullable retry_cause column and an index for the Retrying query. Postgres and Oracle go to version 2, SQL Server and SQLite to version 3. There is no backfill. Jobs that are Scheduled before the upgrade have no cause, so they do not show as Retrying. The Oracle migrator now sets ddl_lock_timeout on its own unpooled session. Without it, the new ALTER TABLE failed with ORA-00054 when other nodes ran the first script at the same time on a cold boot. --- CONTEXT.md | 8 + src/BackWave.Conformance/ConformanceSuite.cs | 190 ++++++++++++++++++ .../Components/JobDetailPanel.razor | 4 + .../Components/JobTable.razor | 30 ++- .../Components/Pages/Failures.razor | 57 ++++-- .../Components/Pages/Jobs.razor | 18 +- src/BackWave.Dashboard/DashboardGlossary.cs | 14 ++ .../DashboardRequestHandler.cs | 22 +- src/BackWave.Oracle/OracleJobStore.cs | 35 +++- src/BackWave.Oracle/OracleMigrator.cs | 15 +- .../Schema/0002_retry_cause.sql | 64 ++++++ src/BackWave.Postgres/PostgresJobStore.cs | 43 +++- src/BackWave.Postgres/PostgresMigrator.cs | 2 +- .../Schema/0002_retry_cause.sql | 21 ++ src/BackWave.Pro.Mcp/Tools/JobTools.cs | 8 + .../Schema/0003_retry_cause.sql | 25 +++ src/BackWave.SqlServer/SqlServerJobStore.cs | 40 +++- src/BackWave.SqlServer/SqlServerMigrator.cs | 2 +- .../Schema/0003_retry_cause.sql | 23 +++ src/BackWave.Sqlite/SqliteJobStore.cs | 14 +- src/BackWave.Sqlite/SqliteMigrator.cs | 23 ++- src/BackWave/Monitor/BackWaveMonitor.cs | 1 + src/BackWave/Monitor/JobSnapshot.cs | 8 + src/BackWave/Storage/IJobStore.cs | 9 + .../Storage/InMemory/InMemoryJobStore.cs | 4 + src/BackWave/Storage/JobRecord.cs | 8 + src/BackWave/Storage/RetryCause.cs | 20 ++ .../DashboardTests.cs | 92 +++++++++ .../OracleRoundTripBudgetTests.cs | 8 +- tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs | 30 ++- .../SchemaGateTests.cs | 27 ++- .../SqlServerCoordinatedMigrationTests.cs | 3 +- .../BackWave.Tests/JobStateWireFormatTests.cs | 7 +- .../RetryCauseWireFormatTests.cs | 37 ++++ .../UpgradeHarnessTests.cs | 10 +- 35 files changed, 838 insertions(+), 84 deletions(-) create mode 100644 src/BackWave.Oracle/Schema/0002_retry_cause.sql create mode 100644 src/BackWave.Postgres/Schema/0002_retry_cause.sql create mode 100644 src/BackWave.SqlServer/Schema/0003_retry_cause.sql create mode 100644 src/BackWave.Sqlite/Schema/0003_retry_cause.sql create mode 100644 src/BackWave/Storage/RetryCause.cs create mode 100644 tests/BackWave.Tests/RetryCauseWireFormatTests.cs 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 0750340..2000fa3 100644 --- a/src/BackWave.Conformance/ConformanceSuite.cs +++ b/src/BackWave.Conformance/ConformanceSuite.cs @@ -1795,6 +1795,196 @@ public async Task Clause_5_9_Reads_FilterAndCount_OverCommittedEffectsOnly() Assert.Contains(new QueueStateCount("other", JobState.Scheduled, 1), counts); } + // ── §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 2281919..f732e31 100644 --- a/src/BackWave.Dashboard/DashboardRequestHandler.cs +++ b/src/BackWave.Dashboard/DashboardRequestHandler.cs @@ -545,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]) @@ -577,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 @@ -716,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 995c299..f5afd73 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); @@ -2615,6 +2617,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]; @@ -3841,7 +3847,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) { @@ -3879,9 +3885,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 46ac03e..0db3717 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 """, @@ -2203,6 +2205,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]; @@ -3477,7 +3483,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) { @@ -3515,6 +3521,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), }; } @@ -3532,6 +3539,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 dceea9f..45e1751 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.")] @@ -93,6 +95,7 @@ _ when sort.Equals("oldest_first", StringComparison.OrdinalIgnoreCase) => JobSor var query = new JobQuery { State = parsedState, + Retrying = retrying ?? false, Queue = queue, WireName = wire_name, ScheduleId = schedule_id, @@ -413,6 +416,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; } @@ -442,6 +449,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 0fd9d4b..a576c5f 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) """, @@ -2381,6 +2384,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]; @@ -3617,7 +3624,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) { @@ -3655,9 +3662,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 c8a207e..f77ee7f 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 """, @@ -2068,6 +2069,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]; @@ -3251,7 +3256,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() { @@ -3273,6 +3278,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 98fa868..3fbff67 100644 --- a/src/BackWave/Monitor/BackWaveMonitor.cs +++ b/src/BackWave/Monitor/BackWaveMonitor.cs @@ -383,6 +383,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 0edebfa..603530f 100644 --- a/src/BackWave/Storage/IJobStore.cs +++ b/src/BackWave/Storage/IJobStore.cs @@ -781,6 +781,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 d70ef88..fc3e2c7 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); @@ -1722,6 +1725,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 31e637e..7f22d6b 100644 --- a/tests/BackWave.Dashboard.Tests/DashboardTests.cs +++ b/tests/BackWave.Dashboard.Tests/DashboardTests.cs @@ -499,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 1b34afa..e5e07d0 100644 --- a/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs +++ b/tests/BackWave.Pro.Mcp.Tests/JobToolsTests.cs @@ -49,7 +49,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}'"); } @@ -390,6 +390,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/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