From e50368aa8402c2aac3f9a1f3ef3908c84b890b01 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Wed, 9 Sep 2026 11:15:30 -0500 Subject: [PATCH 1/2] speed up single-state finalized job lists Listing one finalized state by time can scan and sort the entire job table. PostgreSQL cannot use the partial finalized-time index without an explicit non-null predicate, and singleton ANY prevents it from using the index's time ordering. Build state equality with the list API's existing condition mechanism when custom SQL cannot depend on the existing array argument. Add `finalized_at IS NOT NULL` only for one known finalized state ordered by finalized time. Keep the shared builder and drivers unchanged. Preserve custom SQL, multi-state filters, and deletion queries. Build cursor predicates locally so conversion preserves the caller's conditions. Cover finalized states, both ordering directions, tied timestamps, PostgreSQL pagination, and combined filters across supported drivers. Check custom OR expressions, contradictory conditions, argument binding, and unchanged delete-many query generation. --- CHANGELOG.md | 4 + delete_many_params_test.go | 20 ++ job_list_params.go | 43 ++- job_list_params_test.go | 233 ++++++++++++- .../riverdrivertest/driver_client_test.go | 306 ++++++++++++++++++ 5 files changed, 599 insertions(+), 7 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b875f8a0..43b81d6bc 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] +### Fixed + +- Improved PostgreSQL job listing performance when filtering by one finalized state (`completed`, `cancelled`, or `discarded`) and sorting by finalized time, including in River UI. [PR #1374](https://github.com/riverqueue/river/pull/1374). + ## [0.47.0] - 2026-09-01 ### Added diff --git a/delete_many_params_test.go b/delete_many_params_test.go index 4b2d84db9..07d2ee343 100644 --- a/delete_many_params_test.go +++ b/delete_many_params_test.go @@ -1,10 +1,13 @@ package river import ( + "context" "testing" "github.com/stretchr/testify/require" + "github.com/riverqueue/river/internal/dblist" + "github.com/riverqueue/river/riverdriver/riverpgxv5" "github.com/riverqueue/river/rivertype" ) @@ -29,3 +32,20 @@ func TestJobDeleteManyParams_UnsafeAll(t *testing.T) { NewJobDeleteManyParams().IDs(123).UnsafeAll() }) } + +func TestJobDeleteManyParams_toDBParams(t *testing.T) { + t.Parallel() + + for _, state := range []rivertype.JobState{rivertype.JobStateAvailable, rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { + t.Run(string(state), func(t *testing.T) { + t.Parallel() + + params := NewJobDeleteManyParams().States(state).toDBParams() + require.Empty(t, params.Where) + driverParams, err := dblist.JobMakeDriverParams(context.Background(), params, riverpgxv5.New(nil)) + require.NoError(t, err) + require.Equal(t, "state = any(@state)", driverParams.WhereClause) + require.Equal(t, []string{string(state)}, driverParams.NamedArgs["state"]) + }) + } +} diff --git a/job_list_params.go b/job_list_params.go index aeb24d91f..ded6e878d 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -270,20 +270,51 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { orderBy = append(orderBy, dblist.JobListOrderBy{Expr: "id", Order: sortOrder}) + // Preserve custom SQL and its argument types without trying to parse it. + // In particular, an ungrouped OR may bypass the typed state filter, and + // custom SQL can reference the existing @state array argument. Metadata + // predicates also live in p.where; conservatively keep that path unchanged. + states := p.states + + // Copy conditions so reusing params does not accumulate generated cursor + // predicates or mix them into the caller's custom SQL. + where := append([]dblist.WherePredicate(nil), p.where...) + if len(p.where) == 0 && len(states) == 1 { + // Equality lets Postgres use the timestamp ordering of an index on + // (state, finalized_at). ANY does not establish that state is fixed. + where = append(where, dblist.WherePredicate{ + NamedArgs: map[string]any{"state": string(states[0])}, + SQL: "state = @state", + }) + + // Supported schemas enforce non-null finalized_at for these states. + // Make that explicit so Postgres can use the existing partial index. + if timeField == "finalized_at" { + switch states[0] { + case rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded: + where = append(where, dblist.WherePredicate{SQL: "finalized_at IS NOT NULL"}) + case rivertype.JobStateAvailable, rivertype.JobStatePending, rivertype.JobStateRetryable, rivertype.JobStateRunning, rivertype.JobStateScheduled: + } + } + + // The state filter is already represented in where. + states = nil + } + if p.after != nil { namedArgs := map[string]any{"after_id": p.after.id} if p.after.time.IsZero() { // order by ID only if sortOrder == dblist.SortOrderAsc { - p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id > @after_id)"}) + where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id > @after_id)"}) } else { - p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id < @after_id)"}) + where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: "(id < @after_id)"}) } } else { namedArgs["cursor_time"] = p.after.time if sortOrder == dblist.SortOrderAsc { - p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" > @cursor_time OR ("%s" = @cursor_time AND "id" > @after_id))`, timeField, timeField)}) + where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" > @cursor_time OR ("%s" = @cursor_time AND "id" > @after_id))`, timeField, timeField)}) } else { - p.where = append(p.where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" < @cursor_time OR ("%s" = @cursor_time AND "id" < @after_id))`, timeField, timeField)}) + where = append(where, dblist.WherePredicate{NamedArgs: namedArgs, SQL: fmt.Sprintf(`("%s" < @cursor_time OR ("%s" = @cursor_time AND "id" < @after_id))`, timeField, timeField)}) } } } @@ -296,10 +327,10 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { Priorities: p.priorities, Queues: p.queues, Schema: p.schema, - States: p.states, + States: states, TagsAll: p.tagsAll, TagsAny: p.tagsAny, - Where: p.where, + Where: where, }, nil } diff --git a/job_list_params_test.go b/job_list_params_test.go index 5e6de2f5e..5672866bf 100644 --- a/job_list_params_test.go +++ b/job_list_params_test.go @@ -1,6 +1,7 @@ package river import ( + "context" "encoding/json" "fmt" "testing" @@ -8,6 +9,9 @@ import ( "github.com/stretchr/testify/require" + "github.com/riverqueue/river/internal/dblist" + "github.com/riverqueue/river/riverdriver" + "github.com/riverqueue/river/riverdriver/riverpgxv5" "github.com/riverqueue/river/rivertype" ) @@ -210,7 +214,10 @@ func Test_JobListParams_toDBParams(t *testing.T) { OrderBy(JobListOrderByFinalizedAt, SortOrderDesc). toDBParams() require.NoError(t, err) - require.Equal(t, []rivertype.JobState{rivertype.JobStateCompleted}, dbParams.States) + driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil)) + require.NoError(t, err) + require.Equal(t, "completed", driverParams.NamedArgs["state"]) + require.Equal(t, "state = @state\n AND finalized_at IS NOT NULL", driverParams.WhereClause) }) t.Run("FinalizedAtWithMixedStates", func(t *testing.T) { @@ -283,3 +290,227 @@ func Test_JobListParams_toDBParams(t *testing.T) { require.Empty(t, dbParams.TagsAny) }) } + +func Test_JobListParams_toDBParams_CustomConditions(t *testing.T) { + t.Parallel() + + type testBundle struct { + driver *riverpgxv5.Driver + params *JobListParams + } + + setup := func(t *testing.T) *testBundle { + t.Helper() + + return &testBundle{ + driver: riverpgxv5.New(nil), + params: NewJobListParams().States(rivertype.JobStateCompleted). + OrderBy(JobListOrderByTime, SortOrderDesc), + } + } + + driverParamsFunc := func(t *testing.T, bundle *testBundle) *riverdriver.JobListParams { + t.Helper() + + dbParams, err := bundle.params.toDBParams() + require.NoError(t, err) + driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, bundle.driver) + require.NoError(t, err) + return driverParams + } + + t.Run("ContradictoryFinalizedAt", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // A condition that excludes completed jobs is intentional; don't repair it. + bundle.params = bundle.params.Where("finalized_at IS NULL") + driverParams := driverParamsFunc(t, bundle) + require.Equal(t, "state = any(@state)\n AND finalized_at IS NULL", driverParams.WhereClause) + require.Equal(t, map[string]any{"state": []string{"completed"}}, driverParams.NamedArgs) + }) + + t.Run("ContradictoryState", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // Keep both state conditions, even though no row can satisfy both. + bundle.params = bundle.params.Where("state = @other_state", NamedArgs{"other_state": "available"}) + driverParams := driverParamsFunc(t, bundle) + require.Equal(t, "state = any(@state)\n AND state = @other_state", driverParams.WhereClause) + require.Equal(t, map[string]any{"other_state": "available", "state": []string{"completed"}}, driverParams.NamedArgs) + }) + + t.Run("DuplicateStateArgument", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // The typed state filter owns @state. A custom argument with the same + // name must still report a conflict, rather than replacing the filter. + bundle.params = bundle.params.Where("state = @state", NamedArgs{"state": "available"}) + dbParams, err := bundle.params.toDBParams() + require.NoError(t, err) + _, err = dblist.JobMakeDriverParams(context.Background(), dbParams, bundle.driver) + require.EqualError(t, err, "named argument @state already registered") + }) + + t.Run("ExistingStateArgument", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // Custom SQL can reuse @state, so its value must remain an array. + bundle.params = bundle.params.Where("state = ANY(@state)") + driverParams := driverParamsFunc(t, bundle) + require.Equal(t, "state = any(@state)\n AND state = ANY(@state)", driverParams.WhereClause) + require.Equal(t, map[string]any{"state": []string{"completed"}}, driverParams.NamedArgs) + }) + + t.Run("GroupedOr", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + bundle.params = bundle.params.Where("(state = @other_state OR finalized_at IS NULL)", NamedArgs{"other_state": "completed"}) + driverParams := driverParamsFunc(t, bundle) + require.Equal(t, "state = any(@state)\n AND (state = @other_state OR finalized_at IS NULL)", driverParams.WhereClause) + require.Equal(t, map[string]any{"other_state": "completed", "state": []string{"completed"}}, driverParams.NamedArgs) + }) + + t.Run("Metadata", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // Metadata shares Where's internal condition list, so it also retains + // the array filter and receives no inferred null condition. + bundle.params = bundle.params.Metadata(`{"selected":true}`) + driverParams := driverParamsFunc(t, bundle) + require.Equal(t, "state = any(@state)\n AND metadata @> @metadata_fragment::jsonb", driverParams.WhereClause) + require.Equal(t, map[string]any{"metadata_fragment": `{"selected":true}`, "state": []string{"completed"}}, driverParams.NamedArgs) + }) + + t.Run("Pagination", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // Appending the cursor must preserve the custom OR's existing grouping + // and leave its named argument independent of the cursor arguments. + cursorTime := time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC) + bundle.params = bundle.params. + Where("state = @other_state OR finalized_at IS NULL", NamedArgs{"other_state": "completed"}). + After(&JobListCursor{id: 42, time: cursorTime}) + driverParams := driverParamsFunc(t, bundle) + const cursorSQL = `("finalized_at" < @cursor_time OR ("finalized_at" = @cursor_time AND "id" < @after_id))` + require.Equal(t, "state = any(@state)\n AND state = @other_state OR finalized_at IS NULL\n AND "+cursorSQL, driverParams.WhereClause) + require.Equal(t, map[string]any{ + "after_id": int64(42), + "cursor_time": cursorTime, + "other_state": "completed", + "state": []string{"completed"}, + }, driverParams.NamedArgs) + }) + + t.Run("RepeatedConversion", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + bundle.params = bundle.params.Where("state = ANY(@state)"). + After(&JobListCursor{id: 42, time: time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC)}) + + // Conversion must not store generated conditions in the caller's params + // or append another cursor each time the same params are used. + first := driverParamsFunc(t, bundle) + second := driverParamsFunc(t, bundle) + require.Equal(t, first, second) + require.Equal(t, []dblist.WherePredicate{{SQL: "state = ANY(@state)"}}, bundle.params.where) + }) + + t.Run("UngroupedOr", func(t *testing.T) { + t.Parallel() + + bundle := setup(t) + + // This OR can admit non-finalized rows despite the typed state filter. + // Adding parentheses or a non-null condition would change the results. + bundle.params = bundle.params.Where("false OR finalized_at IS NULL") + driverParams := driverParamsFunc(t, bundle) + require.Equal(t, "state = any(@state)\n AND false OR finalized_at IS NULL", driverParams.WhereClause) + require.Equal(t, map[string]any{"state": []string{"completed"}}, driverParams.NamedArgs) + }) +} + +func Test_JobListParams_toDBParamsFinalizedIndex(t *testing.T) { + t.Parallel() + + for _, state := range []rivertype.JobState{rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { + for _, field := range []JobListOrderByField{JobListOrderByFinalizedAt, JobListOrderByTime} { + for _, order := range []SortOrder{SortOrderAsc, SortOrderDesc} { + t.Run(fmt.Sprintf("%s/%s/%d", state, field, order), func(t *testing.T) { + t.Parallel() + + params := NewJobListParams().States(state).OrderBy(field, order) + dbParams, err := params.toDBParams() + require.NoError(t, err) + driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil)) + require.NoError(t, err) + require.Equal(t, "state = @state\n AND finalized_at IS NOT NULL", driverParams.WhereClause) + require.Equal(t, map[string]any{"state": string(state)}, driverParams.NamedArgs) + direction := "ASC" + if order == SortOrderDesc { + direction = "DESC" + } + require.Equal(t, "finalized_at "+direction+", id "+direction, driverParams.OrderByClause) + params = params.After(&JobListCursor{id: 42, time: time.Now().UTC()}) + dbParams, err = params.toDBParams() + require.NoError(t, err) + driverParams, err = dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil)) + require.NoError(t, err) + require.Len(t, dbParams.Where, 3) + require.Equal(t, "state = @state\n AND finalized_at IS NOT NULL\n AND "+dbParams.Where[2].SQL, driverParams.WhereClause) + require.Equal(t, map[string]any{"after_id": int64(42), "cursor_time": params.after.time, "state": string(state)}, driverParams.NamedArgs) + require.Empty(t, params.where) + repeated, err := params.toDBParams() + require.NoError(t, err) + require.Equal(t, dbParams, repeated) + }) + } + } + } +} + +func Test_JobListParams_toDBParamsWithoutFinalizedIndex(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + name string + params *JobListParams + where string + }{ + {"Default", NewJobListParams(), "state = any(@state)"}, + {"FinalizedByID", NewJobListParams().States(rivertype.JobStateCompleted), "state = @state"}, + {"FinalizedByScheduledAt", NewJobListParams().States(rivertype.JobStateCompleted).OrderBy(JobListOrderByScheduledAt, SortOrderAsc), "state = @state"}, + {"MixedStatesByTime", NewJobListParams().States(rivertype.JobStateCompleted, rivertype.JobStateAvailable).OrderBy(JobListOrderByTime, SortOrderDesc), "state = any(@state)"}, + {"MultipleFinalizedStates", NewJobListParams().OrderBy(JobListOrderByFinalizedAt, SortOrderAsc), "state = any(@state)"}, + {"MultipleNonFinalizedStates", NewJobListParams().States(rivertype.JobStateAvailable, rivertype.JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc), "state = any(@state)"}, + {"NonFinalized", NewJobListParams().States(rivertype.JobStateAvailable).OrderBy(JobListOrderByTime, SortOrderAsc), "state = @state"}, + {"Unfiltered", NewJobListParams().States(), "true"}, + {"Unknown", NewJobListParams().States("unknown").OrderBy(JobListOrderByFinalizedAt, SortOrderAsc), "state = @state"}, + {"UnknownAndFinalized", NewJobListParams().States(rivertype.JobStateCompleted, "unknown").OrderBy(JobListOrderByTime, SortOrderDesc), "state = any(@state)"}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + dbParams, err := tt.params.toDBParams() + require.NoError(t, err) + driverParams, err := dblist.JobMakeDriverParams(context.Background(), dbParams, riverpgxv5.New(nil)) + require.NoError(t, err) + require.Equal(t, tt.where, driverParams.WhereClause) + }) + } +} diff --git a/riverdriver/riverdrivertest/driver_client_test.go b/riverdriver/riverdrivertest/driver_client_test.go index 83bb4c9a9..8d1f895f6 100644 --- a/riverdriver/riverdrivertest/driver_client_test.go +++ b/riverdriver/riverdrivertest/driver_client_test.go @@ -4,6 +4,7 @@ import ( "context" "database/sql" "math" + "slices" "testing" "time" @@ -24,6 +25,7 @@ import ( "github.com/riverqueue/river/rivershared/riversharedtest" "github.com/riverqueue/river/rivershared/testfactory" "github.com/riverqueue/river/rivershared/testsignal" + "github.com/riverqueue/river/rivershared/util/sliceutil" "github.com/riverqueue/river/rivershared/util/testutil" "github.com/riverqueue/river/rivershared/util/urlutil" "github.com/riverqueue/river/rivertype" @@ -652,6 +654,265 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.Equal(t, job.ID, listRes.Jobs[0].ID) }) + t.Run("JobListCustomStateConditions", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + args river.NamedArgs + name string + sql string + wantAvailable bool + wantCompleted bool + }{ + {nil, "ContradictoryFinalizedAt", "finalized_at IS NULL", false, false}, + {river.NamedArgs{"other_state": "available"}, "ContradictoryState", "state = @other_state", false, false}, + {river.NamedArgs{"other_state": "completed"}, "GroupedOr", "(state = @other_state OR finalized_at IS NULL)", false, true}, + {nil, "UngroupedOr", "false OR finalized_at IS NULL", true, false}, + {river.NamedArgs{"other_state": "completed"}, "UngroupedOrWithNamedArgument", "state = @other_state OR finalized_at IS NULL", true, true}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + available := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Schema: bundle.schema}) + completed := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{ + FinalizedAt: new(time.Now().UTC().Truncate(time.Second)), Schema: bundle.schema, State: new(rivertype.JobStateCompleted), + }) + params := river.NewJobListParams().States(rivertype.JobStateCompleted). + OrderBy(river.JobListOrderByTime, river.SortOrderDesc).Where(tt.sql, tt.args) + result, err := client.JobList(ctx, params) + require.NoError(t, err) + gotIDs := make([]int64, 0, len(result.Jobs)) + var wantIDs []int64 + for _, job := range result.Jobs { + gotIDs = append(gotIDs, job.ID) + } + if tt.wantAvailable { + wantIDs = append(wantIDs, available.ID) + } + if tt.wantCompleted { + wantIDs = append(wantIDs, completed.ID) + } + require.ElementsMatch(t, wantIDs, gotIDs) + }) + } + }) + + t.Run("JobListCustomStatePagination", func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + if bundle.driver.DatabaseName() != riverdriver.DatabaseNamePostgres { + t.Skip("uses PostgreSQL array syntax and time cursors") + } + now := time.Now().UTC().Truncate(time.Second) + wantIDs := make([]int64, 0, 3) + for range 3 { + job := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{ + FinalizedAt: &now, Metadata: []byte(`{"selected":true}`), Schema: bundle.schema, State: new(rivertype.JobStateCompleted), + }) + wantIDs = append(wantIDs, job.ID) + } + slices.Reverse(wantIDs) + params := river.NewJobListParams().States(rivertype.JobStateCompleted). + OrderBy(river.JobListOrderByTime, river.SortOrderDesc).First(1). + Where("(state = ANY(@state) OR finalized_at IS NULL)"). + Where("id > @minimum_id", river.NamedArgs{"minimum_id": 0}).Metadata(`{"selected":true}`) + var gotIDs []int64 + pageParams := params + for page := range 4 { + result, err := client.JobList(ctx, pageParams) + require.NoError(t, err) + if page == 3 { + require.Empty(t, result.Jobs) + break + } + require.Len(t, result.Jobs, 1) + gotIDs = append(gotIDs, result.Jobs[0].ID) + pageParams = params.After(result.LastCursor) + } + require.Equal(t, wantIDs, gotIDs) + }) + + t.Run("JobListFinalized", func(t *testing.T) { + t.Parallel() + + type testBundle struct { + driver riverdriver.Driver[TTx] + exec riverdriver.Executor + jobs map[rivertype.JobState][]*rivertype.JobRow + now time.Time + schema string + } + + setup := func(t *testing.T) (*river.Client[TTx], *testBundle) { + t.Helper() + + client, bundle := setup(t) + now := time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC) + jobs := make(map[rivertype.JobState][]*rivertype.JobRow) + + // IDs and timestamps deliberately disagree. Each timestamp has three + // jobs, so a two-job page ends partway through a group of equal times. + for _, state := range []rivertype.JobState{rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { + for _, offset := range []time.Duration{time.Second, 0, time.Second, 0, time.Second, 0} { + job := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{ + FinalizedAt: new(now.Add(offset)), + Kind: new("selected"), + Priority: new(2), + Queue: new("selected"), + Schema: bundle.schema, + State: new(state), + Tags: []string{"alpha", "beta"}, + }) + jobs[state] = append(jobs[state], job) + } + } + + return client, &testBundle{ + driver: bundle.driver, + exec: bundle.exec, + jobs: jobs, + now: now, + schema: bundle.schema, + } + } + + t.Run("Filters", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + includeID bool + name string + optsFunc func(*testfactory.JobOpts) + }{ + {false, "IDs", func(opts *testfactory.JobOpts) {}}, + {true, "Kinds", func(opts *testfactory.JobOpts) { opts.Kind = new("other") }}, + {true, "Priorities", func(opts *testfactory.JobOpts) { opts.Priority = new(3) }}, + {true, "Queues", func(opts *testfactory.JobOpts) { opts.Queue = new("other") }}, + {true, "States", func(opts *testfactory.JobOpts) { opts.State = new(rivertype.JobStateDiscarded) }}, + {true, "TagsAll", func(opts *testfactory.JobOpts) { opts.Tags = []string{"beta"} }}, + {true, "TagsAny", func(opts *testfactory.JobOpts) { opts.Tags = []string{"alpha"} }}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + // All filters match the completed jobs. The extra job fails only + // the filter named by this case, so no other filter can hide it. + opts := &testfactory.JobOpts{ + FinalizedAt: &bundle.now, + Kind: new("selected"), + Priority: new(2), + Queue: new("selected"), + Schema: bundle.schema, + State: new(rivertype.JobStateCompleted), + Tags: []string{"alpha", "beta"}, + } + tt.optsFunc(opts) + excludedJob := testfactory.Job(ctx, t, bundle.exec, opts) + wantIDs := sliceutil.Map(bundle.jobs[rivertype.JobStateCompleted], func(job *rivertype.JobRow) int64 { return job.ID }) + filterIDs := slices.Clone(wantIDs) + if tt.includeID { + filterIDs = append(filterIDs, excludedJob.ID) + } + + listRes, err := client.JobList(ctx, river.NewJobListParams(). + IDs(filterIDs...).Kinds("selected").Priorities(2).Queues("selected"). + States(rivertype.JobStateCompleted).TagsAll("alpha").TagsAny("beta", "gamma"). + OrderBy(river.JobListOrderByTime, river.SortOrderDesc)) + require.NoError(t, err) + require.ElementsMatch(t, wantIDs, sliceutil.Map(listRes.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + }) + } + }) + + t.Run("Ordering", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + name string + params *river.JobListParams + want []int + }{ + {"FinalizedAtAsc", river.NewJobListParams().OrderBy(river.JobListOrderByFinalizedAt, river.SortOrderAsc), []int{1, 3}}, + {"FinalizedAtDesc", river.NewJobListParams().OrderBy(river.JobListOrderByFinalizedAt, river.SortOrderDesc), []int{4, 2}}, + {"TimeAsc", river.NewJobListParams().OrderBy(river.JobListOrderByTime, river.SortOrderAsc), []int{1, 3}}, + {"TimeDesc", river.NewJobListParams().OrderBy(river.JobListOrderByTime, river.SortOrderDesc), []int{4, 2}}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + for _, state := range []rivertype.JobState{rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { + jobs := bundle.jobs[state] + listRes, err := client.JobList(ctx, tt.params.States(state).First(2)) + require.NoError(t, err) + require.Equal(t, []int64{jobs[tt.want[0]].ID, jobs[tt.want[1]].ID}, + sliceutil.Map(listRes.Jobs, func(job *rivertype.JobRow) int64 { return job.ID }), "state: %s", state) + } + }) + } + }) + + t.Run("Pagination", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + name string + params *river.JobListParams + want []int + }{ + {"FinalizedAtAsc", river.NewJobListParams().OrderBy(river.JobListOrderByFinalizedAt, river.SortOrderAsc), []int{1, 3, 5, 0, 2, 4}}, + {"FinalizedAtDesc", river.NewJobListParams().OrderBy(river.JobListOrderByFinalizedAt, river.SortOrderDesc), []int{4, 2, 0, 5, 3, 1}}, + {"TimeAsc", river.NewJobListParams().OrderBy(river.JobListOrderByTime, river.SortOrderAsc), []int{1, 3, 5, 0, 2, 4}}, + {"TimeDesc", river.NewJobListParams().OrderBy(river.JobListOrderByTime, river.SortOrderDesc), []int{4, 2, 0, 5, 3, 1}}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + // SQLite stores timestamps as text, but list cursors bind time.Time + // directly. Re-enable once cursor binding uses the stored format. + if bundle.driver.DatabaseName() == riverdriver.DatabaseNameSQLite { + t.Skip("SQLite time cursor arguments do not match the stored timestamp format") + } + + for _, state := range []rivertype.JobState{rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { + params := tt.params.States(state).First(2) + jobs := bundle.jobs[state] + wantIDs := sliceutil.Map(tt.want, func(index int) int64 { return jobs[index].ID }) + + firstPage, err := client.JobList(ctx, params) + require.NoError(t, err) + require.Equal(t, wantIDs[:2], sliceutil.Map(firstPage.Jobs, func(job *rivertype.JobRow) int64 { return job.ID }), "state: %s", state) + + // Resume in the middle of a timestamp group with a serialized cursor. + encoded, err := firstPage.LastCursor.MarshalText() + require.NoError(t, err) + var cursor river.JobListCursor + require.NoError(t, cursor.UnmarshalText(encoded)) + secondPage, err := client.JobList(ctx, params.After(&cursor)) + require.NoError(t, err) + require.Equal(t, wantIDs[2:4], sliceutil.Map(secondPage.Jobs, func(job *rivertype.JobRow) int64 { return job.ID }), "state: %s", state) + + // The next boundary splits the other timestamp group. Use a job-derived cursor. + thirdPage, err := client.JobList(ctx, params.After(river.JobListCursorFromJob(secondPage.Jobs[1]))) + require.NoError(t, err) + require.Equal(t, wantIDs[4:], sliceutil.Map(thirdPage.Jobs, func(job *rivertype.JobRow) int64 { return job.ID }), "state: %s", state) + + emptyPage, err := client.JobList(ctx, params.After(thirdPage.LastCursor)) + require.NoError(t, err) + require.Empty(t, emptyPage.Jobs) + } + }) + } + }) + }) + t.Run("JobListMetadata", func(t *testing.T) { t.Parallel() @@ -673,6 +934,51 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.Equal(t, job.ID, listRes.Jobs[0].ID) }) + t.Run("JobListStateFilters", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + name string + params *river.JobListParams + wantJobIndexes []int + }{ + {"Default", river.NewJobListParams(), []int{0, 1, 2, 3}}, + {"ExplicitEmpty", river.NewJobListParams().States(), []int{0, 1, 2, 3}}, + {"FinalizedDefaults", river.NewJobListParams().OrderBy(river.JobListOrderByFinalizedAt, river.SortOrderDesc), []int{1, 2, 3}}, + {"Mixed", river.NewJobListParams().States(rivertype.JobStateCompleted, rivertype.JobStateAvailable).OrderBy(river.JobListOrderByTime, river.SortOrderDesc), []int{0, 2}}, + {"NonFinalized", river.NewJobListParams().States(rivertype.JobStateAvailable).OrderBy(river.JobListOrderByTime, river.SortOrderDesc), []int{0}}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + now := time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC) + allIDs := make([]int64, 0, 4) + for _, state := range []rivertype.JobState{rivertype.JobStateAvailable, rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { + opts := &testfactory.JobOpts{Schema: bundle.schema, State: new(state)} + if state != rivertype.JobStateAvailable { + opts.FinalizedAt = &now + } + job := testfactory.Job(ctx, t, bundle.exec, opts) + allIDs = append(allIDs, job.ID) + } + wantIDs := make([]int64, 0, len(tt.wantJobIndexes)) + for _, index := range tt.wantJobIndexes { + wantIDs = append(wantIDs, allIDs[index]) + } + + result, err := client.JobList(ctx, tt.params) + require.NoError(t, err) + gotIDs := make([]int64, 0, len(result.Jobs)) + for _, job := range result.Jobs { + gotIDs = append(gotIDs, job.ID) + } + require.ElementsMatch(t, wantIDs, gotIDs) + }) + } + }) + t.Run("JobListTags", func(t *testing.T) { t.Parallel() From 0f01eafd1abf0f63e0ce7b9cf2fc6a2a62fc53c3 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Wed, 9 Sep 2026 20:25:39 -0500 Subject: [PATCH 2/2] fix SQLite job list cursor timestamps SQLite stores job timestamps as formatted text, while JobList passes cursor values directly to database/sql. Different encodings can skip or repeat jobs when a page boundary shares a timestamp. Format `time.Time` and `*time.Time` list arguments with SQLite's existing timestamp helpers. Copy the argument map so conversion leaves reusable parameters intact, including custom conditions and nullable values. Enable finalized-job pagination coverage for SQLite, libSQL, and Turso. Cover scheduled-job pagination, millisecond precision, time zones, nullable arguments, and repeated use of driver parameters. --- CHANGELOG.md | 1 + .../riverdrivertest/driver_client_test.go | 101 ++++++++++++++++-- .../riversqlite/river_sqlite_driver.go | 16 ++- 3 files changed, 107 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 43b81d6bc..c81eecc83 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### 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). - Improved PostgreSQL job listing performance when filtering by one finalized state (`completed`, `cancelled`, or `discarded`) and sorting by finalized time, including in River UI. [PR #1374](https://github.com/riverqueue/river/pull/1374). ## [0.47.0] - 2026-09-01 diff --git a/riverdriver/riverdrivertest/driver_client_test.go b/riverdriver/riverdrivertest/driver_client_test.go index 8d1f895f6..1e361a33c 100644 --- a/riverdriver/riverdrivertest/driver_client_test.go +++ b/riverdriver/riverdrivertest/driver_client_test.go @@ -703,7 +703,7 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, client, bundle := setup(t) if bundle.driver.DatabaseName() != riverdriver.DatabaseNamePostgres { - t.Skip("uses PostgreSQL array syntax and time cursors") + t.Skip("uses PostgreSQL array and JSON containment syntax") } now := time.Now().UTC().Truncate(time.Second) wantIDs := make([]int64, 0, 3) @@ -738,7 +738,6 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, t.Parallel() type testBundle struct { - driver riverdriver.Driver[TTx] exec riverdriver.Executor jobs map[rivertype.JobState][]*rivertype.JobRow now time.Time @@ -749,7 +748,7 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, t.Helper() client, bundle := setup(t) - now := time.Date(2026, 9, 9, 12, 0, 0, 0, time.UTC) + now := time.Date(2026, 9, 9, 12, 0, 0, 123000000, time.UTC) jobs := make(map[rivertype.JobState][]*rivertype.JobRow) // IDs and timestamps deliberately disagree. Each timestamp has three @@ -770,7 +769,6 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, } return client, &testBundle{ - driver: bundle.driver, exec: bundle.exec, jobs: jobs, now: now, @@ -875,12 +873,6 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, client, bundle := setup(t) - // SQLite stores timestamps as text, but list cursors bind time.Time - // directly. Re-enable once cursor binding uses the stored format. - if bundle.driver.DatabaseName() == riverdriver.DatabaseNameSQLite { - t.Skip("SQLite time cursor arguments do not match the stored timestamp format") - } - for _, state := range []rivertype.JobState{rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded} { params := tt.params.States(state).First(2) jobs := bundle.jobs[state] @@ -934,6 +926,48 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.Equal(t, job.ID, listRes.Jobs[0].ID) }) + t.Run("JobListScheduledPagination", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + name string + order river.SortOrder + }{ + {"Ascending", river.SortOrderAsc}, + {"Descending", river.SortOrderDesc}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t) + + // The time comparison must recognize equality so the ID cursor + // can advance through available jobs sharing a scheduled time. + now := time.Date(2026, 9, 9, 12, 0, 0, 123000000, time.UTC) + opts := &testfactory.JobOpts{Kind: new("selected"), ScheduledAt: &now, Schema: bundle.schema} + job1 := testfactory.Job(ctx, t, bundle.exec, opts) + job2 := testfactory.Job(ctx, t, bundle.exec, opts) + wantIDs := []int64{job1.ID, job2.ID} + if tt.order == river.SortOrderDesc { + slices.Reverse(wantIDs) + } + params := river.NewJobListParams().States(rivertype.JobStateAvailable). + OrderBy(river.JobListOrderByScheduledAt, tt.order).First(1). + Where("kind = @kind_name", river.NamedArgs{"kind_name": "selected"}) + + firstPage, err := client.JobList(ctx, params) + require.NoError(t, err) + require.Equal(t, wantIDs[:1], sliceutil.Map(firstPage.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + secondPage, err := client.JobList(ctx, params.After(firstPage.LastCursor)) + require.NoError(t, err) + require.Equal(t, wantIDs[1:], sliceutil.Map(secondPage.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + emptyPage, err := client.JobList(ctx, params.After(secondPage.LastCursor)) + require.NoError(t, err) + require.Empty(t, emptyPage.Jobs) + }) + } + }) + t.Run("JobListStateFilters", func(t *testing.T) { t.Parallel() @@ -1029,6 +1063,53 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.Len(t, listRes.Jobs, 6) }) + t.Run("JobListTimeArguments", func(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + name string + valueFunc func(time.Time) any + wantFinalized bool + }{ + {"NilTimePointer", func(time.Time) any { return (*time.Time)(nil) }, false}, + {"Time", func(value time.Time) any { return value }, true}, + {"TimePointer", func(value time.Time) any { return &value }, true}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + _, bundle := setup(t) + + now := time.Date(2026, 9, 9, 12, 0, 0, 123000000, time.UTC) + available := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Kind: new("selected"), Schema: bundle.schema}) + completed := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{ + FinalizedAt: &now, Kind: new("selected"), Schema: bundle.schema, State: new(rivertype.JobStateCompleted), + }) + wantID := available.ID + if tt.wantFinalized { + wantID = completed.ID + } + + // Match the same instant in another zone, including milliseconds. + // Use the driver directly to verify it leaves reusable arguments intact. + arg := tt.valueFunc(now.In(time.FixedZone("test", -7*60*60))) + params := &riverdriver.JobListParams{ + Max: 100, + NamedArgs: map[string]any{"kind": "selected", "time": arg}, + OrderByClause: "id ASC", + Schema: bundle.schema, + WhereClause: "kind = @kind AND (finalized_at = @time OR (finalized_at IS NULL AND @time IS NULL))", + } + for range 2 { + jobs, err := bundle.exec.JobList(ctx, params) + require.NoError(t, err) + require.Equal(t, []int64{wantID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + } + require.Equal(t, map[string]any{"kind": "selected", "time": arg}, params.NamedArgs) + }) + } + }) + t.Run("JobListTx", func(t *testing.T) { t.Parallel() diff --git a/riverdriver/riversqlite/river_sqlite_driver.go b/riverdriver/riversqlite/river_sqlite_driver.go index c778396d1..27296a442 100644 --- a/riverdriver/riversqlite/river_sqlite_driver.go +++ b/riverdriver/riversqlite/river_sqlite_driver.go @@ -26,6 +26,7 @@ import ( "errors" "fmt" "io/fs" + "maps" "math" "os" "path" @@ -702,10 +703,23 @@ func (e *Executor) JobKindList(ctx context.Context, params *riverdriver.JobKindL } func (e *Executor) JobList(ctx context.Context, params *riverdriver.JobListParams) ([]*rivertype.JobRow, error) { + // List cursors and custom conditions can contain time values. Bind them + // in the same format as stored timestamps so comparisons work in SQLite. + // Keep the caller's arguments intact for repeated or concurrent use. + namedArgs := maps.Clone(params.NamedArgs) + for name, arg := range namedArgs { + switch arg := arg.(type) { + case time.Time: + namedArgs[name] = timeString(arg) + case *time.Time: + namedArgs[name] = timeStringNullable(arg) + } + } + ctx = sqlctemplate.WithReplacements(ctx, map[string]sqlctemplate.Replacement{ "order_by_clause": {Value: params.OrderByClause}, "where_clause": {Value: params.WhereClause}, - }, params.NamedArgs) + }, namedArgs) jobs, err := dbsqlc.New().JobList(schemaTemplateParam(ctx, params.Schema), e.dbtx, int64(params.Max)) if err != nil {