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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### 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).

### Fixed

- 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).
Expand Down
21 changes: 12 additions & 9 deletions client.go
Original file line number Diff line number Diff line change
Expand Up @@ -1773,15 +1773,6 @@ func insertParamsFromConfigArgsAndOptions(archetype *baseservice.Archetype, conf
State: rivertype.JobStateAvailable,
Tags: tags,
}
if !uniqueOpts.isEmpty() {
internalUniqueOpts := (*dbunique.UniqueOpts)(&uniqueOpts)
insertParams.UniqueKey, err = dbunique.UniqueKey(archetype.Time, internalUniqueOpts, insertParams)
if err != nil {
return nil, err
}
insertParams.UniqueStates = internalUniqueOpts.StateBitmask()
}

switch {
case !insertOpts.ScheduledAt.IsZero():
insertParams.ScheduledAt = &insertOpts.ScheduledAt
Expand All @@ -1799,6 +1790,18 @@ func insertParamsFromConfigArgsAndOptions(archetype *baseservice.Archetype, conf
insertParams.State = rivertype.JobStatePending
}

// Compute the unique key only after the effective scheduled time is known
// so that a ByPeriod key describes the period in which the job is
// scheduled to run rather than the period in which it was inserted.
if !uniqueOpts.isEmpty() {
internalUniqueOpts := (*dbunique.UniqueOpts)(&uniqueOpts)
insertParams.UniqueKey, err = dbunique.UniqueKey(archetype.Time, internalUniqueOpts, insertParams)
if err != nil {
return nil, err
}
insertParams.UniqueStates = internalUniqueOpts.StateBitmask()
}

return insertParams, nil
}

Expand Down
145 changes: 145 additions & 0 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package river

import (
"context"
"crypto/sha256"
"encoding/json"
"errors"
"fmt"
Expand Down Expand Up @@ -9531,6 +9532,75 @@ func TestInsertParamsFromJobArgsAndOptions(t *testing.T) {
require.Equal(t, internalUniqueOpts.StateBitmask(), params.UniqueStates)
})

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

var (
chicago = time.FixedZone("CDT", -5*60*60)
now = time.Date(2026, time.August, 12, 12, 34, 56, 0, time.UTC)
uniqueOpts = UniqueOpts{ByPeriod: time.Hour}
wantKey = sha256.Sum256([]byte("&kind=noOp&period=2026-08-12T12:00:00Z"))
)

utcArchetype := riversharedtest.BaseServiceArchetype(t)
utcArchetype.Time.StubNow(now)

utcParams, err := insertParamsFromConfigArgsAndOptions(utcArchetype, config, noOpArgs{}, &InsertOpts{UniqueOpts: uniqueOpts})
require.NoError(t, err)
require.Equal(t, wantKey[:], utcParams.UniqueKey)

zonedArchetype := riversharedtest.BaseServiceArchetype(t)
zonedArchetype.Time.StubNow(now.In(chicago))

zonedParams, err := insertParamsFromConfigArgsAndOptions(zonedArchetype, config, noOpArgs{}, &InsertOpts{UniqueOpts: uniqueOpts})
require.NoError(t, err)
require.Equal(t, wantKey[:], zonedParams.UniqueKey)

scheduledParams, err := insertParamsFromConfigArgsAndOptions(utcArchetype, config, noOpArgs{}, &InsertOpts{
ScheduledAt: now.Add(10 * time.Minute).In(chicago),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.Equal(t, wantKey[:], scheduledParams.UniqueKey)
})

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

archetype := riversharedtest.BaseServiceArchetype(t)
now := time.Date(2026, time.August, 12, 12, 34, 56, 0, time.UTC)
archetype.Time.StubNow(now)
uniqueOpts := UniqueOpts{ByPeriod: 24 * time.Hour}

