Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Fixed `river bench` inserting every benchmark job with a `num` arg of `0` instead of numbering jobs sequentially. [PR #1379](https://github.com/riverqueue/river/pull/1379).
- Fixed `rivermigrate` leaving `river_migration` rows behind after migrating a non-main migration line down through its version 1, which caused a later up migration of that line to skip version 1. With `MigrateTx`, rows for every removed version were left behind. [PR #1378](https://github.com/riverqueue/river/pull/1378).
- Fixed the job completer panicking when a job it was finalizing had its state changed concurrently, like being moved to `pending` out of band, or being rescued while the completer's update waited on the row lock (in which case PostgreSQL returns the job's pre-update `running` row). Such jobs are now skipped without emitting a completion event. [PR #1383](https://github.com/riverqueue/river/pull/1383).
- 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

Expand Down
11 changes: 7 additions & 4 deletions riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

54 changes: 54 additions & 0 deletions riverdriver/riverdrivertest/client_cancel_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package riverdrivertest

import (
"context"
"sync"
"testing"
"time"

Expand Down Expand Up @@ -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)
}
8 changes: 8 additions & 0 deletions riverdriver/riverdrivertest/driver_client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
11 changes: 7 additions & 4 deletions riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
11 changes: 7 additions & 4 deletions riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading