From 4ff1b65c595f8de75acb1ef67a3aff2d2a6754e7 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Thu, 24 Sep 2026 08:29:18 -0500 Subject: [PATCH 1/2] derive unique periods from UTC scheduled time `UniqueOpts.ByPeriod` keys are computed before an insert's effective `ScheduledAt` is applied, so a job scheduled for tomorrow is keyed by today's period. Re-enqueueing "tomorrow at 9am" at different times of day therefore inserts duplicates, and the key disagrees with the period the job actually runs in. The period bound is also formatted in the time's own location. A process whose local time zone isn't UTC, or a caller passing a `ScheduledAt` in another zone, embeds an offset like `-05:00` in the key, so two processes in different zones compute different keys for the same period and both insert. Compute the unique key after the effective scheduled time is known, and normalize the period bound to UTC before formatting it. Truncation already operates on absolute time, so only the textual form changes. Document both behaviors on `ByPeriod` and add a changelog entry that notes the one-time key change for scheduled jobs and non-UTC processes during a rolling upgrade. Cover non-UTC clocks, non-UTC `ScheduledAt`, and zone-independent period boundaries in key construction; scheduled jobs deduplicating by scheduled period across insertion times while unscheduled and next-period jobs don't conflict; and inserts in different zones deduplicating end to end across every driver. --- CHANGELOG.md | 4 + client.go | 21 +-- client_test.go | 145 ++++++++++++++++++ insert_opts.go | 13 +- internal/dbunique/db_unique.go | 7 +- internal/dbunique/db_unique_test.go | 86 +++++++++++ .../riverdrivertest/driver_client_test.go | 47 ++++++ 7 files changed, 309 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 598b28191..4fd232894 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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). diff --git a/client.go b/client.go index bf07bfc0f..f1dce2b27 100644 --- a/client.go +++ b/client.go @@ -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 @@ -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 } diff --git a/client_test.go b/client_test.go index bb7ba8150..98b44a4ef 100644 --- a/client_test.go +++ b/client_test.go @@ -2,6 +2,7 @@ package river import ( "context" + "crypto/sha256" "encoding/json" "errors" "fmt" @@ -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() @@ -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() diff --git a/insert_opts.go b/insert_opts.go index 64fb770d5..d726fae40 100644 --- a/insert_opts.go +++ b/insert_opts.go @@ -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 diff --git a/internal/dbunique/db_unique.go b/internal/dbunique/db_unique.go index bbf486616..cf6037316 100644 --- a/internal/dbunique/db_unique.go +++ b/internal/dbunique/db_unique.go @@ -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)) } diff --git a/internal/dbunique/db_unique_test.go b/internal/dbunique/db_unique_test.go index 66c67566b..8952df91a 100644 --- a/internal/dbunique/db_unique_test.go +++ b/internal/dbunique/db_unique_test.go @@ -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() diff --git a/riverdriver/riverdrivertest/driver_client_test.go b/riverdriver/riverdrivertest/driver_client_test.go index 1e361a33c..f1d720ebb 100644 --- a/riverdriver/riverdrivertest/driver_client_test.go +++ b/riverdriver/riverdrivertest/driver_client_test.go @@ -418,6 +418,53 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.Equal(t, rivertype.JobStateCancelled, event.Job.State) }) + t.Run("InsertUniqueByPeriod", func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + type JobArgs struct { + testutil.JobArgsReflectKind[JobArgs] + } + + river.AddWorker(bundle.config.Workers, river.WorkFunc(func(ctx context.Context, job *river.Job[JobArgs]) error { + return nil + })) + + var ( + chicago = time.FixedZone("CDT", -5*60*60) + + // Far enough in the future that the test can't cross a period + // boundary while it runs. + scheduledAt = time.Now().UTC().Add(72 * time.Hour).Truncate(24 * time.Hour).Add(9 * time.Hour) + uniqueOpts = river.UniqueOpts{ByPeriod: 24 * time.Hour} + ) + + insertRes0, err := client.Insert(ctx, &JobArgs{}, &river.InsertOpts{ + ScheduledAt: scheduledAt.In(chicago), + UniqueOpts: uniqueOpts, + }) + require.NoError(t, err) + require.False(t, insertRes0.UniqueSkippedAsDuplicate) + + // Same UTC period expressed in UTC is a duplicate. + insertRes1, err := client.Insert(ctx, &JobArgs{}, &river.InsertOpts{ + ScheduledAt: scheduledAt.Add(10 * time.Hour), + UniqueOpts: uniqueOpts, + }) + require.NoError(t, err) + require.True(t, insertRes1.UniqueSkippedAsDuplicate) + require.Equal(t, insertRes0.Job.ID, insertRes1.Job.ID) + + // The next period is not. + insertRes2, err := client.Insert(ctx, &JobArgs{}, &river.InsertOpts{ + ScheduledAt: scheduledAt.Add(24 * time.Hour), + UniqueOpts: uniqueOpts, + }) + require.NoError(t, err) + require.False(t, insertRes2.UniqueSkippedAsDuplicate) + }) + t.Run("JobDelete", func(t *testing.T) { t.Parallel() From 46a692b7b4b30ff216c85565a1f478ea469606dd Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Thu, 24 Sep 2026 10:49:57 -0500 Subject: [PATCH 2/2] align benchmark cutoff with database precision `JobGetAvailable_LargeMetadata` uses the same nanosecond timestamp for insertion and fetch. PostgreSQL stores timestamps at microsecond precision, so rounding can place the scheduled time after the fetch cutoff and make the benchmark intermittently fetch no job. Truncate the benchmark timestamp to microseconds before insertion so the stored scheduled time matches the fetch cutoff exactly. --- riverdriver/riverdrivertest/benchmark.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/riverdriver/riverdrivertest/benchmark.go b/riverdriver/riverdrivertest/benchmark.go index 8c4d0c8b7..d0118d3e4 100644 --- a/riverdriver/riverdrivertest/benchmark.go +++ b/riverdriver/riverdrivertest/benchmark.go @@ -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{