From 3ccb7080f0240db5f701c63630b5e09ed4152a55 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Sun, 27 Sep 2026 16:48:26 -0500 Subject: [PATCH] fix SQLite batch unique-key handling SQLite gives every job in a returning batch the same nonce. When two jobs share an active unique key, the second conflict returns the first job's nonce and both results claim insertion even though only one row exists. PostgreSQL rejects this batch. Reject repeated active unique keys before SQLite writes the batch and give each row its own 16-character nonce. This makes batch conflicts fail atomically like PostgreSQL while keeping the no-returning path free of nonce metadata. Exercise both client batch APIs and the shared driver suite, including keys outside their enforced states and per-row nonce formatting. --- CHANGELOG.md | 1 + .../riverdrivertest/driver_client_test.go | 51 ++++++++++++ riverdriver/riverdrivertest/job_insert.go | 81 +++++++++++++++++++ .../riversqlite/river_sqlite_driver.go | 37 ++++++--- riverdriver/unique_insert.go | 13 ++- riverdriver/unique_insert_test.go | 28 +++++++ 6 files changed, 198 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 405fbe9a5..6238c6183 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Fixed YugabyteDB clients relying on notifications when `LISTEN/NOTIFY` is unavailable or disabled. Clients now automatically poll for running job cancellations and queue pause, resume, and metadata changes with the default `PollOnly: false`, and skip unsupported notification broadcasts. Native notifications require YugabyteDB 2025.2.3 or later with `ysql_yb_enable_listen_notify=true` on both Masters and TServers. [PR #1347](https://github.com/riverqueue/river/pull/1347). - Fixed job cancellations received during a fetch being lost before the fetched jobs started. Matching jobs now receive cancellation before work begins. [PR #1397](https://github.com/riverqueue/river/pull/1397). - Fixed SQLite `JobCancel` and `JobCancelTx` notifying running workers through the shared control outbox, so their contexts are cancelled when the transaction commits. [PR #1398](https://github.com/riverqueue/river/pull/1398). +- Fixed SQLite `InsertMany` and `InsertManyTx` reporting multiple inserted jobs when a batch contains the same active unique key more than once. Such batches now fail atomically, matching PostgreSQL. [PR #1399](https://github.com/riverqueue/river/pull/1399). - Fixed `UniqueOpts.ByArgs` skipping distinct jobs or failing inserts when JSON keys contain path syntax (like `user.id`), are empty, or come from unnamed tags like `json:",omitempty"`. Unaffected unique keys remain unchanged; affected jobs may be inserted again after upgrading or by old and new clients during a rolling upgrade. [PR #1387](https://github.com/riverqueue/river/pull/1387). - Fixed SQLite job list pagination skipping or repeating jobs by formatting cursor timestamps consistently with stored timestamps. [PR #1374](https://github.com/riverqueue/river/pull/1374). - Fixed SQLite drivers deleting jobs in a finalized state whose retention period was set to -1 (keep forever), like `Config.DiscardedJobRetentionPeriod: -1`, whenever another state's retention period was finite. [PR #1389](https://github.com/riverqueue/river/pull/1389). diff --git a/riverdriver/riverdrivertest/driver_client_test.go b/riverdriver/riverdrivertest/driver_client_test.go index 43c0a7b63..a1b6af22a 100644 --- a/riverdriver/riverdrivertest/driver_client_test.go +++ b/riverdriver/riverdrivertest/driver_client_test.go @@ -413,6 +413,57 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, }) } + t.Run("InsertManyDuplicateUniqueKeysInBatch", func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + params := []river.InsertManyParams{ + {Args: noOpArgs{Name: "same"}, InsertOpts: &river.InsertOpts{UniqueOpts: river.UniqueOpts{ByArgs: true}}}, + {Args: noOpArgs{Name: "same"}, InsertOpts: &river.InsertOpts{UniqueOpts: river.UniqueOpts{ByArgs: true}}}, + } + + results, err := client.InsertMany(ctx, params) + jobs, getErr := bundle.exec.JobGetByKindMany(ctx, &riverdriver.JobGetByKindManyParams{ + Kind: []string{(noOpArgs{}).Kind()}, + Schema: bundle.schema, + }) + require.NoError(t, getErr) + if bundle.driver.DatabaseName() == riverdriver.DatabaseNameSQLite { + require.ErrorContains(t, err, "unique key appears more than once in batch") + } else { + require.ErrorContains(t, err, "ON CONFLICT DO UPDATE command cannot affect row a second time") + } + require.Empty(t, results) + require.Empty(t, jobs) + }) + + t.Run("InsertManyTxDuplicateUniqueKeysInBatch", func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + tx, execTx := beginTx(ctx, t, bundle) + + results, err := client.InsertManyTx(ctx, tx, []river.InsertManyParams{ + {Args: noOpArgs{Name: "same"}, InsertOpts: &river.InsertOpts{UniqueOpts: river.UniqueOpts{ByArgs: true}}}, + {Args: noOpArgs{Name: "same"}, InsertOpts: &river.InsertOpts{UniqueOpts: river.UniqueOpts{ByArgs: true}}}, + }) + if bundle.driver.DatabaseName() == riverdriver.DatabaseNameSQLite { + require.ErrorContains(t, err, "unique key appears more than once in batch") + } else { + require.ErrorContains(t, err, "ON CONFLICT DO UPDATE command cannot affect row a second time") + } + require.Empty(t, results) + require.NoError(t, execTx.Rollback(ctx)) + + jobs, err := bundle.exec.JobGetByKindMany(ctx, &riverdriver.JobGetByKindManyParams{ + Kind: []string{(noOpArgs{}).Kind()}, + Schema: bundle.schema, + }) + require.NoError(t, err) + require.Empty(t, jobs) + }) + // Keys containing gjson/sjson path syntax (and the empty key) are distinct // keys when unique by all args, so args differing in their values aren't // duplicates. diff --git a/riverdriver/riverdrivertest/job_insert.go b/riverdriver/riverdrivertest/job_insert.go index 9bc8cf64a..e60a68c36 100644 --- a/riverdriver/riverdrivertest/job_insert.go +++ b/riverdriver/riverdrivertest/job_insert.go @@ -81,9 +81,17 @@ func exerciseJobInsert[TTx any](ctx context.Context, t *testing.T, require.NoError(t, err) require.Len(t, resultRows, len(insertParams)) + uniqueNonces := make(map[string]bool, len(resultRows)) for i, result := range resultRows { require.False(t, result.UniqueSkippedAsDuplicate) job := result.Job + if bundle.driver.DatabaseName() == riverdriver.DatabaseNameSQLite { + nonce, ok := riverdriver.UniqueInsertMetadataNonce(job.Metadata) + require.True(t, ok) + require.Regexp(t, `^[0-9a-f]{16}$`, nonce) + require.False(t, uniqueNonces[nonce]) + uniqueNonces[nonce] = true + } // SQLite needs to set a special metadata key to be able to // check for duplicates. Remove this for purposes of comparing @@ -228,6 +236,79 @@ func exerciseJobInsert[TTx any](ctx context.Context, t *testing.T, require.Equal(t, results1[0].Job.ID, results2[0].Job.ID) }) + t.Run("UniqueConflictWithinBatch", func(t *testing.T) { + t.Parallel() + + driver, schema := driverWithSchema(ctx, t, nil) + exec := driver.GetExecutor() + job := &riverdriver.JobInsertFastParams{ + EncodedArgs: []byte(`{"encoded": "args"}`), + Kind: "test_kind", + MaxAttempts: rivercommon.MaxAttemptsDefault, + Priority: rivercommon.PriorityDefault, + Queue: rivercommon.QueueDefault, + State: rivertype.JobStateAvailable, + Tags: []string{}, + UniqueKey: []byte("unique-key-within-batch"), + UniqueStates: 0xff, + } + + results, err := exec.JobInsertFastMany(ctx, &riverdriver.JobInsertFastManyParams{ + Jobs: []*riverdriver.JobInsertFastParams{job, job}, + Schema: schema, + }) + if driver.DatabaseName() == riverdriver.DatabaseNameSQLite { + require.ErrorContains(t, err, "unique key appears more than once in batch") + } else { + require.ErrorContains(t, err, "ON CONFLICT DO UPDATE command cannot affect row a second time") + } + require.Empty(t, results) + + jobs, err := exec.JobGetByKindMany(ctx, &riverdriver.JobGetByKindManyParams{ + Kind: []string{job.Kind}, + Schema: schema, + }) + require.NoError(t, err) + require.Empty(t, jobs) + }) + + t.Run("UniqueKeyOutsideEnforcedState", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + jobs := []*riverdriver.JobInsertFastParams{ + { + EncodedArgs: []byte(`{"encoded": "args"}`), + Kind: "test_kind", + MaxAttempts: rivercommon.MaxAttemptsDefault, + Priority: rivercommon.PriorityDefault, + Queue: rivercommon.QueueDefault, + State: rivertype.JobStateAvailable, + Tags: []string{}, + UniqueKey: []byte("unique-key-outside-state"), + UniqueStates: 0x01, + }, + { + EncodedArgs: []byte(`{"encoded": "args"}`), + Kind: "test_kind", + MaxAttempts: rivercommon.MaxAttemptsDefault, + Priority: rivercommon.PriorityDefault, + Queue: rivercommon.QueueDefault, + State: rivertype.JobStateScheduled, + Tags: []string{}, + UniqueKey: []byte("unique-key-outside-state"), + UniqueStates: 0x01, + }, + } + + results, err := exec.JobInsertFastMany(ctx, &riverdriver.JobInsertFastManyParams{Jobs: jobs}) + require.NoError(t, err) + require.Len(t, results, 2) + require.False(t, results[0].UniqueSkippedAsDuplicate) + require.False(t, results[1].UniqueSkippedAsDuplicate) + require.NotEqual(t, results[0].Job.ID, results[1].Job.ID) + }) + t.Run("BinaryNonUTF8UniqueKey", func(t *testing.T) { t.Parallel() diff --git a/riverdriver/riversqlite/river_sqlite_driver.go b/riverdriver/riversqlite/river_sqlite_driver.go index 112f686a0..485a3b661 100644 --- a/riverdriver/riversqlite/river_sqlite_driver.go +++ b/riverdriver/riversqlite/river_sqlite_driver.go @@ -612,12 +612,28 @@ func (e *Executor) JobGetStuck(ctx context.Context, params *riverdriver.JobGetSt func (e *Executor) JobInsertFastMany(ctx context.Context, params *riverdriver.JobInsertFastManyParams) ([]*riverdriver.JobInsertFastResult, error) { // We use a special `(xmax != 0)` trick in Postgres to determine whether an // upserted row was inserted or skipped, but as far as I can find, there's no - // such trick possible in SQLite. Instead, we roll a random nonce and insert - // it to metadata. If the same nonce comes back, we know we really inserted - // the row. If not, we're getting an existing row back. - uniqueNonce := randutil.Hex(8) + // such trick possible in SQLite. Instead, we insert a random nonce into + // each row's metadata. A returned nonce from this batch identifies an + // inserted row; another nonce identifies an existing row. + uniqueNonces := make([]string, len(params.Jobs)) + uniqueNoncesInBatch := make(map[string]bool, len(params.Jobs)) + uniqueKeysInBatch := make(map[string]bool, len(params.Jobs)) + for i, job := range params.Jobs { + // PostgreSQL rejects a statement that affects the same unique row twice. + // SQLite allows it, so reject repeated keys covered by the unique index. + if len(job.UniqueKey) > 0 && job.UniqueStates&uniquestates.UniqueStatesToBitmask([]rivertype.JobState{job.State}) != 0 { + key := string(job.UniqueKey) + if uniqueKeysInBatch[key] { + return nil, errors.New("unique key appears more than once in batch") + } + uniqueKeysInBatch[key] = true + } + + uniqueNonces[i] = randutil.Hex(8) + uniqueNoncesInBatch[uniqueNonces[i]] = true + } - jobsParam, err := sqliteJobInsertFastManyJobsParam(params.Jobs, uniqueNonce) + jobsParam, err := sqliteJobInsertFastManyJobsParam(params.Jobs, uniqueNonces) if err != nil { return nil, err } @@ -632,16 +648,17 @@ func (e *Executor) JobInsertFastMany(ctx context.Context, params *riverdriver.Jo if err != nil { return nil, err } + returnedNonce, hasNonce := riverdriver.UniqueInsertMetadataNonce(job.Metadata) return &riverdriver.JobInsertFastResult{ Job: job, - UniqueSkippedAsDuplicate: riverdriver.UniqueInsertMetadataIsDuplicate(job.Metadata, uniqueNonce), + UniqueSkippedAsDuplicate: !hasNonce || !uniqueNoncesInBatch[returnedNonce], }, nil }) } func (e *Executor) JobInsertFastManyNoReturning(ctx context.Context, params *riverdriver.JobInsertFastManyParams) (int, error) { - jobsParam, err := sqliteJobInsertFastManyJobsParam(params.Jobs, "") + jobsParam, err := sqliteJobInsertFastManyJobsParam(params.Jobs, nil) if err != nil { return 0, err } @@ -1574,14 +1591,14 @@ func durationAsString(duration time.Duration) string { return strconv.FormatFloat(duration.Seconds(), 'f', 3, 64) + " seconds" } -func sqliteJobInsertFastManyJobsParam(jobs []*riverdriver.JobInsertFastParams, uniqueNonce string) ([]byte, error) { +func sqliteJobInsertFastManyJobsParam(jobs []*riverdriver.JobInsertFastParams, uniqueNonces []string) ([]byte, error) { jobsParam := make([]map[string]any, len(jobs)) for i, job := range jobs { metadata := sliceutil.FirstNonEmpty(job.Metadata, []byte("{}")) - if uniqueNonce != "" { + if uniqueNonces != nil { var err error - metadata, err = riverdriver.UniqueInsertMetadataWithNonce(metadata, uniqueNonce) + metadata, err = riverdriver.UniqueInsertMetadataWithNonce(metadata, uniqueNonces[i]) if err != nil { return nil, err } diff --git a/riverdriver/unique_insert.go b/riverdriver/unique_insert.go index bdeaa07d6..c03db5c9e 100644 --- a/riverdriver/unique_insert.go +++ b/riverdriver/unique_insert.go @@ -56,16 +56,23 @@ func (m UniqueInsertMode) SQL() string { // from a proposed insert, indicating that an existing row was returned // instead. func UniqueInsertMetadataIsDuplicate(metadata []byte, nonce string) bool { + metadataNonce, ok := UniqueInsertMetadataNonce(metadata) + return !ok || metadataNonce != nonce +} + +// UniqueInsertMetadataNonce returns the nonce in metadata, or false if the +// metadata has no valid nonce. +func UniqueInsertMetadataNonce(metadata []byte) (string, bool) { var metadataMap map[string]json.RawMessage if err := json.Unmarshal(metadata, &metadataMap); err != nil { - return true + return "", false } var metadataNonce string if err := json.Unmarshal(metadataMap[UniqueInsertMetadataKey], &metadataNonce); err != nil { - return true + return "", false } - return metadataNonce != nonce + return metadataNonce, true } // UniqueInsertMetadataWithNonce returns metadata with nonce set under diff --git a/riverdriver/unique_insert_test.go b/riverdriver/unique_insert_test.go index 3c78de985..0a764ddf5 100644 --- a/riverdriver/unique_insert_test.go +++ b/riverdriver/unique_insert_test.go @@ -34,6 +34,34 @@ func TestUniqueInsertMetadataIsDuplicate(t *testing.T) { }) } +func TestUniqueInsertMetadataNonce(t *testing.T) { + t.Parallel() + + t.Run("InvalidMetadata", func(t *testing.T) { + t.Parallel() + + nonce, ok := UniqueInsertMetadataNonce([]byte(`{`)) + require.False(t, ok) + require.Empty(t, nonce) + }) + + t.Run("MatchingNonce", func(t *testing.T) { + t.Parallel() + + nonce, ok := UniqueInsertMetadataNonce([]byte(`{"river:unique_nonce":"nonce"}`)) + require.True(t, ok) + require.Equal(t, "nonce", nonce) + }) + + t.Run("MissingNonce", func(t *testing.T) { + t.Parallel() + + nonce, ok := UniqueInsertMetadataNonce([]byte(`{"existing":123}`)) + require.False(t, ok) + require.Empty(t, nonce) + }) +} + func TestUniqueInsertMetadataWithNonce(t *testing.T) { t.Parallel()