From 9531245d307b4c3ff00687c0da0c20b1e23688a3 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Thu, 24 Sep 2026 12:16:46 -0500 Subject: [PATCH] number `river bench` jobs sequentially `river bench` assigns every inserted job `num: 0` because its insert loops update copies of the args values. This hides each job's intended sequence number in both fixed-count and continuous modes. Write a numbered `BenchmarkArgs` directly into each reused insert parameter in both loops. Carry the counter between batches, and trim the final fixed-count batch before numbering so only prepared jobs advance it. Check persisted args across a real batch boundary, including a short final batch. The assertions report the first incorrect number without dumping thousands of values. --- CHANGELOG.md | 1 + cmd/river/riverbench/river_bench.go | 23 ++----- cmd/river/riverbench/river_bench_test.go | 86 ++++++++++++++++++++++++ 3 files changed, 93 insertions(+), 17 deletions(-) create mode 100644 cmd/river/riverbench/river_bench_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 7fc147770..8d314ebc0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -34,6 +34,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Fixed `JobRescuer` overwriting jobs that complete, leave the running state, or are claimed again by another worker after being fetched for rescue, preserving their state, errors, metadata, and timestamps across PostgreSQL and SQLite drivers. Fixes [#1302](https://github.com/riverqueue/river/issues/1302). [PR #1373](https://github.com/riverqueue/river/pull/1373). - Fixed SQLite notification listeners delivering notifications from before a subscription or from an unsubscribe gap. Notification reads now fetch subscribed topics in bounded batches, and cleanup deletes expired notifications in batches of 10,000 rows (reduced to 1,000 after repeated timeouts), with pauses between batches to reduce write lock contention. [PR #1381](https://github.com/riverqueue/river/pull/1381). - Fixed the `Job appears to be stuck` log line reporting the client-level `JobTimeout` instead of the worker-level timeout when a worker overrides `Timeout`. [PR #1394](https://github.com/riverqueue/river/pull/1394). +- Fixed `river bench` inserting every benchmark job with a `num` arg of `0` instead of numbering jobs sequentially. [PR #1379](https://github.com/riverqueue/river/pull/1379). ## [0.47.0] - 2026-09-01 diff --git a/cmd/river/riverbench/river_bench.go b/cmd/river/riverbench/river_bench.go index e5a0d868a..476fa4b79 100644 --- a/cmd/river/riverbench/river_bench.go +++ b/cmd/river/riverbench/river_bench.go @@ -347,7 +347,6 @@ func (b *Benchmarker[TTx]) insertJobs( // We'll be reusing the same batch for all inserts because (1) we can // get away with it, and (2) to avoid needless allocations. insertParamsBatch = make([]river.InsertManyParams, insertBatchSize) - jobArgsBatch = make([]BenchmarkArgs, insertBatchSize) jobNum int ) @@ -355,19 +354,14 @@ func (b *Benchmarker[TTx]) insertJobs( var numInsertedThisRound int for { - for _, jobArgs := range jobArgsBatch { - jobNum++ - jobArgs.Num = jobNum - } - - for i := range insertParamsBatch { - insertParamsBatch[i].Args = jobArgsBatch[i] - } - numLeft := numTotalJobs - numInsertedThisRound if numLeft < insertBatchSize { insertParamsBatch = insertParamsBatch[0:numLeft] } + for i := range insertParamsBatch { + jobNum++ + insertParamsBatch[i].Args = BenchmarkArgs{Num: jobNum} + } start := time.Now() if _, err := client.InsertMany(ctx, insertParamsBatch); err != nil { @@ -411,7 +405,6 @@ func (b *Benchmarker[TTx]) insertJobsContinuously( // We'll be reusing the same batch for all inserts because (1) we can // get away with it, and (2) to avoid needless allocations. insertParamsBatch = make([]river.InsertManyParams, insertBatchSize) - jobArgsBatch = make([]BenchmarkArgs, insertBatchSize) jobNum int ) @@ -430,13 +423,9 @@ func (b *Benchmarker[TTx]) insertJobsContinuously( var numInsertedThisRound int for { - for _, jobArgs := range jobArgsBatch { - jobNum++ - jobArgs.Num = jobNum - } - for i := range insertParamsBatch { - insertParamsBatch[i].Args = jobArgsBatch[i] + jobNum++ + insertParamsBatch[i].Args = BenchmarkArgs{Num: jobNum} } if _, err := client.InsertMany(ctx, insertParamsBatch); err != nil { diff --git a/cmd/river/riverbench/river_bench_test.go b/cmd/river/riverbench/river_bench_test.go new file mode 100644 index 000000000..b8fbca979 --- /dev/null +++ b/cmd/river/riverbench/river_bench_test.go @@ -0,0 +1,86 @@ +package riverbench + +import ( + "context" + "encoding/json" + "slices" + "sync/atomic" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/require" + + "github.com/riverqueue/river" + "github.com/riverqueue/river/riverdbtest" + "github.com/riverqueue/river/riverdriver/riverpgxv5" + "github.com/riverqueue/river/rivershared/riversharedtest" +) + +func TestBenchmarkerInsertJobs(t *testing.T) { + t.Parallel() + + ctx := context.Background() + + type testBundle struct { + benchmarker *Benchmarker[pgx.Tx] + client *river.Client[pgx.Tx] + } + + setup := func(ctx context.Context, t *testing.T) *testBundle { + t.Helper() + + var ( + driver = riverpgxv5.New(riversharedtest.DBPool(ctx, t)) + logger = riversharedtest.Logger(t) + schema = riverdbtest.TestSchema(ctx, t, driver, nil) + ) + + // Insert-only client since jobs should stay in the database so that + // their args can be inspected. + client, err := river.NewClient(driver, &river.Config{ + Logger: logger, + Schema: schema, + }) + require.NoError(t, err) + + return &testBundle{ + benchmarker: NewBenchmarker(driver, &Config{Logger: logger, Schema: schema}), + client: client, + } + } + + t.Run("NumbersJobsSequentiallyAcrossBatches", func(t *testing.T) { + t.Parallel() + + bundle := setup(ctx, t) + + // Read the persisted args across a batch boundary. This catches + // changes to a copy of the args that never reach InsertMany. + const numTotalJobs = insertBatchSize + 2 + + var ( + minJobsReady = make(chan struct{}) + numJobsInserted atomic.Int64 + numJobsLeft atomic.Int64 + ) + + bundle.benchmarker.insertJobs(ctx, bundle.client, minJobsReady, &numJobsInserted, &numJobsLeft, numTotalJobs, make(chan struct{})) + require.Equal(t, int64(numTotalJobs), numJobsInserted.Load()) + + listRes, err := bundle.client.JobList(ctx, river.NewJobListParams().Kinds((BenchmarkArgs{}).Kind()).First(numTotalJobs)) + require.NoError(t, err) + require.Len(t, listRes.Jobs, numTotalJobs) + + jobNums := make([]int, len(listRes.Jobs)) + for i, job := range listRes.Jobs { + var args BenchmarkArgs + require.NoError(t, json.Unmarshal(job.EncodedArgs, &args)) + jobNums[i] = args.Num + } + slices.Sort(jobNums) + + for i, jobNum := range jobNums { + require.Equal(t, i+1, jobNum, "job number at sorted index %d", i) + } + }) +}