From 348aa0138115e75d6bbf6d679cdccc3caf767e3c Mon Sep 17 00:00:00 2001 From: Jack Danger Date: Mon, 14 Sep 2026 20:09:21 -0500 Subject: [PATCH 1/3] Test concurrent `JobCancel` calls see the committed cancelled row MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Pin the documented "Returns the up-to-date JobRow" guarantee: two callers cancel the same job at the same moment, twenty rounds, and each loser must report `cancelled` with the winner's `finalized_at` — not its own pre-commit snapshot of the row. These tests fail until the locking fallback read lands. Signed-off-by: Jack Danger --- client_test.go | 57 ++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 57 insertions(+) diff --git a/client_test.go b/client_test.go index 4a9ae1612..241c49ba5 100644 --- a/client_test.go +++ b/client_test.go @@ -1476,6 +1476,63 @@ func Test_Client_Common(t *testing.T) { client.producersByQueueName[QueueDefault].testSignals.QueueControlEventTriggered.RequireEmpty() }) + t.Run("ConcurrentCancelSingleFinalizerAndFreshReturn", func(t *testing.T) { + t.Parallel() + + client, _ := setup(t) + + subscribeChan := subscribe(t, client) + startClient(ctx, t, client) + + // Scheduled far out so no producer can ever work it. + insertRes, err := client.Insert(ctx, &noOpArgs{}, &InsertOpts{ScheduledAt: time.Now().Add(5 * time.Minute)}) + require.NoError(t, err) + + const cancelRounds = 20 + + var firstFinalizedAt *time.Time + + for range cancelRounds { + var ( + group sync.WaitGroup + rows [2]*rivertype.JobRow + errs [2]error + ) + + group.Go(func() { + rows[0], errs[0] = client.JobCancel(ctx, insertRes.Job.ID) + }) + group.Go(func() { + rows[1], errs[1] = client.JobCancel(ctx, insertRes.Job.ID) + }) + group.Wait() + + for i := range 2 { + require.NoError(t, errs[i]) + require.Equal(t, rivertype.JobStateCancelled, rows[i].State) + } + + require.Equal(t, *rows[0].FinalizedAt, *rows[1].FinalizedAt) + + if firstFinalizedAt == nil { + firstFinalizedAt = rows[0].FinalizedAt + } + require.Equal(t, *firstFinalizedAt, *rows[0].FinalizedAt, + "finalized_at must be written exactly once; later cancels must not re-stamp it") + } + + finalRow, err := client.JobGet(ctx, insertRes.Job.ID) + require.NoError(t, err) + require.Equal(t, rivertype.JobStateCancelled, finalRow.State) + require.Equal(t, *firstFinalizedAt, *finalRow.FinalizedAt) + + select { + case event := <-subscribeChan: + t.Fatalf("expected no job events for operator cancels of a never-worked job, got: %v", event) + case <-time.After(500 * time.Millisecond): + } + }) + t.Run("AlternateSchema", func(t *testing.T) { t.Parallel() From ded4f6d007c5fb632cd95240f9ce07daff95444c Mon Sep 17 00:00:00 2001 From: Jack Danger Date: Mon, 28 Sep 2026 21:58:16 -0700 Subject: [PATCH 2/3] Fix stale `JobCancel` return to the loser of a concurrent cancel race MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `JobCancel` documents "Returns the up-to-date JobRow". When two cancels race, the loser's update matches zero rows and the query falls through to its UNION fallback arm — a plain (non-locking) scan that runs on the loser's pre-commit snapshot, so the loser is returned e.g. still "scheduled" though "cancelled" is durably committed. Lock that fallback read (FOR UPDATE) so it re-reads the committed row: the same mechanism `locked_job` already applies at the top of the same statement, so no new lock ordering, wait, or deadlock surface. At REPEATABLE READ/SERIALIZABLE the locking read turns a raced cancel into a serialization error instead of a stale row; River runs READ COMMITTED. Regenerated both Postgres dialects with sqlc; SQLite's JobCancel is a single UPDATE .. RETURNING and unaffected. Signed-off-by: Jack Danger --- CHANGELOG.md | 1 + .../riverdatabasesql/internal/dbsqlc/river_job.sql.go | 11 +++++++---- riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql | 11 +++++++---- .../riverpgxv5/internal/dbsqlc/river_job.sql.go | 11 +++++++---- 4 files changed, 22 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2e0952855..ef6950e26 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,6 +41,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Fixed SQLite notification listeners delivering notifications from before a subscription or from an unsubscribe gap. Notification reads now fetch subscribed topics in bounded batches, and cleanup deletes expired notifications in batches of 10,000 rows (reduced to 1,000 after repeated timeouts), with pauses between batches to reduce write lock contention. [PR #1381](https://github.com/riverqueue/river/pull/1381). - Fixed SQLite reusing the ID of a deleted job when that job held the largest ID, which could cause an ID observed earlier to refer to an unrelated job later. [PR #1390](https://github.com/riverqueue/river/pull/1390). - Fixed the `Job appears to be stuck` log line reporting the client-level `JobTimeout` instead of the worker-level timeout when a worker overrides `Timeout`. [PR #1394](https://github.com/riverqueue/river/pull/1394). +- Fixed `JobCancel` returning a stale pre-commit row to the loser of a concurrent-cancel race. The query's fallback read now takes a row lock (`FOR UPDATE`), matching the documented "returns the up-to-date `JobRow`" contract. The analogous shape in `JobRetry` is known and will follow separately. [PR #1409](https://github.com/riverqueue/river/pull/1409). ## [0.47.0] - 2026-09-01 diff --git a/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go b/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go index eeee90039..4c329629e 100644 --- a/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go @@ -48,10 +48,13 @@ updated_job AS ( WHERE river_job.id = notification.id RETURNING river_job.id, river_job.args, river_job.attempt, river_job.attempted_at, river_job.attempted_by, river_job.created_at, river_job.errors, river_job.finalized_at, river_job.kind, river_job.max_attempts, river_job.metadata, river_job.priority, river_job.queue, river_job.state, river_job.scheduled_at, river_job.tags, river_job.unique_key, river_job.unique_states ) -SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states -FROM /* TEMPLATE: schema */river_job -WHERE id = $1::bigint - AND id NOT IN (SELECT id FROM updated_job) +SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states FROM ( + SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states + FROM /* TEMPLATE: schema */river_job + WHERE id = $1::bigint + FOR UPDATE +) AS fallback_job +WHERE fallback_job.id NOT IN (SELECT id FROM updated_job) UNION SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states FROM updated_job diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql index e045848b2..b60cc1135 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql @@ -72,10 +72,13 @@ updated_job AS ( WHERE river_job.id = notification.id RETURNING river_job.* ) -SELECT * -FROM /* TEMPLATE: schema */river_job -WHERE id = @id::bigint - AND id NOT IN (SELECT id FROM updated_job) +SELECT * FROM ( + SELECT * + FROM /* TEMPLATE: schema */river_job + WHERE id = @id::bigint + FOR UPDATE +) AS fallback_job +WHERE fallback_job.id NOT IN (SELECT id FROM updated_job) UNION SELECT * FROM updated_job; diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go index b3a230bea..79451a5ea 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go @@ -48,10 +48,13 @@ updated_job AS ( WHERE river_job.id = notification.id RETURNING river_job.id, river_job.args, river_job.attempt, river_job.attempted_at, river_job.attempted_by, river_job.created_at, river_job.errors, river_job.finalized_at, river_job.kind, river_job.max_attempts, river_job.metadata, river_job.priority, river_job.queue, river_job.state, river_job.scheduled_at, river_job.tags, river_job.unique_key, river_job.unique_states ) -SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states -FROM /* TEMPLATE: schema */river_job -WHERE id = $1::bigint - AND id NOT IN (SELECT id FROM updated_job) +SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states FROM ( + SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states + FROM /* TEMPLATE: schema */river_job + WHERE id = $1::bigint + FOR UPDATE +) AS fallback_job +WHERE fallback_job.id NOT IN (SELECT id FROM updated_job) UNION SELECT id, args, attempt, attempted_at, attempted_by, created_at, errors, finalized_at, kind, max_attempts, metadata, priority, queue, state, scheduled_at, tags, unique_key, unique_states FROM updated_job From b4cce99026dc6f6dd747d9c22c5ec634fd78f105 Mon Sep 17 00:00:00 2001 From: Jack Danger Date: Tue, 29 Sep 2026 17:15:58 -0700 Subject: [PATCH 3/3] Move concurrent-cancel race test to the shared driver suite Per review: the test now lives in riverdrivertest and runs under every driver fixture (pgxv5, both database/sql Postgres dialects, and the SQLite family). The fixed-wait negative-event assertion does not survive the move; noted inline. Signed-off-by: Jack Danger --- client_test.go | 57 ------------------- .../riverdrivertest/client_cancel_test.go | 54 ++++++++++++++++++ .../riverdrivertest/driver_client_test.go | 8 +++ 3 files changed, 62 insertions(+), 57 deletions(-) diff --git a/client_test.go b/client_test.go index 241c49ba5..4a9ae1612 100644 --- a/client_test.go +++ b/client_test.go @@ -1476,63 +1476,6 @@ func Test_Client_Common(t *testing.T) { client.producersByQueueName[QueueDefault].testSignals.QueueControlEventTriggered.RequireEmpty() }) - t.Run("ConcurrentCancelSingleFinalizerAndFreshReturn", func(t *testing.T) { - t.Parallel() - - client, _ := setup(t) - - subscribeChan := subscribe(t, client) - startClient(ctx, t, client) - - // Scheduled far out so no producer can ever work it. - insertRes, err := client.Insert(ctx, &noOpArgs{}, &InsertOpts{ScheduledAt: time.Now().Add(5 * time.Minute)}) - require.NoError(t, err) - - const cancelRounds = 20 - - var firstFinalizedAt *time.Time - - for range cancelRounds { - var ( - group sync.WaitGroup - rows [2]*rivertype.JobRow - errs [2]error - ) - - group.Go(func() { - rows[0], errs[0] = client.JobCancel(ctx, insertRes.Job.ID) - }) - group.Go(func() { - rows[1], errs[1] = client.JobCancel(ctx, insertRes.Job.ID) - }) - group.Wait() - - for i := range 2 { - require.NoError(t, errs[i]) - require.Equal(t, rivertype.JobStateCancelled, rows[i].State) - } - - require.Equal(t, *rows[0].FinalizedAt, *rows[1].FinalizedAt) - - if firstFinalizedAt == nil { - firstFinalizedAt = rows[0].FinalizedAt - } - require.Equal(t, *firstFinalizedAt, *rows[0].FinalizedAt, - "finalized_at must be written exactly once; later cancels must not re-stamp it") - } - - finalRow, err := client.JobGet(ctx, insertRes.Job.ID) - require.NoError(t, err) - require.Equal(t, rivertype.JobStateCancelled, finalRow.State) - require.Equal(t, *firstFinalizedAt, *finalRow.FinalizedAt) - - select { - case event := <-subscribeChan: - t.Fatalf("expected no job events for operator cancels of a never-worked job, got: %v", event) - case <-time.After(500 * time.Millisecond): - } - }) - t.Run("AlternateSchema", func(t *testing.T) { t.Parallel() diff --git a/riverdriver/riverdrivertest/client_cancel_test.go b/riverdriver/riverdrivertest/client_cancel_test.go index 601c5c882..ecd226de8 100644 --- a/riverdriver/riverdrivertest/client_cancel_test.go +++ b/riverdriver/riverdrivertest/client_cancel_test.go @@ -2,6 +2,7 @@ package riverdrivertest import ( "context" + "sync" "testing" "time" @@ -100,3 +101,56 @@ func exerciseClientCancelRunningJob[TTx any](ctx context.Context, t *testing.T, require.Equal(t, insertRes.Job.ID, event.Job.ID) require.Equal(t, rivertype.JobStateCancelled, event.Job.State) } + +func exerciseClientCancelConcurrentRaceFreshReturn[TTx any](ctx context.Context, t *testing.T, driver riverdriver.Driver[TTx], schema string) { + t.Helper() + + config := newTestConfig(t, schema) + + client, err := river.NewClient(driver, config) + require.NoError(t, err) + + // Scheduled far out so no producer can ever work it; the root suite's + // fixed-wait negative-event assertion is dropped (remote-backend + // flake-prone). + insertRes, err := client.Insert(ctx, &noOpArgs{}, &river.InsertOpts{ScheduledAt: time.Now().Add(5 * time.Minute)}) + require.NoError(t, err) + + const cancelRounds = 20 + + var firstFinalizedAt *time.Time + + for range cancelRounds { + var ( + group sync.WaitGroup + rows [2]*rivertype.JobRow + errs [2]error + ) + + group.Go(func() { + rows[0], errs[0] = client.JobCancel(ctx, insertRes.Job.ID) + }) + group.Go(func() { + rows[1], errs[1] = client.JobCancel(ctx, insertRes.Job.ID) + }) + group.Wait() + + for i := range 2 { + require.NoError(t, errs[i]) + require.Equal(t, rivertype.JobStateCancelled, rows[i].State) + } + + require.Equal(t, *rows[0].FinalizedAt, *rows[1].FinalizedAt) + + if firstFinalizedAt == nil { + firstFinalizedAt = rows[0].FinalizedAt + } + require.Equal(t, *firstFinalizedAt, *rows[0].FinalizedAt, + "finalized_at must be written exactly once; later cancels must not re-stamp it") + } + + finalRow, err := client.JobGet(ctx, insertRes.Job.ID) + require.NoError(t, err) + require.Equal(t, rivertype.JobStateCancelled, finalRow.State) + require.Equal(t, *firstFinalizedAt, *finalRow.FinalizedAt) +} diff --git a/riverdriver/riverdrivertest/driver_client_test.go b/riverdriver/riverdrivertest/driver_client_test.go index a1b6af22a..3689ad228 100644 --- a/riverdriver/riverdrivertest/driver_client_test.go +++ b/riverdriver/riverdrivertest/driver_client_test.go @@ -634,6 +634,14 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.False(t, insertRes2.UniqueSkippedAsDuplicate) }) + t.Run("JobCancelConcurrentRaceFreshReturn", func(t *testing.T) { + t.Parallel() + + _, bundle := setupConfig(t) + + exerciseClientCancelConcurrentRaceFreshReturn(ctx, t, bundle.driver, bundle.schema) + }) + t.Run("JobDelete", func(t *testing.T) { t.Parallel()