From 6893b969697d300e55ef18967078a109f42567c9 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sun, 25 Feb 2024 20:45:48 -0800 Subject: [PATCH 1/6] 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 395b4d06..7db1a2b7 100644 --- a/client.go +++ b/client.go @@ -1360,9 +1360,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 { @@ -1370,10 +1370,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 @@ -1385,14 +1392,22 @@ 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 } diff --git a/client_test.go b/client_test.go index 62d0ada6..9be0e49c 100644 --- a/client_test.go +++ b/client_test.go @@ -640,7 +640,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") }) @@ -1387,12 +1387,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 })) }) @@ -1406,12 +1406,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 })) }) @@ -1425,12 +1425,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 })) }) @@ -1451,11 +1451,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 })) } @@ -1477,11 +1477,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 })) } @@ -1496,17 +1496,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) @@ -1514,12 +1514,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) { @@ -1530,22 +1530,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) { @@ -1557,11 +1566,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 })) @@ -1575,7 +1584,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 dc8ce82cff94789c71cfb1f609ec69ae880c9c3e Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Mon, 26 Feb 2024 09:25:04 -0800 Subject: [PATCH 2/6] 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 7db1a2b7..e4860f6d 100644 --- a/client.go +++ b/client.go @@ -1351,6 +1351,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. @@ -1360,9 +1367,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 { @@ -1370,17 +1377,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 @@ -1392,22 +1400,23 @@ 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 } diff --git a/client_test.go b/client_test.go index 9be0e49c..c9806901 100644 --- a/client_test.go +++ b/client_test.go @@ -640,9 +640,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) { @@ -1387,14 +1387,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 @@ -1406,14 +1406,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) { @@ -1425,14 +1425,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) { @@ -1451,13 +1451,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 })) } }) @@ -1477,13 +1477,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 })) } }) @@ -1496,14 +1496,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) { @@ -1516,10 +1516,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) { @@ -1535,26 +1535,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) { @@ -1566,14 +1566,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) { @@ -1584,9 +1584,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 0169106c0bc99575694b7a61b285c3e0e9a79921 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Mon, 26 Feb 2024 09:38:28 -0800 Subject: [PATCH 3/6] 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 c9806901..c8b64074 100644 --- a/client_test.go +++ b/client_test.go @@ -1538,23 +1538,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 74ad2874fd41b8feb1ceb45fb4d8075a94836b76 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 15 Mar 2024 08:59:24 -0700 Subject: [PATCH 4/6] 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 e4860f6d..5536f9f6 100644 --- a/client.go +++ b/client.go @@ -1354,8 +1354,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 @@ -1386,7 +1386,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 } @@ -1416,7 +1416,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 cdbc529952871f277986acd8a8d0f605371eba7e Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 15 Mar 2024 09:04:52 -0700 Subject: [PATCH 5/6] Add unreleased changelog --- CHANGELOG.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index c01cd288..31fb77cc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - 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 49f1208ab46b7690e2bb093a37d96c11ce58e30e Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 15 Mar 2024 09:06:21 -0700 Subject: [PATCH 6/6] 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 9fe868a8..9a3500ec 100644 --- a/client_test.go +++ b/client_test.go @@ -1542,23 +1542,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) {