From 3c171a16c93ee2a19f8c8656546f01754a295e5d Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sun, 25 Feb 2024 20:45:48 -0800 Subject: [PATCH 1/7] Change JobList ergnomics - add filtering by states, return cursor from function --- client.go | 29 ++++++-- client_test.go | 63 +++++++++-------- internal/dblist/db_list.go | 9 +-- internal/dblist/db_list_test.go | 12 ++-- job_list_params.go | 116 ++++++++++++++++++++++---------- job_list_params_test.go | 8 +-- 6 files changed, 154 insertions(+), 83 deletions(-) diff --git a/client.go b/client.go index d5aea11a..c74c1fbb 100644 --- a/client.go +++ b/client.go @@ -1455,9 +1455,9 @@ 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) ([]*rivertype.JobRow, *JobListCursor, error) { if !c.driver.HasPool() { - return nil, errNoDriverDBPool + return nil, nil, errNoDriverDBPool } if params == nil { @@ -1465,10 +1465,17 @@ func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) ([]*ri } dbParams, err := params.toDBParams() if err != nil { - return nil, err + return nil, nil, err } - return dblist.JobList(ctx, c.driver.GetExecutor(), dbParams) + res, err := dblist.JobList(ctx, c.driver.GetExecutor(), dbParams) + if err != nil { + return nil, nil, err + } + if len(res) > 0 { + return res, JobListCursorFromJob(res[len(res)-1], params.sortField), nil + } + return res, nil, nil } // JobListTx returns a paginated list of jobs matching the provided filters. The @@ -1480,16 +1487,24 @@ 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) ([]*rivertype.JobRow, *JobListCursor, error) { if params == nil { params = NewJobListParams() } + dbParams, err := params.toDBParams() if err != nil { - return nil, err + return nil, nil, err } - return dblist.JobList(ctx, c.driver.UnwrapExecutor(tx), dbParams) + res, err := dblist.JobList(ctx, c.driver.UnwrapExecutor(tx), dbParams) + if err != nil { + return nil, nil, err + } + if len(res) > 0 { + return res, JobListCursorFromJob(res[len(res)-1], params.sortField), nil + } + return res, nil, nil } // PeriodicJobs returns the currently configured set of periodic jobs for the diff --git a/client_test.go b/client_test.go index d5feac21..1cac1b4f 100644 --- a/client_test.go +++ b/client_test.go @@ -702,7 +702,7 @@ func Test_Client_Stop(t *testing.T) { require.NoError(t, client.Stop(ctx)) - runningJobs, err := client.JobList(ctx, NewJobListParams().State(rivertype.JobStateRunning)) + runningJobs, _, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.NoError(t, err) require.Empty(t, runningJobs, "expected no jobs to be left running") }) @@ -1449,12 +1449,12 @@ 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")) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().Kinds("test_kind_2")) + jobs, _, 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 })) }) @@ -1468,12 +1468,12 @@ 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")) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().Queues("queue_2")) + jobs, _, 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 })) }) @@ -1487,12 +1487,12 @@ 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)) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().State(JobStateRunning)) + jobs, _, 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 })) }) @@ -1513,11 +1513,11 @@ 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)) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().State(state).OrderBy(JobListOrderByTime, SortOrderDesc)) + jobs, _, 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 })) } @@ -1539,11 +1539,11 @@ 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)) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().State(state).OrderBy(JobListOrderByTime, SortOrderDesc)) + jobs, _, 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 })) } @@ -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)) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().State(JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc)) + jobs, _, 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 })) }) - 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) + jobs, _, 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(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))) + jobs, cursor, 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, cursor.id, job2.ID) - jobs, err = client.JobList(ctx, NewJobListParams().State(rivertype.JobStateRunning).After(JobListCursorFromJob(job3))) + jobs, cursor, 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, cursor.id, job4.ID) - jobs, err = client.JobList(ctx, NewJobListParams().State(rivertype.JobStateCompleted).After(JobListCursorFromJob(job5))) + jobs, cursor, 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, cursor.id, job6.ID) + + jobs, cursor, 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(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, cursor.sortField, JobListOrderByScheduledAt) + require.Equal(t, cursor.id, job2.ID) }) t.Run("MetadataOnly", func(t *testing.T) { @@ -1619,11 +1628,11 @@ 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"}`)) + jobs, _, 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 })) - jobs, err = client.JobList(ctx, NewJobListParams().State("").Metadata(`{"baz": "value"}`).OrderBy(JobListOrderByTime, SortOrderDesc)) + jobs, _, 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 })) @@ -1637,7 +1646,7 @@ func Test_Client_JobList(t *testing.T) { ctx, cancel := context.WithCancel(ctx) cancel() // cancel immediately - jobs, err := client.JobList(ctx, NewJobListParams().State(JobStateRunning)) + jobs, _, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.ErrorIs(t, context.Canceled, err) require.Empty(t, jobs) }) diff --git a/internal/dblist/db_list.go b/internal/dblist/db_list.go index 7050f5e7..8e6f9875 100644 --- a/internal/dblist/db_list.go +++ b/internal/dblist/db_list.go @@ -6,6 +6,7 @@ import ( "fmt" "strings" + "github.com/lib/pq" "github.com/riverqueue/river/riverdriver" "github.com/riverqueue/river/rivertype" ) @@ -42,7 +43,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 +82,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..83fd2d36 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -15,19 +15,36 @@ 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 JobListOrderByAttemptedAt: + if job.AttemptedAt != nil { + time = *job.AttemptedAt + } + case JobListOrderByFinalizedAt: + if job.FinalizedAt != nil { + time = *job.FinalizedAt + } + case JobListOrderByScheduledAt: + time = job.ScheduledAt + } 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 +63,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 +76,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 +92,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 +110,20 @@ const ( ) // JobListOrderByField specifies the field to sort by. -type JobListOrderByField int +type JobListOrderByField string const ( // 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" + // JobListOrderByCreatedAt specifies that the sort should be by created_at. + JobListOrderByCreatedAt JobListOrderByField = "created_at" + // JobListOrderByScheduledAt specifies that the sort should be by scheduled_at. + JobListOrderByScheduledAt JobListOrderByField = "scheduled_at" + // JobListOrderByAttemptedAt specifies that the sort should be by attempted_at. + JobListOrderByAttemptedAt JobListOrderByField = "attempted_at" + // JobListOrderByFinalizedAt specifies that the sort should be by finalized_at. + JobListOrderByFinalizedAt JobListOrderByField = "finalized_at" ) // JobListParams specifies the parameters for a JobList query. It must be @@ -111,7 +139,7 @@ type JobListParams struct { queues []string sortField JobListOrderByField sortOrder SortOrder - state rivertype.JobState + states []rivertype.JobState } // NewJobListParams creates a new JobListParams to return available jobs sorted @@ -121,7 +149,15 @@ func NewJobListParams() *JobListParams { paginationCount: 100, sortField: JobListOrderByTime, sortOrder: SortOrderAsc, - state: rivertype.JobStateAvailable, + states: []rivertype.JobState{ + JobStateAvailable, + JobStateCancelled, + JobStateCompleted, + JobStateDiscarded, + JobStateRetryable, + JobStateRunning, + JobStateScheduled, + }, } } @@ -134,7 +170,7 @@ func (p *JobListParams) copy() *JobListParams { 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,13 @@ 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") + timeField := "created_at" + 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 +232,7 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { OrderBy: orderBy, Priorities: nil, Queues: p.queues, - State: p.state, + States: p.states, } return dbParams, nil @@ -251,16 +290,23 @@ func (p *JobListParams) Queues(queues ...string) *JobListParams { // specified field and direction. func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) *JobListParams { result := p.copy() + switch field { + case JobListOrderByTime, JobListOrderByCreatedAt, JobListOrderByScheduledAt, JobListOrderByAttemptedAt, JobListOrderByFinalizedAt: + result.sortField = field + 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)) + 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) From e0f3efd546720ab10538ae0d124fb1787ae2f491 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Mon, 26 Feb 2024 09:25:04 -0800 Subject: [PATCH 2/7] Update JobList response type to be combined struct of result and cursor --- client.go | 39 +++++++++++++-------- client_test.go | 94 +++++++++++++++++++++++++------------------------- 2 files changed, 71 insertions(+), 62 deletions(-) diff --git a/client.go b/client.go index c74c1fbb..0fe639e0 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 + Cursor *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,9 +1462,9 @@ func validateQueueName(queueName string) error { // if err != nil { // // handle error // } -func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) ([]*rivertype.JobRow, *JobListCursor, error) { +func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) (*JobListResult, error) { if !c.driver.HasPool() { - return nil, nil, errNoDriverDBPool + return nil, errNoDriverDBPool } if params == nil { @@ -1465,17 +1472,18 @@ func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) ([]*ri } dbParams, err := params.toDBParams() if err != nil { - return nil, nil, err + return nil, err } - res, err := dblist.JobList(ctx, c.driver.GetExecutor(), dbParams) + jobs, err := dblist.JobList(ctx, c.driver.GetExecutor(), dbParams) if err != nil { - return nil, nil, err + return nil, err } - if len(res) > 0 { - return res, JobListCursorFromJob(res[len(res)-1], params.sortField), nil + res := &JobListResult{Jobs: jobs} + if len(jobs) > 0 { + res.Cursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) } - return res, nil, nil + return res, nil } // JobListTx returns a paginated list of jobs matching the provided filters. The @@ -1487,24 +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, *JobListCursor, 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, nil, err + return nil, err } - res, err := dblist.JobList(ctx, c.driver.UnwrapExecutor(tx), dbParams) + jobs, err := dblist.JobList(ctx, c.driver.UnwrapExecutor(tx), dbParams) if err != nil { - return nil, nil, err + return nil, err } - if len(res) > 0 { - return res, JobListCursorFromJob(res[len(res)-1], params.sortField), nil + res := &JobListResult{Jobs: jobs} + if len(jobs) > 0 { + res.Cursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) } - return res, nil, nil + 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 1cac1b4f..fbfb3697 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().States(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().States(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().States(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().States(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().States(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().States(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().States(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,14 +1558,14 @@ 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().States(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().States(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("WithNilParamsFiltersToAllStatesByDefault", func(t *testing.T) { @@ -1578,10 +1578,10 @@ func Test_Client_JobList(t *testing.T) { job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable), ScheduledAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) 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, job3.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) { @@ -1597,26 +1597,26 @@ func Test_Client_JobList(t *testing.T) { 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, cursor, err := client.JobList(ctx, NewJobListParams().States(rivertype.JobStateAvailable).After(JobListCursorFromJob(job1, JobListOrderByTime))) + 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, cursor.id, job2.ID) + require.Equal(t, []int64{job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, res.Cursor.id, job2.ID) - jobs, cursor, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning).After(JobListCursorFromJob(job3, JobListOrderByTime))) + 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, cursor.id, job4.ID) + require.Equal(t, []int64{job4.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, res.Cursor.id, job4.ID) - jobs, cursor, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateCompleted).After(JobListCursorFromJob(job5, JobListOrderByTime))) + 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, cursor.id, job6.ID) + require.Equal(t, []int64{job6.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) + require.Equal(t, res.Cursor.id, job6.ID) - jobs, cursor, err = client.JobList(ctx, NewJobListParams().OrderBy(JobListOrderByScheduledAt, SortOrderAsc).After(JobListCursorFromJob(job4, JobListOrderByScheduledAt))) + 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(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, cursor.sortField, JobListOrderByScheduledAt) - require.Equal(t, cursor.id, job2.ID) + 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, res.Cursor.sortField, JobListOrderByScheduledAt) + require.Equal(t, res.Cursor.id, job2.ID) }) t.Run("MetadataOnly", func(t *testing.T) { @@ -1628,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().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().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) { @@ -1646,9 +1646,9 @@ func Test_Client_JobList(t *testing.T) { ctx, cancel := context.WithCancel(ctx) cancel() // cancel immediately - jobs, _, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) require.ErrorIs(t, context.Canceled, err) - require.Empty(t, jobs) + require.Nil(t, res) }) } From 5a915cfdc7c7098428bd3a33ed124e0c1014fe6f Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Mon, 26 Feb 2024 09:38:28 -0800 Subject: [PATCH 3/7] golangci-lint --- client_test.go | 10 +++++----- internal/dblist/db_list.go | 1 + job_list_params.go | 5 ++++- 3 files changed, 10 insertions(+), 6 deletions(-) diff --git a/client_test.go b/client_test.go index fbfb3697..69c06485 100644 --- a/client_test.go +++ b/client_test.go @@ -1600,23 +1600,23 @@ func Test_Client_JobList(t *testing.T) { 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(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, res.Cursor.id, job2.ID) + require.Equal(t, job2.ID, res.Cursor.id) 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(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, res.Cursor.id, job4.ID) + require.Equal(t, job4.ID, res.Cursor.id) 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(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, res.Cursor.id, job6.ID) + require.Equal(t, job6.ID, res.Cursor.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, res.Cursor.sortField, JobListOrderByScheduledAt) - require.Equal(t, res.Cursor.id, job2.ID) + require.Equal(t, JobListOrderByScheduledAt, res.Cursor.sortField) + require.Equal(t, job2.ID, res.Cursor.id) }) t.Run("MetadataOnly", func(t *testing.T) { diff --git a/internal/dblist/db_list.go b/internal/dblist/db_list.go index 8e6f9875..79f7f07a 100644 --- a/internal/dblist/db_list.go +++ b/internal/dblist/db_list.go @@ -7,6 +7,7 @@ import ( "strings" "github.com/lib/pq" + "github.com/riverqueue/river/riverdriver" "github.com/riverqueue/river/rivertype" ) diff --git a/job_list_params.go b/job_list_params.go index 83fd2d36..b9710f6a 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -38,6 +38,9 @@ func JobListCursorFromJob(job *rivertype.JobRow, sortField JobListOrderByField) } case JobListOrderByScheduledAt: time = job.ScheduledAt + case JobListOrderByCreatedAt: + default: + // stick with created_at } return &JobListCursor{ id: job.ID, @@ -190,7 +193,7 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { return nil, errors.New("invalid sort order") } - timeField := "created_at" + var timeField string if len(p.states) > 0 && p.sortField == JobListOrderByTime { timeField = jobListTimeFieldForState(p.states[0]) } else { From 2326ee7ca3196ab6d91e581db48ad8c28d35a8ad Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 15 Mar 2024 08:59:24 -0700 Subject: [PATCH 4/7] Sort JobListOrderBy consts, rename Cursor to LastCursor --- client.go | 8 ++++---- job_list_params.go | 14 +++++++------- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/client.go b/client.go index 0fe639e0..63dd9479 100644 --- a/client.go +++ b/client.go @@ -1449,8 +1449,8 @@ func validateQueueName(queueName string) error { // 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 - Cursor *JobListCursor + Jobs []*rivertype.JobRow + LastCursor *JobListCursor } // JobList returns a paginated list of jobs matching the provided filters. The @@ -1481,7 +1481,7 @@ func (c *Client[TTx]) JobList(ctx context.Context, params *JobListParams) (*JobL } res := &JobListResult{Jobs: jobs} if len(jobs) > 0 { - res.Cursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) + res.LastCursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) } return res, nil } @@ -1511,7 +1511,7 @@ func (c *Client[TTx]) JobListTx(ctx context.Context, tx TTx, params *JobListPara } res := &JobListResult{Jobs: jobs} if len(jobs) > 0 { - res.Cursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) + res.LastCursor = JobListCursorFromJob(jobs[len(jobs)-1], params.sortField) } return res, nil } diff --git a/job_list_params.go b/job_list_params.go index b9710f6a..037bfa05 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -116,17 +116,17 @@ const ( type JobListOrderByField string const ( - // JobListOrderByTime specifies that the sort should be by time. The specific - // time field used will vary by the first specified job state. - JobListOrderByTime JobListOrderByField = "time" - // JobListOrderByCreatedAt specifies that the sort should be by created_at. - JobListOrderByCreatedAt JobListOrderByField = "created_at" - // JobListOrderByScheduledAt specifies that the sort should be by scheduled_at. - JobListOrderByScheduledAt JobListOrderByField = "scheduled_at" // JobListOrderByAttemptedAt specifies that the sort should be by attempted_at. JobListOrderByAttemptedAt JobListOrderByField = "attempted_at" + // JobListOrderByCreatedAt specifies that the sort should be by created_at. + JobListOrderByCreatedAt JobListOrderByField = "created_at" // JobListOrderByFinalizedAt specifies that the sort should be by finalized_at. 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 the first specified job state. + JobListOrderByTime JobListOrderByField = "time" ) // JobListParams specifies the parameters for a JobList query. It must be From bb56e9f4cc04e9ee565f28ba806ab4cfb67789a3 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 15 Mar 2024 09:04:52 -0700 Subject: [PATCH 5/7] Add unreleased changelog --- CHANGELOG.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index b8793d8a..1bdf2e73 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -46,6 +46,8 @@ Although it comes with a number of improvements, there's nothing particularly no - River uses a new job completer that batches up completion work so that large numbers of them can be performed more efficiently. In a purely synthetic (i.e. mostly unrealistic) benchmark, River's job throughput increases ~4.5x. [PR #258](https://github.com/riverqueue/river/pull/258). - Changed default client IDs to be a combination of hostname and the time which the client started. This can still be changed by specifying `Config.ID`. [PR #255](https://github.com/riverqueue/river/pull/255). - Notifier refactored for better robustness and testability. [PR #253](https://github.com/riverqueue/river/pull/253). +- JobList/JobListTx now support querying Jobs by a list of Job Kinds and States (breaking change). Also allows for filtering by specific timestamp values. [PR #236](https://github.com/riverqueue/river/pull/236). + ## [0.0.25] - 2024-03-01 From 94099762a6585b198bd2772ae1c67c96fd508071 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 15 Mar 2024 09:06:21 -0700 Subject: [PATCH 6/7] Update changelog, cursor -> lastcursor --- client_test.go | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/client_test.go b/client_test.go index 69c06485..59d7bb5a 100644 --- a/client_test.go +++ b/client_test.go @@ -1600,23 +1600,23 @@ func Test_Client_JobList(t *testing.T) { 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(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, job2.ID, res.Cursor.id) + require.Equal(t, job2.ID, res.LastCursor.id) 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(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, job4.ID, res.Cursor.id) + require.Equal(t, job4.ID, res.LastCursor.id) 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(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - require.Equal(t, job6.ID, res.Cursor.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.Cursor.sortField) - require.Equal(t, job2.ID, res.Cursor.id) + require.Equal(t, JobListOrderByScheduledAt, res.LastCursor.sortField) + require.Equal(t, job2.ID, res.LastCursor.id) }) t.Run("MetadataOnly", func(t *testing.T) { From f61b83b0f32debd06792df1f4dda37f3b960c6d4 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Thu, 28 Mar 2024 23:43:55 -0500 Subject: [PATCH 7/7] remove inefficient created, attempted filters. Limits on finalized River does not keep indexes on `created_at` and `attempted_at` fields and we do not want to do so for performance reasons. We shouldn't offer a built-in API for sorting by these fields which will be guaranteed to have poor performance on large tables. Additionally, while `finalized_at` _is_ indexed, it's indexed via a _partial_ index on non-null values (i.e. only for finalized jobs). We can still make it possible to order based upon this field, but only when also filtering to only finalized states. This still _may_ be inefficient at times, but it should hopefully be able to at least avoid a full table scan. --- CHANGELOG.md | 6 ++++-- job_list_params.go | 51 +++++++++++++++++++++++++++++++++++----------- 2 files changed, 43 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1bdf2e73..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 @@ -46,8 +50,6 @@ Although it comes with a number of improvements, there's nothing particularly no - River uses a new job completer that batches up completion work so that large numbers of them can be performed more efficiently. In a purely synthetic (i.e. mostly unrealistic) benchmark, River's job throughput increases ~4.5x. [PR #258](https://github.com/riverqueue/river/pull/258). - Changed default client IDs to be a combination of hostname and the time which the client started. This can still be changed by specifying `Config.ID`. [PR #255](https://github.com/riverqueue/river/pull/255). - Notifier refactored for better robustness and testability. [PR #253](https://github.com/riverqueue/river/pull/253). -- JobList/JobListTx now support querying Jobs by a list of Job Kinds and States (breaking change). Also allows for filtering by specific timestamp values. [PR #236](https://github.com/riverqueue/river/pull/236). - ## [0.0.25] - 2024-03-01 diff --git a/job_list_params.go b/job_list_params.go index 037bfa05..2881be74 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -28,19 +28,14 @@ func JobListCursorFromJob(job *rivertype.JobRow, sortField JobListOrderByField) switch sortField { case JobListOrderByTime: time = jobListTimeValue(job) - case JobListOrderByAttemptedAt: - if job.AttemptedAt != nil { - time = *job.AttemptedAt - } case JobListOrderByFinalizedAt: if job.FinalizedAt != nil { time = *job.FinalizedAt } case JobListOrderByScheduledAt: time = job.ScheduledAt - case JobListOrderByCreatedAt: default: - // stick with created_at + panic("invalid sort field") } return &JobListCursor{ id: job.ID, @@ -116,11 +111,11 @@ const ( type JobListOrderByField string const ( - // JobListOrderByAttemptedAt specifies that the sort should be by attempted_at. - JobListOrderByAttemptedAt JobListOrderByField = "attempted_at" - // JobListOrderByCreatedAt specifies that the sort should be by created_at. - JobListOrderByCreatedAt JobListOrderByField = "created_at" - // JobListOrderByFinalizedAt specifies that the sort should be by finalized_at. + // 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" @@ -138,6 +133,7 @@ type JobListParams struct { after *JobListCursor kinds []string metadataFragment string + overrodeState bool paginationCount int32 queues []string sortField JobListOrderByField @@ -169,6 +165,7 @@ 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, @@ -193,6 +190,23 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { return nil, errors.New("invalid sort order") } + 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]) @@ -291,11 +305,23 @@ 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, JobListOrderByCreatedAt, JobListOrderByScheduledAt, JobListOrderByAttemptedAt, JobListOrderByFinalizedAt: + 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") } @@ -309,6 +335,7 @@ func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) func (p *JobListParams) States(states ...rivertype.JobState) *JobListParams { result := p.copy() result.states = make([]rivertype.JobState, len(states)) + result.overrodeState = true copy(result.states, states) return result }