diff --git a/CHANGELOG.md b/CHANGELOG.md index 604d35e1a..8cf5a59d4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - 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 SQLite job cleanup stalling when excluded queues fill the oldest batch, and added support for included queue filters so per-queue retention works on SQLite. [PR #1417](https://github.com/riverqueue/river/pull/1417). +- 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 0c57698e0..6e651cbfe 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql @@ -290,7 +290,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 @@ -340,7 +341,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 1666bbbfd..c18c912e0 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go @@ -800,7 +800,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 ` @@ -907,7 +908,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 `