params, err := insertParamsFromConfigArgsAndOptions(archetype, config, noOpArgs{}, &InsertOpts{
ScheduledAt: now.Add(48*time.Hour + 5*time.Minute),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)

// The period comes from the scheduled time rather than insertion time.
wantKey := sha256.Sum256([]byte("&kind=noOp&period=2026-08-14T00:00:00Z"))
require.Equal(t, wantKey[:], params.UniqueKey)

// A job scheduled later in the same period gets the same key even when
// it's inserted at a different time.
archetype.Time.StubNow(now.Add(20 * time.Hour))
laterParams, err := insertParamsFromConfigArgsAndOptions(archetype, config, noOpArgs{}, &InsertOpts{
ScheduledAt: now.Add(48*time.Hour + 9*time.Hour),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.Equal(t, wantKey[:], laterParams.UniqueKey)

// An unscheduled job uses the insertion time's period.
unscheduledParams, err := insertParamsFromConfigArgsAndOptions(archetype, config, noOpArgs{}, &InsertOpts{
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
unscheduledKey := sha256.Sum256([]byte("&kind=noOp&period=2026-08-13T00:00:00Z"))
require.Equal(t, unscheduledKey[:], unscheduledParams.UniqueKey)
})

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

Expand Down Expand Up @@ -9850,6 +9920,81 @@ func TestUniqueOpts(t *testing.T) {
require.Equal(t, insertRes0.Job.ID, insertRes1.Job.ID)
})

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

client, _ := setup(t)

var (
chicago = time.FixedZone("CDT", -5*60*60)
scheduledAt = client.baseService.Time.Now().Add(48 * time.Hour).Truncate(24 * time.Hour).Add(9 * time.Hour)
uniqueOpts = UniqueOpts{ByPeriod: 24 * time.Hour}
)

insertRes0, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{
ScheduledAt: scheduledAt.UTC(),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.False(t, insertRes0.UniqueSkippedAsDuplicate)

// The same UTC period expressed in another zone is a duplicate, and so
// is an insert from a process whose clock is in another zone.
client.baseService.Time.StubNow(client.baseService.Time.Now().In(chicago))
insertRes1, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{
ScheduledAt: scheduledAt.Add(time.Hour).In(chicago),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.True(t, insertRes1.UniqueSkippedAsDuplicate)
require.Equal(t, insertRes0.Job.ID, insertRes1.Job.ID)
})

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

client, _ := setup(t)

var (
now = client.baseService.Time.Now()
scheduledAt = now.Add(48 * time.Hour).Truncate(24 * time.Hour).Add(9 * time.Hour)
uniqueOpts = UniqueOpts{ByPeriod: 24 * time.Hour}
)

insertRes0, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{
ScheduledAt: scheduledAt,
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.False(t, insertRes0.UniqueSkippedAsDuplicate)

// Inserting at a different time for the same scheduled period is a
// duplicate even though the insertion times fall in different periods.
client.baseService.Time.StubNow(now.Add(24 * time.Hour))
insertRes1, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{
ScheduledAt: scheduledAt.Add(6 * time.Hour),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.True(t, insertRes1.UniqueSkippedAsDuplicate)
require.Equal(t, insertRes0.Job.ID, insertRes1.Job.ID)

// Neither a job scheduled in the next period nor an unscheduled job in
// the current period conflicts with the scheduled job.
insertRes2, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{
ScheduledAt: scheduledAt.Add(24 * time.Hour),
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.False(t, insertRes2.UniqueSkippedAsDuplicate)

insertRes3, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{
UniqueOpts: uniqueOpts,
})
require.NoError(t, err)
require.False(t, insertRes3.UniqueSkippedAsDuplicate)
})

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

Expand Down
13 changes: 9 additions & 4 deletions insert_opts.go
Original file line number Diff line number Diff line change
Expand Up @@ -158,10 +158,15 @@ type UniqueOpts struct {
// }
ByArgs bool

// ByPeriod defines uniqueness within a given period. On an insert time is
// rounded down to the nearest multiple of the given period, and a job is
// only inserted if there isn't an existing job that will run between then
// and the next multiple of the period.
// ByPeriod defines uniqueness within a given period. On an insert, the
// job's scheduled time (ScheduledAt, or the current time for jobs that
// aren't scheduled) is rounded down to the nearest multiple of the given
// period, and a job is only inserted if there isn't an existing job that
// will run between then and the next multiple of the period.
//
// Periods are measured in UTC, so the same period produces the same
// unique key regardless of the time zone of the inserting process or of
// a provided ScheduledAt.
//
// Default is no unique period, meaning that as long as any other unique
// property is enabled, uniqueness will be enforced across all jobs of the
Expand Down
7 changes: 6 additions & 1 deletion internal/dbunique/db_unique.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,12 @@ func buildUniqueKeyString(timeGen rivertype.TimeGenerator, uniqueOpts *UniqueOpt
}

if uniqueOpts.ByPeriod != time.Duration(0) {
lowerPeriodBound := ptrutil.ValOrDefaultFunc(params.ScheduledAt, timeGen.Now).Truncate(uniqueOpts.ByPeriod)
// Format the period bound in UTC. RFC 3339 formatting includes the
// time's location, so without normalization the same period produces
// different keys for processes or callers in different time zones.
// Truncation operates on absolute time, so the bound itself is
// independent of location.
lowerPeriodBound := ptrutil.ValOrDefaultFunc(params.ScheduledAt, timeGen.Now).Truncate(uniqueOpts.ByPeriod).UTC()
sb.WriteString("&period=" + lowerPeriodBound.Format(time.RFC3339))
}

Expand Down
86 changes: 86 additions & 0 deletions internal/dbunique/db_unique_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -519,6 +519,92 @@ func TestUniqueKey(t *testing.T) {
}
}

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

var (
// A fixed zone avoids depending on the host's time zone database.
chicago = time.FixedZone("CDT", -5*60*60)
nowUTC = time.Date(2026, time.September, 24, 12, 34, 56, 0, time.UTC)
uniqueOpts = &UniqueOpts{ByPeriod: time.Hour}
wantPreHash = "&kind=worker_1&period=2026-09-24T12:00:00Z"
)

type testBundle struct {
params *rivertype.JobInsertParams
timeGen *riversharedtest.TimeStub
}

setup := func(t *testing.T) *testBundle {
t.Helper()

timeGen := &riversharedtest.TimeStub{}
timeGen.StubNow(nowUTC)

return &testBundle{
params: &rivertype.JobInsertParams{
Args: JobArgsStaticKind{kind: "worker_1"},
EncodedArgs: []byte(`{}`),
Kind: "worker_1",
Queue: "default",
},
timeGen: timeGen,
}
}

requirePreHash := func(t *testing.T, bundle *testBundle, want string) {
t.Helper()

preHash, err := buildUniqueKeyString(bundle.timeGen, uniqueOpts, bundle.params)
require.NoError(t, err)
require.Equal(t, want, preHash)

wantKey := sha256.Sum256([]byte(want))
uniqueKey, err := UniqueKey(bundle.timeGen, uniqueOpts, bundle.params)
require.NoError(t, err)
require.Equal(t, wantKey[:], uniqueKey)
}

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

bundle := setup(t)
bundle.timeGen.StubNow(nowUTC.In(chicago))

requirePreHash(t, bundle, wantPreHash)
})

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

