From ff417dff91fb5f4a37783fd19f7e05c31e1b5b35 Mon Sep 17 00:00:00 2001 From: Jack Danger Date: Thu, 1 Oct 2026 06:58:31 +0000 Subject: [PATCH] Fix unique skip rewriting the existing job's kind under `ExcludeKind` The unique insert upsert returns the conflicting row with `DO UPDATE SET kind = EXCLUDED.kind`. With `UniqueOpts.ExcludeKind`, jobs of different kinds can share a unique key, so a skipped insert of kind B rewrote the existing kind A job to kind B. That job was then worked by B's worker with A's args, or failed as an unknown kind. Set the kind to the existing row's own instead, which keeps the update (and the returned row) but makes it a no-op. Changed in `JobInsertFastMany` for Postgres and SQLite, and in SQLite's unused `JobInsertFast` for consistency. Signed-off-by: Jack Danger --- CHANGELOG.md | 1 + .../internal/dbsqlc/river_job.sql.go | 3 +- riverdriver/riverdrivertest/job_insert.go | 43 +++++++++++++++++++ .../riverpgxv5/internal/dbsqlc/river_job.sql | 3 +- .../internal/dbsqlc/river_job.sql.go | 3 +- .../riversqlite/internal/dbsqlc/river_job.sql | 6 ++- .../internal/dbsqlc/river_job.sql.go | 6 ++- 7 files changed, 58 insertions(+), 7 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d4dcc0efc..79eef4c0a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed - Fixed `rivertest.Worker` to honor a configured `Config.JobStuckThreshold` for stuck job detection. Previously, it always used an internal 5 second threshold, so the `Job appears to be stuck` log line was emitted at a different time than it would be under a real client. [PR #1418](https://github.com/riverqueue/river/pull/1418). +- Fixed a unique insert skipped as a duplicate of a job of a different kind (possible with `UniqueOpts.ExcludeKind`) changing the existing job's kind to its own. [PR #1421](https://github.com/riverqueue/river/pull/1421). ## [0.48.0] - 2026-09-30 diff --git a/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go b/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go index 8531a9d85..37d3476ab 100644 --- a/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go @@ -757,7 +757,8 @@ ON CONFLICT (unique_key) AND unique_states IS NOT NULL AND /* TEMPLATE: schema */river_job_state_in_bitmask(unique_states, state) -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under ` + "`" + `ExcludeKind` + "`" + `. + DO UPDATE SET kind = river_job.kind 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, /* TEMPLATE_BEGIN: unique_skipped_as_duplicate */ (xmax != 0) /* TEMPLATE_END */ AS unique_skipped_as_duplicate diff --git a/riverdriver/riverdrivertest/job_insert.go b/riverdriver/riverdrivertest/job_insert.go index a911d8cc0..172f2b971 100644 --- a/riverdriver/riverdrivertest/job_insert.go +++ b/riverdriver/riverdrivertest/job_insert.go @@ -236,6 +236,49 @@ func exerciseJobInsert[TTx any](ctx context.Context, t *testing.T, require.Equal(t, results1[0].Job.ID, results2[0].Job.ID) }) + // With UniqueOpts.ExcludeKind, jobs of different kinds share a unique + // key. A skipped insert must return the incumbent row unmodified rather + // than rewriting its kind to the kind of the skipped job. + t.Run("UniqueConflictDifferentKindLeavesIncumbentKind", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + + insertParams := func(kind string) *riverdriver.JobInsertFastManyParams { + return &riverdriver.JobInsertFastManyParams{ + Jobs: []*riverdriver.JobInsertFastParams{ + { + EncodedArgs: []byte(`{"encoded": "args"}`), + Kind: kind, + MaxAttempts: rivercommon.MaxAttemptsDefault, + Priority: rivercommon.PriorityDefault, + Queue: rivercommon.QueueDefault, + State: rivertype.JobStateAvailable, + Tags: []string{}, + UniqueKey: []byte("unique-key-exclude-kind"), + UniqueStates: 0xff, + }, + }, + } + } + + results1, err := exec.JobInsertFastMany(ctx, insertParams("kind_a")) + require.NoError(t, err) + require.Len(t, results1, 1) + require.False(t, results1[0].UniqueSkippedAsDuplicate) + + results2, err := exec.JobInsertFastMany(ctx, insertParams("kind_b")) + require.NoError(t, err) + require.Len(t, results2, 1) + require.True(t, results2[0].UniqueSkippedAsDuplicate) + require.Equal(t, results1[0].Job.ID, results2[0].Job.ID) + require.Equal(t, "kind_a", results2[0].Job.Kind) + + job, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: results1[0].Job.ID}) + require.NoError(t, err) + require.Equal(t, "kind_a", job.Kind) + }) + t.Run("UniqueConflictWithinBatch", func(t *testing.T) { t.Parallel() diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql index f61ceed01..d87361114 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql @@ -329,7 +329,8 @@ ON CONFLICT (unique_key) AND unique_states IS NOT NULL AND /* TEMPLATE: schema */river_job_state_in_bitmask(unique_states, state) -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under `ExcludeKind`. + DO UPDATE SET kind = river_job.kind RETURNING sqlc.embed(river_job), /* TEMPLATE_BEGIN: unique_skipped_as_duplicate */ (xmax != 0) /* TEMPLATE_END */ AS unique_skipped_as_duplicate; diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go index 2fc4d014b..ac4a5f076 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go @@ -733,7 +733,8 @@ ON CONFLICT (unique_key) AND unique_states IS NOT NULL AND /* TEMPLATE: schema */river_job_state_in_bitmask(unique_states, state) -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under ` + "`" + `ExcludeKind` + "`" + `. + DO UPDATE SET kind = river_job.kind 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, /* TEMPLATE_BEGIN: unique_skipped_as_duplicate */ (xmax != 0) /* TEMPLATE_END */ AS unique_skipped_as_duplicate diff --git a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql index 2fbc2c734..62a487aa1 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql @@ -299,7 +299,8 @@ ON CONFLICT (unique_key) ELSE 0 END >= 1 -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under `ExcludeKind`. + DO UPDATE SET kind = river_job.kind RETURNING *; -- name: JobInsertFastMany :many @@ -349,7 +350,8 @@ ON CONFLICT (unique_key) ELSE 0 END >= 1 -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under `ExcludeKind`. + DO UPDATE SET kind = river_job.kind RETURNING *; -- name: JobInsertFastNoReturning :execrows diff --git a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go index c46d0ae51..4f3bd2739 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go @@ -813,7 +813,8 @@ ON CONFLICT (unique_key) ELSE 0 END >= 1 -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under ` + "`" + `ExcludeKind` + "`" + `. + DO UPDATE SET kind = river_job.kind RETURNING id, json(args), attempt, attempted_at, json(attempted_by), created_at, json(errors), finalized_at, kind, max_attempts, json(metadata), priority, queue, state, scheduled_at, json(tags), unique_key, unique_states ` @@ -920,7 +921,8 @@ ON CONFLICT (unique_key) ELSE 0 END >= 1 -- Something needs to be updated for a row to be returned on a conflict. - DO UPDATE SET kind = EXCLUDED.kind + -- Keep the existing kind, which may differ under ` + "`" + `ExcludeKind` + "`" + `. + DO UPDATE SET kind = river_job.kind RETURNING id, json(args), attempt, attempted_at, json(attempted_by), created_at, json(errors), finalized_at, kind, max_attempts, json(metadata), priority, queue, state, scheduled_at, json(tags), unique_key, unique_states `