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()