diff --git a/CHANGELOG.md b/CHANGELOG.md index b8793d8a..62501c68 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] +### Changed + +- **Breaking change:** JobList/JobListTx now support querying Jobs by a list of Job Kinds and States. Also allows for filtering by specific timestamp values. Thank you Jos Kraaijeveld (@thatjos)! 🙏🏻 [PR #236](https://github.com/riverqueue/river/pull/236). + ## [0.3.0] - 2024-04-15 ### Added diff --git a/client.go b/client.go index d5aea11a..63dd9479 100644 --- a/client.go +++ b/client.go @@ -1446,6 +1446,13 @@ func validateQueueName(queueName string) error { return nil } +// JobListResult is the result of a job list operation. It contains a list of +// jobs and a cursor for fetching the next page of results. +type JobListResult struct { + Jobs []*rivertype.JobRow + LastCursor *JobListCursor +} + // JobList returns a paginated list of jobs matching the provided filters. The // provided context is used for the underlying Postgres query and can be used to // cancel the operation or apply a timeout. @@ -1455,7 +1462,7 @@ func validateQueueName(queueName string) error { // if err != nil { // // handle error // } -func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) ([]*rivertype.JobRow, error) { +func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) (*JobListResult, error) { if !c.driver.HasPool() { return nil, errNoDriverDBPool } @@ -1468,7 +1475,15 @@ func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) ([]*ri return nil, err } - return dblist.JobList(ctx, c.driver.GetExecutor(), dbParams) + jobs, err := dblist.JobList(ctx, c.driver.GetExecutor(), dbParams) + if err != nil { + return nil, err + } + res := &JobListResult{Jobs: jobs} + if len(jobs) > 0 { + res.LastCursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) + } + return res, nil } // JobListTx returns a paginated list of jobs matching the provided filters. The @@ -1480,16 +1495,25 @@ func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) ([]*ri // if err != nil { // // handle error // } -func (c *Client[TTx]) JobListTx(ctx context.Context, tx TTx, params *JobListParams) ([]*rivertype.JobRow, error) { +func (c *Client[TTx]) JobListTx(ctx context.Context, tx TTx, params *JobListParams) (*JobListResult, error) { if params == nil { params = NewJobListParams() } + dbParams, err := params.toDBParams() if err != nil { return nil, err } - return dblist.JobList(ctx, c.driver.UnwrapExecutor(tx), dbParams) + jobs, err := dblist.JobList(ctx, c.driver.UnwrapExecutor(tx), dbParams) + if err != nil { + return nil, err + } + res := &JobListResult{Jobs: jobs} + if len(jobs) > 0 { + res.LastCursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) + } + return res, nil } // PeriodicJobs returns the currently configured set of periodic jobs for the diff --git a/client_test.go b/client_test.go index d5feac21..59d7bb5a 100644 --- a/client_test.go +++ b/client_test.go @@ -702,9 +702,9 @@ func Test_Client_Stop(t *testing.T) { require.NoError(t, client.Stop(ctx)) - runningJobs, err := client.JobList(ctx, NewJobListParams().State(rivertype.JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.NoError(t, err) - require.Empty(t, runningJobs, "expected no jobs to be left running") + require.Empty(t, res.Jobs, "expected no jobs to be left running") }) t.Run("WithSubscriber", func(t *testing.T) { @@ -1449,14 +1449,14 @@ func Test_Client_JobList(t *testing.T) { job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Kind: ptrutil.Ptr("test_kind_1")}) job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Kind: ptrutil.Ptr("test_kind_2")}) - jobs, err := client.JobList(ctx, NewJobListParams().Kinds("test_kind_1")) + res, err := client.JobList(ctx, NewJobListParams().Kinds("test_kind_1")) require.NoError(t, err) // jobs ordered by ScheduledAt ASC by default - require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().Kinds("test_kind_2")) + res, err = client.JobList(ctx, NewJobListParams().Kinds("test_kind_2")) require.NoError(t, err) - require.Equal(t, []int64{job3.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job3.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) t.Run("FiltersByQueue", func(t *testing.T) { //nolint:dupl @@ -1468,14 +1468,14 @@ func Test_Client_JobList(t *testing.T) { job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Queue: ptrutil.Ptr("queue_1")}) job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Queue: ptrutil.Ptr("queue_2")}) - jobs, err := client.JobList(ctx, NewJobListParams().Queues("queue_1")) + res, err := client.JobList(ctx, NewJobListParams().Queues("queue_1")) require.NoError(t, err) // jobs ordered by ScheduledAt ASC by default - require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().Queues("queue_2")) + res, err = client.JobList(ctx, NewJobListParams().Queues("queue_2")) require.NoError(t, err) - require.Equal(t, []int64{job3.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job3.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) t.Run("FiltersByState", func(t *testing.T) { @@ -1487,14 +1487,14 @@ func Test_Client_JobList(t *testing.T) { job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)}) job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning)}) - jobs, err := client.JobList(ctx, NewJobListParams().State(JobStateAvailable)) + res, err := client.JobList(ctx, NewJobListParams().States(JobStateAvailable)) require.NoError(t, err) // jobs ordered by ScheduledAt ASC by default - require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State(JobStateRunning)) + res, err = client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.NoError(t, err) - require.Equal(t, []int64{job3.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job3.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) t.Run("SortsAvailableRetryableAndScheduledJobsByScheduledAt", func(t *testing.T) { @@ -1513,13 +1513,13 @@ func Test_Client_JobList(t *testing.T) { job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(dbState), ScheduledAt: &now}) job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(dbState), ScheduledAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) - jobs, err := client.JobList(ctx, NewJobListParams().State(state)) + res, err := client.JobList(ctx, NewJobListParams().States(state)) require.NoError(t, err) - require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State(state).OrderBy(JobListOrderByTime, SortOrderDesc)) + res, err = client.JobList(ctx, NewJobListParams().States(state).OrderBy(JobListOrderByTime, SortOrderDesc)) require.NoError(t, err) - require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) } }) @@ -1539,13 +1539,13 @@ func Test_Client_JobList(t *testing.T) { job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(dbState), FinalizedAt: ptrutil.Ptr(now.Add(-10 * time.Second))}) job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(dbState), FinalizedAt: ptrutil.Ptr(now.Add(-15 * time.Second))}) - jobs, err := client.JobList(ctx, NewJobListParams().State(state)) + res, err := client.JobList(ctx, NewJobListParams().States(state)) require.NoError(t, err) - require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State(state).OrderBy(JobListOrderByTime, SortOrderDesc)) + res, err = client.JobList(ctx, NewJobListParams().States(state).OrderBy(JobListOrderByTime, SortOrderDesc)) require.NoError(t, err) - require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) } }) @@ -1558,17 +1558,17 @@ func Test_Client_JobList(t *testing.T) { job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), AttemptedAt: &now}) job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), AttemptedAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) - jobs, err := client.JobList(ctx, NewJobListParams().State(JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.NoError(t, err) - require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State(JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc)) + res, err = client.JobList(ctx, NewJobListParams().States(JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc)) require.NoError(t, err) // Sort order was explicitly reversed: - require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) - t.Run("WithNilParamsFiltersToAvailableByDefault", func(t *testing.T) { + t.Run("WithNilParamsFiltersToAllStatesByDefault", func(t *testing.T) { t.Parallel() client, bundle := setup(t) @@ -1576,12 +1576,12 @@ func Test_Client_JobList(t *testing.T) { now := time.Now().UTC() job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable), ScheduledAt: &now}) job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable), ScheduledAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) - _ = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning)}) + job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), ScheduledAt: ptrutil.Ptr(now.Add(-2 * time.Second))}) - jobs, err := client.JobList(ctx, nil) + res, err := client.JobList(ctx, nil) require.NoError(t, err) // sort order is switched by ScheduledAt values: - require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job2.ID, job3.ID, job1.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) t.Run("PaginatesWithAfter", func(t *testing.T) { @@ -1592,22 +1592,31 @@ func Test_Client_JobList(t *testing.T) { now := time.Now().UTC() job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable), ScheduledAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable), ScheduledAt: &now}) - job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), AttemptedAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) - job4 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), AttemptedAt: &now}) - job5 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCompleted), FinalizedAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) - job6 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCompleted), FinalizedAt: &now}) + job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), ScheduledAt: ptrutil.Ptr(now.Add(-5 * time.Second)), AttemptedAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) + job4 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning), ScheduledAt: ptrutil.Ptr(now.Add(-6 * time.Second)), AttemptedAt: &now}) + job5 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCompleted), ScheduledAt: ptrutil.Ptr(now.Add(-7 * time.Second)), FinalizedAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) + job6 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateCompleted), ScheduledAt: ptrutil.Ptr(now.Add(-7 * time.Second)), FinalizedAt: &now}) - jobs, err := client.JobList(ctx, NewJobListParams().After(JobListCursorFromJob(job1))) + res, err := client.JobList(ctx, NewJobListParams().States(rivertype.JobStateAvailable).After(JobListCursorFromJob(job1, JobListOrderByTime))) require.NoError(t, err) - require.Equal(t, []int64{job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, job2.ID, res.LastCursor.id) - jobs, err = client.JobList(ctx, NewJobListParams().State(rivertype.JobStateRunning).After(JobListCursorFromJob(job3))) + res, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning).After(JobListCursorFromJob(job3, JobListOrderByTime))) require.NoError(t, err) - require.Equal(t, []int64{job4.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job4.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, job4.ID, res.LastCursor.id) - jobs, err = client.JobList(ctx, NewJobListParams().State(rivertype.JobStateCompleted).After(JobListCursorFromJob(job5))) + res, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateCompleted).After(JobListCursorFromJob(job5, JobListOrderByTime))) require.NoError(t, err) - require.Equal(t, []int64{job6.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job6.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, job6.ID, res.LastCursor.id) + + res, err = client.JobList(ctx, NewJobListParams().OrderBy(JobListOrderByScheduledAt, SortOrderAsc).After(JobListCursorFromJob(job4, JobListOrderByScheduledAt))) + require.NoError(t, err) + require.Equal(t, []int64{job1.ID, job3.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, JobListOrderByScheduledAt, res.LastCursor.sortField) + require.Equal(t, job2.ID, res.LastCursor.id) }) t.Run("MetadataOnly", func(t *testing.T) { @@ -1619,14 +1628,14 @@ func Test_Client_JobList(t *testing.T) { job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Metadata: []byte(`{"baz": "value"}`)}) job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Metadata: []byte(`{"baz": "value"}`)}) - jobs, err := client.JobList(ctx, NewJobListParams().State("").Metadata(`{"foo": "bar"}`)) + res, err := client.JobList(ctx, NewJobListParams().Metadata(`{"foo": "bar"}`)) require.NoError(t, err) - require.Equal(t, []int64{job1.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job1.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State("").Metadata(`{"baz": "value"}`).OrderBy(JobListOrderByTime, SortOrderDesc)) + res, err = client.JobList(ctx, NewJobListParams().Metadata(`{"baz": "value"}`).OrderBy(JobListOrderByTime, SortOrderDesc)) require.NoError(t, err) // Sort order was explicitly reversed: - require.Equal(t, []int64{job3.ID, job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, []int64{job3.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) t.Run("WithCancelledContext", func(t *testing.T) { @@ -1637,9 +1646,9 @@ func Test_Client_JobList(t *testing.T) { ctx, cancel := context.WithCancel(ctx) cancel() // cancel immediately - jobs, err := client.JobList(ctx, NewJobListParams().State(JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.ErrorIs(t, context.Canceled, err) - require.Empty(t, jobs) + require.Nil(t, res) }) } diff --git a/internal/dblist/db_list.go b/internal/dblist/db_list.go index 7050f5e7..79f7f07a 100644 --- a/internal/dblist/db_list.go +++ b/internal/dblist/db_list.go @@ -6,6 +6,8 @@ import ( "fmt" "strings" + "github.com/lib/pq" + "github.com/riverqueue/river/riverdriver" "github.com/riverqueue/river/rivertype" ) @@ -42,7 +44,7 @@ type JobListParams struct { OrderBy []JobListOrderBy Priorities []int16 Queues []string - State rivertype.JobState + States []rivertype.JobState } func JobList(ctx context.Context, exec riverdriver.Executor, params *JobListParams) ([]*rivertype.JobRow, error) { @@ -81,10 +83,10 @@ func JobList(ctx context.Context, exec riverdriver.Executor, params *JobListPara namedArgs["queues"] = params.Queues } - if params.State != "" { + if len(params.States) > 0 { writeWhereOrAnd() - conditionsBuilder.WriteString("state = @state::river_job_state") - namedArgs["state"] = params.State + conditionsBuilder.WriteString("state = any(@states::river_job_state[])") + namedArgs["states"] = pq.Array(params.States) } if params.Conditions != "" { diff --git a/internal/dblist/db_list_test.go b/internal/dblist/db_list_test.go index 9e25ae32..2d2f2929 100644 --- a/internal/dblist/db_list_test.go +++ b/internal/dblist/db_list_test.go @@ -40,7 +40,7 @@ func TestJobListNoJobs(t *testing.T) { bundle := setup() _, err := JobList(ctx, bundle.exec, &JobListParams{ - State: rivertype.JobStateCompleted, + States: []rivertype.JobState{rivertype.JobStateCompleted}, LimitCount: 1, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderAsc}}, }) @@ -55,7 +55,7 @@ func TestJobListNoJobs(t *testing.T) { _, err := JobList(ctx, bundle.exec, &JobListParams{ Conditions: "queue = 'test' AND priority = 1 AND args->>'foo' = @foo", NamedArgs: pgx.NamedArgs{"foo": "bar"}, - State: rivertype.JobStateCompleted, + States: []rivertype.JobState{rivertype.JobStateCompleted}, LimitCount: 1, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderAsc}}, }) @@ -123,7 +123,7 @@ func TestJobListWithJobs(t *testing.T) { params := &JobListParams{ LimitCount: 3, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, - State: rivertype.JobStateAvailable, + States: []rivertype.JobState{rivertype.JobStateAvailable}, } execTest(ctx, t, bundle, params, func(jobs []*rivertype.JobRow, err error) { @@ -150,7 +150,7 @@ func TestJobListWithJobs(t *testing.T) { LimitCount: 2, NamedArgs: map[string]any{"paths1": []string{"job_num"}, "value1": 2}, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, - State: rivertype.JobStateAvailable, + States: []rivertype.JobState{rivertype.JobStateAvailable}, } execTest(ctx, t, bundle, params, func(jobs []*rivertype.JobRow, err error) { @@ -172,7 +172,7 @@ func TestJobListWithJobs(t *testing.T) { LimitCount: 2, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, Kinds: []string{"alternate_kind"}, - State: rivertype.JobStateAvailable, + States: []rivertype.JobState{rivertype.JobStateAvailable}, } execTest(ctx, t, bundle, params, func(jobs []*rivertype.JobRow, err error) { @@ -194,7 +194,7 @@ func TestJobListWithJobs(t *testing.T) { LimitCount: 2, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, Queues: []string{"priority"}, - State: rivertype.JobStateAvailable, + States: []rivertype.JobState{rivertype.JobStateAvailable}, } execTest(ctx, t, bundle, params, func(jobs []*rivertype.JobRow, err error) { diff --git a/job_list_params.go b/job_list_params.go index 14c540ca..2881be74 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -15,19 +15,34 @@ import ( // JobListCursor is used to specify a starting point for a paginated // job list query. type JobListCursor struct { - id int64 - kind string - queue string - time time.Time + id int64 + kind string + queue string + sortField JobListOrderByField + time time.Time } // JobListCursorFromJob creates a JobListCursor from a JobRow. -func JobListCursorFromJob(job *rivertype.JobRow) *JobListCursor { +func JobListCursorFromJob(job *rivertype.JobRow, sortField JobListOrderByField) *JobListCursor { + time := job.CreatedAt + switch sortField { + case JobListOrderByTime: + time = jobListTimeValue(job) + case JobListOrderByFinalizedAt: + if job.FinalizedAt != nil { + time = *job.FinalizedAt + } + case JobListOrderByScheduledAt: + time = job.ScheduledAt + default: + panic("invalid sort field") + } return &JobListCursor{ - id: job.ID, - kind: job.Kind, - queue: job.Queue, - time: jobListTimeValue(job), + id: job.ID, + kind: job.Kind, + queue: job.Queue, + sortField: sortField, + time: time, } } @@ -46,10 +61,11 @@ func (c *JobListCursor) UnmarshalText(text []byte) error { return err } *c = JobListCursor{ - id: wrapperValue.ID, - kind: wrapperValue.Kind, - queue: wrapperValue.Queue, - time: wrapperValue.Time, + id: wrapperValue.ID, + kind: wrapperValue.Kind, + queue: wrapperValue.Queue, + sortField: JobListOrderByField(wrapperValue.SortField), + time: wrapperValue.Time, } return nil } @@ -58,10 +74,11 @@ func (c *JobListCursor) UnmarshalText(text []byte) error { // opaque string. func (c JobListCursor) MarshalText() ([]byte, error) { wrapperValue := jobListPaginationCursorJSON{ - ID: c.id, - Kind: c.kind, - Queue: c.queue, - Time: c.time, + ID: c.id, + Kind: c.kind, + Queue: c.queue, + SortField: string(c.sortField), + Time: c.time, } data, err := json.Marshal(wrapperValue) if err != nil { @@ -73,10 +90,11 @@ func (c JobListCursor) MarshalText() ([]byte, error) { } type jobListPaginationCursorJSON struct { - ID int64 `json:"id"` - Kind string `json:"kind"` - Queue string `json:"queue"` - Time time.Time `json:"time"` + ID int64 `json:"id"` + Kind string `json:"kind"` + Queue string `json:"queue"` + SortField string `json:"sort_field"` + Time time.Time `json:"time"` } // SortOrder specifies the direction of a sort. @@ -90,12 +108,20 @@ const ( ) // JobListOrderByField specifies the field to sort by. -type JobListOrderByField int +type JobListOrderByField string const ( + // JobListOrderByFinalizedAt specifies that the sort should be by + // finalized_at. + // + // This option must be used in conjunction with filtering by only finalized + // job states. + JobListOrderByFinalizedAt JobListOrderByField = "finalized_at" + // JobListOrderByScheduledAt specifies that the sort should be by scheduled_at. + JobListOrderByScheduledAt JobListOrderByField = "scheduled_at" // JobListOrderByTime specifies that the sort should be by time. The specific - // time field used will vary by job state. - JobListOrderByTime JobListOrderByField = iota + // time field used will vary by the first specified job state. + JobListOrderByTime JobListOrderByField = "time" ) // JobListParams specifies the parameters for a JobList query. It must be @@ -107,11 +133,12 @@ type JobListParams struct { after *JobListCursor kinds []string metadataFragment string + overrodeState bool paginationCount int32 queues []string sortField JobListOrderByField sortOrder SortOrder - state rivertype.JobState + states []rivertype.JobState } // NewJobListParams creates a new JobListParams to return available jobs sorted @@ -121,7 +148,15 @@ func NewJobListParams() *JobListParams { paginationCount: 100, sortField: JobListOrderByTime, sortOrder: SortOrderAsc, - state: rivertype.JobStateAvailable, + states: []rivertype.JobState{ + JobStateAvailable, + JobStateCancelled, + JobStateCompleted, + JobStateDiscarded, + JobStateRetryable, + JobStateRunning, + JobStateScheduled, + }, } } @@ -130,11 +165,12 @@ func (p *JobListParams) copy() *JobListParams { after: p.after, kinds: append([]string(nil), p.kinds...), metadataFragment: p.metadataFragment, + overrodeState: p.overrodeState, paginationCount: p.paginationCount, queues: append([]string(nil), p.queues...), sortField: p.sortField, sortOrder: p.sortOrder, - state: p.state, + states: append([]rivertype.JobState(nil), p.states...), } } @@ -154,10 +190,30 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { return nil, errors.New("invalid sort order") } - if p.sortField != JobListOrderByTime { - return nil, errors.New("invalid sort field") + if p.sortField == JobListOrderByFinalizedAt { + currentNonFinalizedStates := make([]rivertype.JobState, 0, len(p.states)) + for _, state := range p.states { + //nolint:exhaustive + switch state { + case JobStateCancelled, JobStateCompleted, JobStateDiscarded: + default: + currentNonFinalizedStates = append(currentNonFinalizedStates, state) + } + } + // This indicates the user overrode the States list with only non-finalized + // states prior to then requesting FinalizedAt ordering. + if len(currentNonFinalizedStates) == 0 { + return nil, fmt.Errorf("cannot order by finalized_at with non-finalized state filters %+v", currentNonFinalizedStates) + } + } + + var timeField string + if len(p.states) > 0 && p.sortField == JobListOrderByTime { + timeField = jobListTimeFieldForState(p.states[0]) + } else { + timeField = string(p.sortField) } - timeField := jobListTimeFieldForState(p.state) + orderBy = append(orderBy, []dblist.JobListOrderBy{ {Expr: timeField, Order: sortOrder}, {Expr: "id", Order: sortOrder}, @@ -193,7 +249,7 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { OrderBy: orderBy, Priorities: nil, Queues: p.queues, - State: p.state, + States: p.states, } return dbParams, nil @@ -249,18 +305,38 @@ func (p *JobListParams) Queues(queues ...string) *JobListParams { // OrderBy returns an updated filter set that will sort the results using the // specified field and direction. +// +// If ordering by FinalizedAt, the States filter will be set to only include +// finalized job states unless it has already been overridden. func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) *JobListParams { result := p.copy() + switch field { + case JobListOrderByTime, JobListOrderByScheduledAt: + result.sortField = field + case JobListOrderByFinalizedAt: + result.sortField = field + if !p.overrodeState { + result.states = []rivertype.JobState{ + JobStateCancelled, + JobStateCompleted, + JobStateDiscarded, + } + } + default: + panic("invalid order by field") + } result.sortField = field result.sortOrder = direction return result } -// State returns an updated filter set that will only return jobs in the given -// state. -func (p *JobListParams) State(state rivertype.JobState) *JobListParams { +// States returns an updated filter set that will only return jobs in the given +// states. +func (p *JobListParams) States(states ...rivertype.JobState) *JobListParams { result := p.copy() - result.state = state + result.states = make([]rivertype.JobState, len(states)) + result.overrodeState = true + copy(result.states, states) return result } diff --git a/job_list_params_test.go b/job_list_params_test.go index 8aa9c474..f8533167 100644 --- a/job_list_params_test.go +++ b/job_list_params_test.go @@ -35,7 +35,7 @@ func Test_JobListCursor_JobListCursorFromJob(t *testing.T) { ScheduledAt: now.Add(-10 * time.Second), } - cursor := JobListCursorFromJob(jobRow) + cursor := JobListCursorFromJob(jobRow, JobListOrderByTime) require.Equal(t, jobRow.ID, cursor.id) require.Equal(t, jobRow.Kind, cursor.kind) require.Equal(t, jobRow.Queue, cursor.queue) @@ -65,7 +65,7 @@ func Test_JobListCursor_JobListCursorFromJob(t *testing.T) { ScheduledAt: now.Add(-10 * time.Second), } - cursor := JobListCursorFromJob(jobRow) + cursor := JobListCursorFromJob(jobRow, JobListOrderByTime) require.Equal(t, jobRow.ID, cursor.id) require.Equal(t, jobRow.Kind, cursor.kind) require.Equal(t, jobRow.Queue, cursor.queue) @@ -87,7 +87,7 @@ func Test_JobListCursor_JobListCursorFromJob(t *testing.T) { ScheduledAt: now.Add(-10 * time.Second), } - cursor := JobListCursorFromJob(jobRow) + cursor := JobListCursorFromJob(jobRow, JobListOrderByTime) require.Equal(t, jobRow.ID, cursor.id) require.Equal(t, jobRow.Kind, cursor.kind) require.Equal(t, jobRow.Queue, cursor.queue) @@ -107,7 +107,7 @@ func Test_JobListCursor_JobListCursorFromJob(t *testing.T) { ScheduledAt: now.Add(-10 * time.Second), } - cursor := JobListCursorFromJob(jobRow) + cursor := JobListCursorFromJob(jobRow, JobListOrderByTime) require.Equal(t, jobRow.ID, cursor.id) require.Equal(t, jobRow.Kind, cursor.kind) require.Equal(t, jobRow.Queue, cursor.queue)