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) + } + }) +}