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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- If you're on Postgres, you can ignore it with no adverse affects.
- If you're on SQLite, it rebuilds `river_job` to add an `AUTOINCREMENT` keyword to the primary key, preventing a possible edge case where generated job IDs could be reused after deletion. It's not necessary to run the migration for River to work, but it's a good idea to get it in when convenient. [PR #1390](https://github.com/riverqueue/river/pull/1390).

⚠️ **Breaking behavior change:** `UniqueOpts{ExcludeKind: true}` with no other unique fields set was previously a silent no-op; it's now rejected at insert time with an explicit error. [PR #1404](https://github.com/riverqueue/river/pull/1404).

### Added

- Added `Config.FetchOnlyKnownKinds` to restrict job fetching to registered worker kinds, including aliases. Clients with different workers can share a queue while leaving unknown jobs available without consuming attempts. Disabled by default; leader election and stuck-job rescue behavior are unchanged. [PR #1396](https://github.com/riverqueue/river/pull/1396).
Expand All @@ -21,6 +23,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed

- `UniqueOpts.ByPeriod` now derives a job's period from its effective scheduled time (`InsertOpts.ScheduledAt` when set, otherwise the insertion time), so scheduled jobs are deduplicated against other jobs scheduled in the same period rather than against jobs inserted in the same period. Periods are also now always measured in UTC, so processes and `ScheduledAt` values in different time zones produce the same unique key for the same period. Unique keys for scheduled `ByPeriod` jobs, and for any `ByPeriod` job inserted from a process whose local time zone isn't UTC, differ from those produced by previous versions. During a rolling upgrade, old and new clients may therefore each insert one job for such a period; jobs that aren't scheduled and are inserted from UTC processes are unaffected. [PR #1377](https://github.com/riverqueue/river/pull/1377).
- `UniqueOpts{ExcludeKind: true}` alone is now rejected at insert time instead of being silently ignored. `UniqueOpts.isEmpty()` now considers `ExcludeKind`, in line with the equivalent handling in internal/dbunique. [PR #1404](https://github.com/riverqueue/river/pull/1404).

### Fixed

Expand Down
34 changes: 34 additions & 0 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9963,6 +9963,40 @@ func TestInsertParamsFromJobArgsAndOptions(t *testing.T) {
require.EqualError(t, err, "UniqueOpts.ByPeriod should not be less than 1 second")
require.Nil(t, insertParams)
})

t.Run("UniqueOptsExcludeKindOnlyValidated", func(t *testing.T) {
t.Parallel()

insertParams, err := insertParamsFromConfigArgsAndOptions(
archetype,
config,
noOpArgs{},
&InsertOpts{UniqueOpts: UniqueOpts{ExcludeKind: true}},
)
require.EqualError(t, err, "UniqueOpts.ExcludeKind requires ByArgs, ByQueue, or ByPeriod")
require.Nil(t, insertParams)
})

t.Run("UniqueOptsExcludeKindRemovesKindFromKey", func(t *testing.T) {
t.Parallel()

argsKindA := JobArgsStaticKind{kind: "kind_a"}
argsKindB := JobArgsStaticKind{kind: "kind_b"}

// With ExcludeKind, two different kinds with identical encoded args share a key.
paramsA, err := insertParamsFromConfigArgsAndOptions(archetype, config, argsKindA, &InsertOpts{UniqueOpts: UniqueOpts{ByArgs: true, ExcludeKind: true}})
require.NoError(t, err)
paramsB, err := insertParamsFromConfigArgsAndOptions(archetype, config, argsKindB, &InsertOpts{UniqueOpts: UniqueOpts{ByArgs: true, ExcludeKind: true}})
require.NoError(t, err)
require.Equal(t, paramsA.UniqueKey, paramsB.UniqueKey, "unique keys should be identical across kinds with ExcludeKind")

paramsAWithKind, err := insertParamsFromConfigArgsAndOptions(archetype, config, argsKindA, &InsertOpts{UniqueOpts: UniqueOpts{ByArgs: true}})
require.NoError(t, err)
paramsBWithKind, err := insertParamsFromConfigArgsAndOptions(archetype, config, argsKindB, &InsertOpts{UniqueOpts: UniqueOpts{ByArgs: true}})
require.NoError(t, err)
require.NotEqual(t, paramsAWithKind.UniqueKey, paramsBWithKind.UniqueKey)
require.NotEqual(t, paramsA.UniqueKey, paramsAWithKind.UniqueKey)
})
}

