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

Expand Down
37 changes: 27 additions & 10 deletions riverdriver/riversqlite/river_sqlite_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand All @@ -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
}
Expand Down Expand Up @@ -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
}
Expand Down
13 changes: 10 additions & 3 deletions riverdriver/unique_insert.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 28 additions & 0 deletions riverdriver/unique_insert_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
Loading