diff --git a/src/BackWave.SqlServer/SqlServerJobStore.cs b/src/BackWave.SqlServer/SqlServerJobStore.cs index f120d27..efb006b 100644 --- a/src/BackWave.SqlServer/SqlServerJobStore.cs +++ b/src/BackWave.SqlServer/SqlServerJobStore.cs @@ -467,11 +467,10 @@ public async ValueTask ClaimBatchAsync( // Tags hydrate in one batched round-trip (ADR 0022) — but only when tags are actually in // use (issue 0169). T-SQL's OUTPUT forbids the correlated subquery the Postgres RETURNING - // uses to fold tags into the claim, and capturing the claim into a table variable to read - // them in-transaction widens the lease's row locks enough to deadlock the transition write, - // so SQL Server instead GATES the existing post-commit hydration: under the no-tags - // configuration the job_tags table is empty, the gate skips the round-trip entirely, and the - // claim hot path pays nothing. See TagsInUseAsync for the cheap, sound presence signal. + // uses to fold tags into the claim, so SQL Server instead GATES the existing post-commit + // hydration: under the no-tags configuration the job_tags table is empty, the gate skips the + // round-trip entirely, and the claim hot path pays nothing. See TagsInUseAsync for the cheap, + // sound presence signal. var tagged = claimed.Count == 0 || !await TagsInUseAsync(connection, cancellationToken).ConfigureAwait(false) ? claimed : await WithTagsAsync(connection, claimed, cancellationToken).ConfigureAwait(false); @@ -556,14 +555,22 @@ private async ValueTask> ClaimQueueAsync( } // The single contended operation: UPDLOCK/READPAST is the dialect's skip-locked. + // + // Transition Log (§5.12): one Leased entry per claimed job at its post-claim Attempt, in the + // same batch and transaction as the lease write. The entries go in while the candidates hold + // only U locks, before the UPDATE takes X - see the note in ExpireLeasesUntracedAsync. + var recordTransitions = _historyPolicy != JobHistoryPolicy.Off; await using var claim = Cmd( - """ - WITH candidates AS ( - SELECT TOP (@take) job_id - FROM backwave.jobs WITH (UPDLOCK, READPAST, ROWLOCK) - WHERE queue = @queue AND state = 0 AND due_time <= @now - ORDER BY due_time, [sequence] - ) + $""" + {DeclareTransitionBatch} + INSERT INTO @batch (job_id, state, attempt, detail) + SELECT TOP (@take) job_id, 2, attempt + 1, NULL + FROM backwave.jobs WITH (UPDLOCK, READPAST, ROWLOCK) + WHERE queue = @queue AND state = 0 AND due_time <= @now + ORDER BY due_time, [sequence]; + + {(recordTransitions ? InsertTransitionsFromBatch : "")} + UPDATE j SET state = 2, attempt = j.attempt + 1, lease_owner = @worker, lease_expiry = @expiry OUTPUT inserted.job_id, inserted.wire_name, inserted.payload, inserted.queue, @@ -572,8 +579,8 @@ UPDATE j inserted.terminal_cause, inserted.schedule_id, inserted.parents_remaining, 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 + FROM @batch c + INNER LOOP JOIN backwave.jobs j ON j.job_id = c.job_id; """, connection, transaction); claim.Parameters.AddWithValue("queue", queue); @@ -582,8 +589,14 @@ FROM backwave.jobs j claim.Parameters.AddWithValue("worker", request.WorkerId); claim.Parameters.AddWithValue("expiry", request.Now + request.LeaseDuration); var queueClaims = new List(); + var maxNewOrdinal = -1L; await using (var reader = await claim.ExecuteReaderCountedAsync(cancellationToken).ConfigureAwait(false)) { + if (recordTransitions) + { + maxNewOrdinal = await ReadMaxOrdinalAsync(reader, cancellationToken).ConfigureAwait(false); + await reader.NextResultAsync(cancellationToken).ConfigureAwait(false); + } while (await reader.ReadAsync(cancellationToken).ConfigureAwait(false)) { var job = ReadJob(reader); @@ -598,12 +611,10 @@ FROM backwave.jobs j } // Crash after the lease write, before commit: rollback must un-lease every row (issue 0034). await FailpointAsync("claim", cancellationToken).ConfigureAwait(false); - // Transition Log (§5.12): one Leased entry per claimed job at its post-claim Attempt, in - // ONE set-based INSERT in this same transaction (atomic with the lease write). - await RecordTransitionsBatchAsync( - connection, transaction, - [.. queueClaims.Select(j => (j.JobId, JobState.Leased, j.Attempt, (string?)null))], - request.Now, cancellationToken).ConfigureAwait(false); + await PruneTransitionsAsync( + connection, transaction, maxNewOrdinal, + () => JsonSerializer.Serialize(queueClaims.Select(j => new TransitionRow(j.JobId, (int)JobState.Leased, j.Attempt, null)).ToArray()), + cancellationToken).ConfigureAwait(false); await transaction.CommitAsync(cancellationToken).ConfigureAwait(false); // OUTPUT does not guarantee order; the contract's per-Queue (DueTime, enqueue @@ -943,8 +954,12 @@ private async ValueTask> ReportOutcomesUntrac JobOutcome.Unroutable unroutable => (6, unroutable.Reason, null, now), _ => throw new ArgumentOutOfRangeException(nameof(batch)), }; + // Failure Detail rides only a failing transition, and only on the full history rung. + var detail = row.Outcome is JobOutcome.Failure && _historyPolicy == JobHistoryPolicy.TransitionsAndFailureDetail + ? options.Bounds.ClampFailureDetail(row.FailureDetail) + : null; rows[i] = new OutcomeRow( - row.JobId, row.WorkerId, row.Attempt, target.State, target.Cause, target.Due, target.TerminalAt); + row.JobId, row.WorkerId, row.Attempt, target.State, target.Cause, target.Due, target.TerminalAt, detail); } var payload = JsonSerializer.Serialize(rows); @@ -952,21 +967,44 @@ private async ValueTask> ReportOutcomesUntrac await using var transaction = (SqlTransaction)await connection .BeginTransactionAsync(cancellationToken).ConfigureAwait(false); - // One fenced multi-row UPDATE: OPENJSON unpacks the payload into a set, and the WHERE applies - // the per-(worker, attempt) Effect-Once fence to every row INDEPENDENTLY. A row whose lease is - // 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). + // One fenced batch: OPENJSON unpacks the payload into a set, and the WHERE applies the + // per-(worker, attempt) Effect-Once fence to every row INDEPENDENTLY. A row whose lease is no + // longer live simply fails to join and changes nothing (StaleLease). The fenced rows go into + // @batch under UPDLOCK, so they hold only U locks until the UPDATE. The UPDATE returns each + // applied row 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). retry_cause records a // handler-failure retry and is left alone on every terminal row. + // + // Transition Log (§5.12): one entry per fenced row for its resulting state at this Attempt, + // written between the fence and the UPDATE - see the note in ExpireLeasesUntracedAsync. The + // history policy Off appends nothing. var matched = new Dictionary(); + var recordTransitions = _historyPolicy != JobHistoryPolicy.Off; + var maxNewOrdinal = -1L; // 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 // cardinality from OPENJSON or from a VALUES list, picks a merge join over a full scan of // backwave.jobs, and takes a U lock on every row it passes. Two concurrent writers each // holding rows the other must scan past then deadlock (§5.5). await using (var update = Cmd( - """ + $""" + DECLARE @batch TABLE ( + job_id uniqueidentifier PRIMARY KEY, state int, attempt int, detail nvarchar(max), + cause nvarchar(max), due datetimeoffset, terminal_at datetimeoffset); + INSERT INTO @batch (job_id, state, attempt, detail, cause, due, terminal_at) + SELECT d.job_id, d.state, d.attempt, d.detail, d.cause, d.due, d.terminal_at + FROM OPENJSON(@payload) + WITH (job_id uniqueidentifier '$.JobId', worker nvarchar(450) '$.WorkerId', + attempt int '$.Attempt', state int '$.State', cause nvarchar(max) '$.Cause', + due datetimeoffset '$.Due', terminal_at datetimeoffset '$.TerminalAt', + detail nvarchar(max) '$.Detail') d + INNER LOOP JOIN backwave.jobs j WITH (UPDLOCK, ROWLOCK) ON j.job_id = d.job_id + WHERE j.state = 2 AND j.lease_owner = d.worker AND j.attempt = d.attempt + AND j.lease_expiry > @now; + + {(recordTransitions ? InsertTransitionsFromBatch : "")} + UPDATE j SET state = d.state, lease_owner = NULL, @@ -977,24 +1015,28 @@ UPDATE j 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', - attempt int '$.Attempt', state int '$.State', cause nvarchar(max) '$.Cause', - due datetimeoffset '$.Due', terminal_at datetimeoffset '$.TerminalAt') d - INNER LOOP JOIN backwave.jobs j ON j.job_id = d.job_id - WHERE j.state = 2 AND j.lease_owner = d.worker AND j.attempt = d.attempt - AND j.lease_expiry > @now + FROM @batch d + INNER LOOP JOIN backwave.jobs j ON j.job_id = d.job_id; """, connection, transaction)) { update.Parameters.Add("payload", SqlDbType.NVarChar, -1).Value = payload; update.Parameters.AddWithValue("now", now); await using var reader = await update.ExecuteReaderCountedAsync(cancellationToken).ConfigureAwait(false); + if (recordTransitions) + { + maxNewOrdinal = await ReadMaxOrdinalAsync(reader, cancellationToken).ConfigureAwait(false); + await reader.NextResultAsync(cancellationToken).ConfigureAwait(false); + } while (await reader.ReadAsync(cancellationToken).ConfigureAwait(false)) { matched[reader.GetGuid(0)] = reader.GetInt32(1); } } + await PruneTransitionsAsync( + connection, transaction, maxNewOrdinal, + () => JsonSerializer.Serialize(rows.Where(r => matched.ContainsKey(r.JobId)).ToArray()), + cancellationToken).ConfigureAwait(false); // Output and Tag deltas land ONLY for matched rows — a fenced-out (StaleLease) row leaves // nothing a stale node buffered. Output persists only on a Success outcome; Tags union onto @@ -1020,22 +1062,6 @@ AND j.lease_expiry > @now } } - // Transition Log (§5.12): one entry per matched row for its resulting state at this Attempt, - // in ONE set-based INSERT atomic with the outcome write. Failure Detail rides only a failing - // transition; every other outcome records null. The batch honors the history policy (Off - // appends nothing), so the noop-drain hot path adds no transition statements at all. - var transitionRows = new List<(Guid JobId, JobState State, int Attempt, string? FailureDetail)>(matched.Count); - foreach (var row in batch) - { - if (matched.TryGetValue(row.JobId, out var newState)) - { - transitionRows.Add((row.JobId, (JobState)newState, row.Attempt, - row.Outcome is JobOutcome.Failure ? row.FailureDetail : null)); - } - } - await RecordTransitionsBatchAsync(connection, transaction, transitionRows, now, cancellationToken) - .ConfigureAwait(false); - // First-level child-latch resolution for the matched TERMINAL ids only (a retry row stays // non-terminal and gates nothing). One lookup finds the few terminal parents that actually // gate a Dependency; the existing per-parent cascade then resolves each recursively, so a deep @@ -1099,7 +1125,7 @@ await ResolveChildLatchesAsync( // Property names are the OPENJSON '$.X' paths above; the CLR types map to the WITH column types. private sealed record OutcomeRow( Guid JobId, string WorkerId, int Attempt, int State, string? Cause, - DateTimeOffset? Due, DateTimeOffset? TerminalAt); + DateTimeOffset? Due, DateTimeOffset? TerminalAt, string? Detail); /// /// The latch (invariant I2), inside the same transaction as the terminal transition. @@ -1939,7 +1965,7 @@ private async Task RecordTransitionsBatchAsync( // OUTPUT the assigned ordinals so the prune can be skipped entirely (below) when no job in // the batch has reached the cap — the common 2-transition job pays no DELETE round-trip. - var maxNewOrdinal = -1L; + long maxNewOrdinal; // Materialize the payload into a keyed table variable before the INSERT. OPENJSON carries no // cardinality, and the key lets the optimizer cost the join to job_transitions. // The key does not settle the FK check to jobs. That plan stays the optimizer's to pick, and @@ -1949,44 +1975,75 @@ private async Task RecordTransitionsBatchAsync( // records transitions while it holds only U locks passes through concurrent scans, because U // and S are compatible - see the note in ExpireLeasesUntracedAsync. await using (var insert = Cmd( - """ - DECLARE @batch TABLE (job_id uniqueidentifier PRIMARY KEY, state int, attempt int, detail nvarchar(max)); + $""" + {DeclareTransitionBatch} INSERT INTO @batch (job_id, state, attempt, detail) SELECT job_id, state, attempt, detail FROM OPENJSON(@payload) WITH (job_id uniqueidentifier '$.JobId', state int '$.State', attempt int '$.Attempt', detail nvarchar(max) '$.Detail'); - INSERT INTO backwave.job_transitions (job_id, ordinal, recorded_at, state, attempt, failure_detail) - OUTPUT inserted.ordinal - SELECT d.job_id, COALESCE(t.maxord, -1) + 1, @now, d.state, d.attempt, d.detail - FROM @batch d - LEFT JOIN ( - SELECT job_id, MAX(ordinal) AS maxord - FROM backwave.job_transitions - WHERE job_id IN (SELECT job_id FROM @batch) - GROUP BY job_id - ) t ON t.job_id = d.job_id + {InsertTransitionsFromBatch} """, connection, transaction)) { insert.Parameters.Add("payload", SqlDbType.NVarChar, -1).Value = payload; insert.Parameters.Add("now", SqlDbType.DateTimeOffset).Value = now; await using var reader = await insert.ExecuteReaderCountedAsync(cancellationToken).ConfigureAwait(false); - while (await reader.ReadAsync(cancellationToken).ConfigureAwait(false)) + maxNewOrdinal = await ReadMaxOrdinalAsync(reader, cancellationToken).ConfigureAwait(false); + } + + await PruneTransitionsAsync(connection, transaction, maxNewOrdinal, () => payload, cancellationToken) + .ConfigureAwait(false); + } + + // The keyed table variable the Transition Log INSERT reads its rows from. + private const string DeclareTransitionBatch = + "DECLARE @batch TABLE (job_id uniqueidentifier PRIMARY KEY, state int, attempt int, detail nvarchar(max));"; + + // Appends one Transition Log entry per @batch row at the job's next ordinal and returns the + // assigned ordinals. The batch recorder, the claim, and the report write through this one + // statement. The caller declares @batch with at least the columns of DeclareTransitionBatch - + // job_id (the key), state, attempt, and detail - and binds @now. + private const string InsertTransitionsFromBatch = + """ + INSERT INTO backwave.job_transitions (job_id, ordinal, recorded_at, state, attempt, failure_detail) + OUTPUT inserted.ordinal + SELECT d.job_id, COALESCE(t.maxord, -1) + 1, @now, d.state, d.attempt, d.detail + FROM @batch d + LEFT JOIN ( + SELECT job_id, MAX(ordinal) AS maxord + FROM backwave.job_transitions + WHERE job_id IN (SELECT job_id FROM @batch) + GROUP BY job_id + ) t ON t.job_id = d.job_id; + """; + + // Reads the ordinals InsertTransitionsFromBatch returned and gives the highest, or -1 for none. + private static async ValueTask ReadMaxOrdinalAsync(SqlDataReader reader, CancellationToken cancellationToken) + { + var maxNewOrdinal = -1L; + while (await reader.ReadAsync(cancellationToken).ConfigureAwait(false)) + { + var ordinal = reader.GetInt64(0); + if (ordinal > maxNewOrdinal) { - var ordinal = reader.GetInt64(0); - if (ordinal > maxNewOrdinal) - { - maxNewOrdinal = ordinal; - } + maxNewOrdinal = ordinal; } } + return maxNewOrdinal; + } - // Per-job-life cap (§7): skip the prune entirely unless some job's new ordinal reached the - // cap — a job nowhere near MaxTransitionsPerJob never pays the DELETE. When some job did - // reach it, one set-based DELETE keeps only the newest MaxTransitionsPerJob per job (the - // correlated MAX no-ops for the jobs still under the cap). + // Per-job-life cap (§7) for a batch whose highest new ordinal is maxNewOrdinal. payload gives the + // batch as a JSON array of rows with a JobId, and only a prune that runs asks for it. + private async Task PruneTransitionsAsync( + SqlConnection connection, SqlTransaction transaction, long maxNewOrdinal, Func payload, + CancellationToken cancellationToken) + { + // Skip the prune entirely unless some job's new ordinal reached the cap — a job nowhere near + // MaxTransitionsPerJob never pays the DELETE. When some job did reach it, one set-based DELETE + // keeps only the newest MaxTransitionsPerJob per job (the correlated MAX no-ops for the jobs + // still under the cap). if (maxNewOrdinal < options.Bounds.MaxTransitionsPerJob) { return; @@ -1999,7 +2056,7 @@ WHERE jt.job_id IN (SELECT job_id FROM OPENJSON(@payload) AND jt.ordinal <= (SELECT MAX(ordinal) FROM backwave.job_transitions x WHERE x.job_id = jt.job_id) - @cap """, connection, transaction); - prune.Parameters.Add("payload", SqlDbType.NVarChar, -1).Value = payload; + prune.Parameters.Add("payload", SqlDbType.NVarChar, -1).Value = payload(); prune.Parameters.AddWithValue("cap", options.Bounds.MaxTransitionsPerJob); await prune.ExecuteNonQueryCountedAsync(cancellationToken).ConfigureAwait(false); } diff --git a/tests/BackWave.SqlServer.Tests/SqlServerConcurrentMaintenanceTests.cs b/tests/BackWave.SqlServer.Tests/SqlServerConcurrentMaintenanceTests.cs index 85feb62..7fb45b4 100644 --- a/tests/BackWave.SqlServer.Tests/SqlServerConcurrentMaintenanceTests.cs +++ b/tests/BackWave.SqlServer.Tests/SqlServerConcurrentMaintenanceTests.cs @@ -71,25 +71,39 @@ public async Task Concurrent_lease_sweeps_never_deadlock(JobHistoryPolicy policy Assert.Equal(0, faults.Terminal); } - // Every writer shares one transition-log INSERT, so its plan is compiled once and then reused. A - // wide lease sweep compiles it at 96 rows, where the optimizer serves the foreign-key check to - // backwave.jobs with a full scan, and every later caller inherits that plan whatever its own batch - // size holds. This test poisons the plan cache that way on purpose, then reports the narrow outcome - // batches a real fleet reports. That is the shape that lost a deadlock before the bounded retry. + // SQL Server caches one plan per batch text, and the report is one batch whose plan every later + // report reuses. When the first report compiles it at 96 rows, the optimizer serves the + // foreign-key check of the transition INSERT to backwave.jobs with a full scan, and every later + // report inherits that scan whatever its own batch size holds. This test clears the plan cache and + // poisons it that way on purpose, then reports the narrow outcome batches a real fleet reports. + // That is the shape that lost a deadlock before the bounded retry. [Fact] public async Task Concurrent_outcome_reports_never_deadlock() { const int Poison = 96, Rounds = 30, Workers = 6, PerWorker = 16; var store = await SqlServerTestDatabase.CreateFreshStoreAsync(JobHistoryPolicy.TransitionsAndFailureDetail); - var deadLetter = new RetryPolicy { MaxAttempts = 1, Backoff = _ => TimeSpan.FromMinutes(1) }.ToDisposition(); + // Without the clear, an earlier test in the run can leave a narrow plan in the cache, and this + // test then passes on any lock order. + await using (var connection = new SqlConnection(SqlServerTestDatabase.ConnectionString)) + { + await connection.OpenAsync(); + await using var clear = new SqlCommand("ALTER DATABASE SCOPED CONFIGURATION CLEAR PROCEDURE_CACHE", connection); + await clear.ExecuteNonQueryAsync(); + } for (var i = 0; i < Poison; i++) { await store.EnqueueAsync(new NewJob(Guid.NewGuid(), "t", "{}"u8.ToArray(), "poison", T0), T0); } - await store.ClaimAsync(new ClaimRequest("sweeper", ["poison"], Poison, Lease, T0)); - // MaxAttempts 1 dead-letters the swept jobs, so they never return to the claimable set. - await store.ExpireLeasesAsync(T0 + Lease + TimeSpan.FromSeconds(1), Poison, ["poison"], deadLetter); + // One claim never returns more than Bounds.MaxClaimBatch rows, so the poison is claimed in passes. + var poison = new List(); + while (poison.Count < Poison) + { + poison.AddRange(await store.ClaimAsync(new ClaimRequest("poisoner", ["poison"], Poison, Lease, T0))); + } + // Success is terminal, so the poison jobs never return to the claimable set. + await store.ReportOutcomesAsync( + [.. poison.Select(job => new OutcomeReport(job.JobId, "poisoner", job.Attempt, new JobOutcome.Success()))], T0); using var faults = new StoreFaultCounter(); var escaped = 0; @@ -260,6 +274,57 @@ await Execute( Assert.Equal(1, relinquished); } + // The two race pins above need the plan to lose the race before the counter moves. These two pin the + // lock order the races depend on, with one writer and no timing window: the claim and the report + // write the Transition Log while they hold only U locks, before the UPDATE takes X. + // + // A rival session holds an S lock on ONE row of the batch. U is compatible with S, so the lock step + // passes it and the transition INSERT runs. X is not, so the UPDATE waits on that row. A dirty read + // taken while the store waits then shows the order: the new entries are already there. A store that + // takes X before the INSERT waits at the same row with no entry written. + [Fact] + public async Task A_claim_writes_its_transitions_before_it_takes_X_locks() + { + const int Batch = 4; + var store = await SqlServerTestDatabase.CreateFreshStoreAsync(); + var jobs = new List(); + for (var i = 0; i < Batch; i++) + { + var id = Guid.NewGuid(); + await store.EnqueueAsync(new NewJob(id, "t", "{}"u8.ToArray(), "default", T0), T0); + jobs.Add(id); + } + + var written = await TransitionsWrittenWhileBlocked( + jobs[^1], + () => store.ClaimAsync(new ClaimRequest("w", ["default"], Batch, Lease, T0)).AsTask(), + JobState.Leased); + + Assert.Equal(Batch, written); + } + + [Fact] + public async Task A_report_writes_its_transitions_before_it_takes_X_locks() + { + const int Batch = 4; + var store = await SqlServerTestDatabase.CreateFreshStoreAsync(); + for (var i = 0; i < Batch; i++) + { + await store.EnqueueAsync(new NewJob(Guid.NewGuid(), "t", "{}"u8.ToArray(), "default", T0), T0); + } + var claimed = await store.ClaimAsync(new ClaimRequest("w", ["default"], Batch, Lease, T0)); + Assert.Equal(Batch, claimed.Count); + + var written = await TransitionsWrittenWhileBlocked( + claimed[^1].JobId, + () => store.ReportOutcomesAsync( + [.. claimed.Select(job => new OutcomeReport(job.JobId, "w", job.Attempt, new JobOutcome.Success()))], + T0).AsTask(), + JobState.Succeeded); + + Assert.Equal(Batch, written); + } + // The two pins above are worth nothing unless a real absorbed deadlock would move the counter. This // provokes one deterministically and walks the whole chain: SQL Server picks the store's transaction // as the victim, the bounded retry replays it, StoreFaultCounter sees the loss, and the caller still @@ -500,6 +565,38 @@ private static async Task WaitUntilBlockedBy(short session) Assert.Fail($"No session blocked on {session} within 30 s, so the deadlock cycle was never set up."); } + // Holds an S lock on one job while the write runs, waits until the write blocks on it, and counts the + // Transition Log entries in the given state that the blocked write already inserted. The rival only + // ever holds the one S lock, so it can never be half of a cycle. + private static async Task TransitionsWrittenWhileBlocked(Guid held, Func write, JobState state) + { + await using var rival = new SqlConnection(SqlServerTestDatabase.ConnectionString); + await rival.OpenAsync(); + var rivalSession = (short)(await Scalar(rival, null, "SELECT @@SPID"))!; + await using var rivalTx = (SqlTransaction)await rival.BeginTransactionAsync(); + await Execute(rival, rivalTx, + "SELECT job_id FROM backwave.jobs WITH (REPEATABLEREAD, ROWLOCK) WHERE job_id = @id", held); + + var writing = Task.Run(write); + int written; + try + { + await WaitUntilBlockedBy(rivalSession); + await using var reader = new SqlConnection(SqlServerTestDatabase.ConnectionString); + await reader.OpenAsync(); + written = (int)(await Scalar(reader, null, + $"SELECT count(*) FROM backwave.job_transitions WITH (NOLOCK) WHERE state = {(int)state}"))!; + } + finally + { + // Release the row whatever happened, so the blocked write can finish and be awaited. + await rivalTx.RollbackAsync(); + } + + await writing.WaitAsync(TimeSpan.FromSeconds(60)); + return written; + } + // One scalar query on a caller's session. private static async Task Scalar( SqlConnection connection, SqlTransaction? transaction, string sql, Guid? id = null) diff --git a/tests/BackWave.SqlServer.Tests/SqlServerRoundTripBudgetTests.cs b/tests/BackWave.SqlServer.Tests/SqlServerRoundTripBudgetTests.cs index 05a73f2..88acaf0 100644 --- a/tests/BackWave.SqlServer.Tests/SqlServerRoundTripBudgetTests.cs +++ b/tests/BackWave.SqlServer.Tests/SqlServerRoundTripBudgetTests.cs @@ -46,13 +46,13 @@ private static NewJob Job(string queue = "budget") => // arithmetic behind each number is in its test. private static readonly Budget Claim = new( - "ClaimBatchAsync of 32 jobs (one queue, cold caches)", Statements: 6); + "ClaimBatchAsync of 32 jobs (one queue, cold caches)", Statements: 5); private static readonly Budget ReportOutcomes = new( - "ReportOutcomesAsync of 32 succeeded rows", Statements: 3); + "ReportOutcomesAsync of 32 succeeded rows", Statements: 2); private static readonly Budget ReportOutcomesWithOutput = new( - "ReportOutcomesAsync of 32 succeeded rows, every one carrying job output", Statements: 35); + "ReportOutcomesAsync of 32 succeeded rows, every one carrying job output", Statements: 34); private static readonly Budget ExpireLeases = new( "ExpireLeasesAsync over 32 expired leases, all rescheduled", Statements: 3); @@ -79,12 +79,11 @@ public async Task Claim_of_a_full_batch_stays_within_its_round_trip_budget() await store.EnqueueAsync(Job(), T0); // also warms the one-time schema check, off the measured path } - // 1 queue-config applock + 1 queue_limits read + 1 claim UPDATE ... OUTPUT - // + 1 batched transition insert + 1 tags-in-use probe + 1 next-due read = 6, independent of - // batch size. The claim is a single UPDATE with OUTPUT, so 32 leased rows come back on the same - // trip that writes them, and the transition insert is set-based over OPENJSON, so 32 log entries - // cost one statement. No prune: the batch recorder issues a DELETE only when some job in it - // reached MaxTransitionsPerJob, and a freshly claimed job is on its second transition. + // 1 queue-config applock + 1 queue_limits read + 1 claim batch + 1 tags-in-use probe + // + 1 next-due read = 5, independent of batch size. The claim batch writes the transition log + // and the lease in one trip, and its UPDATE with OUTPUT returns the 32 leased rows on the same + // trip. No prune: the claim issues a DELETE only when some job in it reached + // MaxTransitionsPerJob, and a freshly claimed job is on its second transition. // // ClaimBatchAsync, not ClaimAsync: the extra statement over the plain claim is the next-due read, // which is this adapter's ONLY idle-wakeup mechanism (SQL Server has no Wake-Up Hint channel), so @@ -115,13 +114,13 @@ public async Task ReportOutcomes_of_a_full_batch_stays_within_its_round_trip_bud // The plain drain: every row succeeded, none carries output or a tag delta, so nothing but the // fenced state write and the transition log runs. - // 1 fenced batch UPDATE ... OUTPUT + 1 batched transition insert + 1 child-latch probe = 3, - // independent of batch size. The fence is applied per row inside that one UPDATE - OPENJSON - // unpacks the payload and the WHERE tests each row's (worker, attempt) independently - and - // OUTPUT reports which rows matched, so the per-row Effect-Once verdict costs no extra trip. + // 1 fenced batch + 1 child-latch probe = 2, independent of batch size. The fenced batch writes + // the transition log and the outcomes in one trip. The fence is applied per row inside it - + // OPENJSON unpacks the payload and the WHERE tests each row's (worker, attempt) independently - + // and OUTPUT reports which rows matched, so the per-row Effect-Once verdict costs no extra trip. // The child-latch probe runs because Succeeded is terminal: one lookup asks whether ANY of the // 32 ids parents a Dependency, and the answer here is no, so nothing cascades. - // As above, no job in this batch is near the cap, so the batch recorder issues no prune DELETE. + // As above, no job in this batch is near the cap, so the report issues no prune DELETE. var batch = claimed .Select(job => new OutcomeReport(job.JobId, "budget-worker", job.Attempt, new JobOutcome.Success())) .ToArray(); @@ -149,7 +148,7 @@ public async Task ReportOutcomes_carrying_job_output_still_costs_one_statement_p var claimed = await store.ClaimAsync(new ClaimRequest("budget-worker", ["budget"], ClaimBatch, Lease, T0)); Assert.Equal(ClaimBatch, claimed.Count); - // The drain budget above plus ONE STATEMENT PER ROW: 3 + 32 = 35. Output is the one write on + // The drain budget above plus ONE STATEMENT PER ROW: 2 + 32 = 34. Output is the one write on // this path that is still a per-row loop on this adapter - the fenced UPDATE cannot carry the // blob, because OPENJSON has no varbinary(max) column type, so a blob would have to go over as // base64 text and be converted back per row.