diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b875f8a..c81eecc8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### 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 ### Added diff --git a/delete_many_params_test.go b/delete_many_params_test.go index 4b2d84db..07d2ee34 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 aeb24d91..ded6e878 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 5e6de2f5..5672866b 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 83bb4c9a..1e361a33 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,257 @@ 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 and JSON containment syntax") + } + 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 { + 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, 123000000, 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{ + 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) + + 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 +926,93 @@ 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() + + 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() @@ -723,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 c778396d..27296a44 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 {