func TestID(t *testing.T) {
Expand Down
11 changes: 10 additions & 1 deletion insert_opts.go
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,10 @@ type UniqueOpts struct {
// ExcludeKind indicates that the job kind should not be included in the
// uniqueness check. This is useful when you want to enforce uniqueness
// across all jobs regardless of kind.
//
// Must be combined with at least one of ByArgs, ByQueue, or ByPeriod;
// otherwise the key would be constant across the entire table and the
// combination is rejected by validate.
ExcludeKind bool
}

Expand All @@ -232,7 +236,8 @@ func (o *UniqueOpts) isEmpty() bool {
return !o.ByArgs &&
o.ByPeriod == time.Duration(0) &&
!o.ByQueue &&
o.ByState == nil
o.ByState == nil &&
!o.ExcludeKind
}

var jobStateAll = rivertype.JobStates() //nolint:gochecknoglobals
Expand Down Expand Up @@ -262,6 +267,10 @@ func (o *UniqueOpts) validate() error {
return errors.New("UniqueOpts.ByPeriod should not be less than 1 second")
}

if o.ExcludeKind && !o.ByArgs && !o.ByQueue && o.ByPeriod == 0 {
return errors.New("UniqueOpts.ExcludeKind requires ByArgs, ByQueue, or ByPeriod")
}

// Job states are typed, but since the underlying type is a string, users
// can put anything they want in there.
for _, state := range o.ByState {
Expand Down
16 changes: 16 additions & 0 deletions insert_opts_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,3 +76,19 @@ func TestUniqueOpts_validate(t *testing.T) {

require.NoError(t, (&UniqueOpts{ByState: rivertype.JobStates()}).validate())
}

func TestUniqueOpts_validateExcludeKind(t *testing.T) {
t.Parallel()

require.EqualError(t, (&UniqueOpts{ExcludeKind: true}).validate(),
"UniqueOpts.ExcludeKind requires ByArgs, ByQueue, or ByPeriod")

require.EqualError(t, (&UniqueOpts{ByState: rivertype.UniqueOptsByStateDefault(), ExcludeKind: true}).validate(),
"UniqueOpts.ExcludeKind requires ByArgs, ByQueue, or ByPeriod")

require.NoError(t, (&UniqueOpts{ByArgs: true, ExcludeKind: true}).validate())
require.NoError(t, (&UniqueOpts{ByQueue: true, ExcludeKind: true}).validate())
require.NoError(t, (&UniqueOpts{ByPeriod: 10 * time.Second, ExcludeKind: true}).validate())
require.NoError(t, (&UniqueOpts{ByArgs: true}).validate())
require.NoError(t, (&UniqueOpts{}).validate())
}
31 changes: 31 additions & 0 deletions job_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (

"github.com/stretchr/testify/require"

"github.com/riverqueue/river/internal/dbunique"
"github.com/riverqueue/river/rivertype"
)

Expand All @@ -17,4 +18,34 @@ func TestUniqueOpts_isEmpty(t *testing.T) {
require.False(t, (&UniqueOpts{ByPeriod: 1 * time.Nanosecond}).isEmpty())
require.False(t, (&UniqueOpts{ByQueue: true}).isEmpty())
require.False(t, (&UniqueOpts{ByState: []rivertype.JobState{rivertype.JobStateAvailable}}).isEmpty())

require.False(t, (&UniqueOpts{ExcludeKind: true}).isEmpty())

states := []rivertype.JobState{
rivertype.JobStateAvailable,
rivertype.JobStatePending,
rivertype.JobStateRunning,
rivertype.JobStateScheduled,
}
for _, byArgs := range []bool{false, true} {
for _, byQueue := range []bool{false, true} {
for _, byPeriod := range []time.Duration{0, 10 * time.Second} {
for _, byState := range [][]rivertype.JobState{nil, states} {
for _, excludeKind := range []bool{false, true} {
opts := UniqueOpts{
ByArgs: byArgs,
ByPeriod: byPeriod,
ByQueue: byQueue,
ByState: byState,
ExcludeKind: excludeKind,
}
internalOpts := (*dbunique.UniqueOpts)(&opts)

require.Equal(t, internalOpts.IsEmpty(), opts.isEmpty(),
"isEmpty and internal IsEmpty should agree for opts %+v", opts)
}
}
}
}
}
}
Loading