bundle := setup(t)
bundle.params.ScheduledAt = new(nowUTC.In(chicago))

requirePreHash(t, bundle, wantPreHash)
})

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

bundle := setup(t)

// 06:59:59 in a UTC-05:00 zone is 11:59:59 UTC, so the job belongs to
// the 11:00 UTC period rather than to a period derived from its local
// wall-clock representation.
bundle.params.ScheduledAt = new(time.Date(2026, time.September, 24, 6, 59, 59, 0, chicago))

requirePreHash(t, bundle, "&kind=worker_1&period=2026-09-24T11:00:00Z")
})

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

bundle := setup(t)

requirePreHash(t, bundle, wantPreHash)
})
}

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

Expand Down
4 changes: 3 additions & 1 deletion riverdriver/riverdrivertest/benchmark.go
Original file line number Diff line number Diff line change
Expand Up @@ -216,7 +216,9 @@ func Benchmark[TTx any](ctx context.Context, b *testing.B,
for _, tc := range testCases {
b.Run(tc.name, func(b *testing.B) {
largeMetadata := makeBenchmarkMetadataWithRiverLogSize(tc.metadataSizeBytes)
now := time.Now().UTC()
// Match PostgreSQL's timestamp precision so the inserted
// scheduled_at can't round past the fetch cutoff.
now := time.Now().UTC().Truncate(time.Microsecond)

insertedJobs, err := exec.JobInsertFullMany(ctx, &riverdriver.JobInsertFullManyParams{
Jobs: []*riverdriver.JobInsertFullParams{
Expand Down
Loading
Loading