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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
23 changes: 6 additions & 17 deletions cmd/river/riverbench/river_bench.go
Original file line number Diff line number Diff line change
Expand Up @@ -347,27 +347,21 @@ 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
)

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 {
Expand Down Expand Up @@ -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
)
Expand All @@ -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 {
Expand Down
86 changes: 86 additions & 0 deletions cmd/river/riverbench/river_bench_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
})
}
Loading