From 07bae0b95d02bb210a00207b4a62b819713b2da6 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Fri, 23 Feb 2024 17:33:26 -0800 Subject: [PATCH 1/8] Fix ergonomics around JobList --- client_test.go | 26 ++++++------- internal/dbadapter/db_adapter.go | 3 -- internal/dbadapter/db_adapter_test.go | 4 -- internal/dblist/job_list.go | 6 +-- internal/dblist/job_list_test.go | 3 -- job_list_params.go | 56 ++++++++++++++++++++------- 6 files changed, 57 insertions(+), 41 deletions(-) diff --git a/client_test.go b/client_test.go index 4bf2258a..86db164b 100644 --- a/client_test.go +++ b/client_test.go @@ -1407,12 +1407,12 @@ func Test_Client_JobList(t *testing.T) { job2 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: dbsqlc.JobStateAvailable}) job3 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: dbsqlc.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 })) }) @@ -1433,11 +1433,11 @@ func Test_Client_JobList(t *testing.T) { job1 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: dbState, ScheduledAt: ptrutil.Ptr(now)}) job2 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: 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 })) } @@ -1459,11 +1459,11 @@ func Test_Client_JobList(t *testing.T) { job1 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: dbState, FinalizedAt: ptrutil.Ptr(now.Add(-10 * time.Second))}) job2 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: 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 })) } @@ -1478,11 +1478,11 @@ func Test_Client_JobList(t *testing.T) { job1 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: dbsqlc.JobStateRunning, AttemptedAt: ptrutil.Ptr(now)}) job2 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{State: dbsqlc.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 })) @@ -1524,11 +1524,11 @@ func Test_Client_JobList(t *testing.T) { require.NoError(t, err) require.Equal(t, []int64{job2.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State(rivertype.JobStateRunning).After(JobListCursorFromJob(jobRow3))) + jobs, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning).After(JobListCursorFromJob(jobRow3))) require.NoError(t, err) require.Equal(t, []int64{job4.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - jobs, err = client.JobList(ctx, NewJobListParams().State(rivertype.JobStateCompleted).After(JobListCursorFromJob(jobRow5))) + jobs, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateCompleted).After(JobListCursorFromJob(jobRow5))) require.NoError(t, err) require.Equal(t, []int64{job6.ID}, sliceutil.Map(jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) @@ -1542,11 +1542,11 @@ func Test_Client_JobList(t *testing.T) { job2 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{Metadata: []byte(`{"baz": "value"}`)}) job3 := insertJob(ctx, client.driver.GetDBPool(), insertJobParams{Metadata: []byte(`{"baz": "value"}`)}) - jobs, err := client.JobList(ctx, NewJobListParams().State("").Metadata(`{"foo": "bar"}`)) + jobs, err := client.JobList(ctx, NewJobListParams().States("").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().States("").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 })) @@ -1560,7 +1560,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/dbadapter/db_adapter.go b/internal/dbadapter/db_adapter.go index ce3ad01a..06948624 100644 --- a/internal/dbadapter/db_adapter.go +++ b/internal/dbadapter/db_adapter.go @@ -21,7 +21,6 @@ import ( "github.com/riverqueue/river/internal/util/sliceutil" "github.com/riverqueue/river/internal/util/valutil" "github.com/riverqueue/river/riverdriver" - "github.com/riverqueue/river/rivertype" ) // When a job has specified unique options, but has not set the ByState @@ -95,7 +94,6 @@ type JobListParams struct { OrderBy []JobListOrderBy Priorities []int16 Queues []string - State rivertype.JobState } // Adapter is an interface to the various database-level operations which River @@ -510,7 +508,6 @@ func (a *StandardAdapter) JobListTx(ctx context.Context, tx pgx.Tx, params JobLi NamedArgs: namedArgs, OrderBy: orderBy, Priorities: params.Priorities, - State: dbsqlc.JobState(params.State), }) if err != nil { return nil, err diff --git a/internal/dbadapter/db_adapter_test.go b/internal/dbadapter/db_adapter_test.go index fc649b3b..5df1ce5c 100644 --- a/internal/dbadapter/db_adapter_test.go +++ b/internal/dbadapter/db_adapter_test.go @@ -21,7 +21,6 @@ import ( "github.com/riverqueue/river/internal/util/ptrutil" "github.com/riverqueue/river/internal/util/sliceutil" "github.com/riverqueue/river/riverdriver" - "github.com/riverqueue/river/rivertype" ) func Test_StandardAdapter_JobCancel(t *testing.T) { @@ -881,7 +880,6 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { params := JobListParams{ LimitCount: 2, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, - State: rivertype.JobStateAvailable, } execTest(ctx, t, adapter, params, bundle.tx, func(jobs []*dbsqlc.RiverJob, err error) { @@ -907,7 +905,6 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { LimitCount: 2, NamedArgs: map[string]any{"paths1": []string{"job_num"}, "value1": 2}, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, - State: rivertype.JobStateAvailable, } execTest(ctx, t, adapter, params, bundle.tx, func(jobs []*dbsqlc.RiverJob, err error) { @@ -929,7 +926,6 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { LimitCount: 2, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, Queues: []string{"priority"}, - State: rivertype.JobStateAvailable, } execTest(ctx, t, adapter, params, bundle.tx, func(jobs []*dbsqlc.RiverJob, err error) { diff --git a/internal/dblist/job_list.go b/internal/dblist/job_list.go index 1ba278a7..43154235 100644 --- a/internal/dblist/job_list.go +++ b/internal/dblist/job_list.go @@ -36,7 +36,6 @@ type JobListOrderBy struct { } type JobListParams struct { - State dbsqlc.JobState Priorities []int16 Conditions string OrderBy []JobListOrderBy @@ -73,10 +72,7 @@ func JobList(ctx context.Context, tx pgx.Tx, arg JobListParams) ([]*dbsqlc.River } var conditions []string - if arg.State != "" { - conditions = append(conditions, "state = @state::river_job_state") - namedArgs["state"] = arg.State - } + if arg.Conditions != "" { conditions = append(conditions, arg.Conditions) } diff --git a/internal/dblist/job_list_test.go b/internal/dblist/job_list_test.go index 0b4714c3..cc05c9f4 100644 --- a/internal/dblist/job_list_test.go +++ b/internal/dblist/job_list_test.go @@ -7,7 +7,6 @@ import ( "github.com/jackc/pgx/v5" "github.com/stretchr/testify/require" - "github.com/riverqueue/river/internal/dbsqlc" "github.com/riverqueue/river/internal/riverinternaltest" ) @@ -21,7 +20,6 @@ func TestJobList(t *testing.T) { tx := riverinternaltest.TestTx(ctx, t) _, err := JobList(ctx, tx, JobListParams{ - State: dbsqlc.JobStateCompleted, LimitCount: 1, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderAsc}}, }) @@ -37,7 +35,6 @@ func TestJobList(t *testing.T) { _, err := JobList(ctx, tx, JobListParams{ Conditions: "queue = 'test' AND priority = 1 AND args->>'foo' = @foo", NamedArgs: pgx.NamedArgs{"foo": "bar"}, - State: dbsqlc.JobStateCompleted, LimitCount: 1, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderAsc}}, }) diff --git a/job_list_params.go b/job_list_params.go index b5b45a41..b8d649d3 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -8,6 +8,7 @@ import ( "strings" "time" + "github.com/lib/pq" "github.com/riverqueue/river/internal/dbadapter" "github.com/riverqueue/river/rivertype" ) @@ -90,12 +91,19 @@ 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 + // JobListOrderByTime specifies that the sort should be inferred from the job state. + JobListOrderByTime JobListOrderByField = "time" + // JobListOrderByCreatedAt specifies that the sort should be by the time the job was created. + JobListOrderByCreatedAt JobListOrderByField = "created_at" + // JobListOrderByScheduledAt specifies that the sort should be by the time the job was scheduled. + JobListOrderByScheduledAt JobListOrderByField = "scheduled_at" + // JobListOrderByAttemptedAt specifies that the sort should be by the last time the job was attempted. + JobListOrderByAttemptedAt JobListOrderByField = "attempted_at" + // JobListOrderByFinalizedAt specifies that the sort should be by the time the job was finalized. + JobListOrderByFinalizedAt JobListOrderByField = "finalized_at" ) // JobListParams specifies the parameters for a JobList query. It must be @@ -110,7 +118,8 @@ type JobListParams struct { queues []string sortField JobListOrderByField sortOrder SortOrder - state rivertype.JobState + kinds []string + states []rivertype.JobState } // NewJobListParams creates a new JobListParams to return available jobs sorted @@ -120,7 +129,7 @@ func NewJobListParams() *JobListParams { paginationCount: 100, sortField: JobListOrderByTime, sortOrder: SortOrderAsc, - state: rivertype.JobStateAvailable, + states: []rivertype.JobState{rivertype.JobStateAvailable}, } } @@ -132,7 +141,8 @@ func (p *JobListParams) copy() *JobListParams { queues: append([]string(nil), p.queues...), sortField: p.sortField, sortOrder: p.sortOrder, - state: p.state, + states: p.states, + kinds: p.kinds, } } @@ -152,10 +162,12 @@ func (p *JobListParams) toDBParams() (*dbadapter.JobListParams, error) { return nil, errors.New("invalid sort order") } - if p.sortField != JobListOrderByTime { - return nil, errors.New("invalid sort field") + timeField := "created_at" + if p.sortField == JobListOrderByTime { + timeField = jobListTimeFieldForState(p.states[0]) + } else { + timeField = string(p.sortField) } - timeField := jobListTimeFieldForState(p.state) orderBy = append(orderBy, []dbadapter.JobListOrderBy{ {Expr: timeField, Order: sortOrder}, {Expr: "id", Order: sortOrder}, @@ -166,6 +178,15 @@ func (p *JobListParams) toDBParams() (*dbadapter.JobListParams, error) { namedArgs["metadata_fragment"] = p.metadataFragment } + if len(p.kinds) > 0 { + conditions = append(conditions, `"kind" = ANY(@kinds)`) + namedArgs["kinds"] = pq.Array(p.kinds) + } + if len(p.states) > 0 { + conditions = append(conditions, "state = ANY(@states)") + namedArgs["states"] = pq.Array(p.states) + } + if p.after != nil { if sortOrder == dbadapter.SortOrderAsc { conditions = append(conditions, fmt.Sprintf(`("%s" > @cursor_time OR ("%s" = @cursor_time AND "id" > @after_id))`, timeField, timeField)) @@ -190,7 +211,6 @@ func (p *JobListParams) toDBParams() (*dbadapter.JobListParams, error) { OrderBy: orderBy, Priorities: nil, Queues: p.queues, - State: p.state, } return dbParams, nil @@ -246,9 +266,19 @@ func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) // State returns an updated filter set that will only return jobs in the given // state. -func (p *JobListParams) State(state rivertype.JobState) *JobListParams { +func (p *JobListParams) States(states ...rivertype.JobState) *JobListParams { + result := p.copy() + result.states = make([]rivertype.JobState, 0, len(p.states)) + copy(result.states, p.states) + return result +} + +// Kinds returns an updated filter set that will only return jobs of the given +// kinds. +func (p *JobListParams) Kinds(kinds ...string) *JobListParams { result := p.copy() - result.state = state + result.kinds = make([]string, 0, len(kinds)) + copy(result.kinds, kinds) return result } From d78d9a0932312fde3f10a6c39228e0183cadc49a Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 17:07:14 -0800 Subject: [PATCH 2/8] Fix test now state filter isn't applied in the db adapter by default --- internal/dbadapter/db_adapter_test.go | 5 ++++- job_list_params.go | 2 +- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/internal/dbadapter/db_adapter_test.go b/internal/dbadapter/db_adapter_test.go index 5df1ce5c..42905e7f 100644 --- a/internal/dbadapter/db_adapter_test.go +++ b/internal/dbadapter/db_adapter_test.go @@ -12,6 +12,7 @@ import ( "time" "github.com/jackc/pgx/v5" + "github.com/lib/pq" "github.com/stretchr/testify/require" "github.com/riverqueue/river/internal/dbsqlc" @@ -21,6 +22,7 @@ import ( "github.com/riverqueue/river/internal/util/ptrutil" "github.com/riverqueue/river/internal/util/sliceutil" "github.com/riverqueue/river/riverdriver" + "github.com/riverqueue/river/rivertype" ) func Test_StandardAdapter_JobCancel(t *testing.T) { @@ -880,11 +882,12 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { params := JobListParams{ LimitCount: 2, OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, + Conditions: "state = any(@states)", + NamedArgs: map[string]any{"states": pq.Array([]rivertype.JobState{rivertype.JobStateAvailable})}, } execTest(ctx, t, adapter, params, bundle.tx, func(jobs []*dbsqlc.RiverJob, err error) { require.NoError(t, err) - // job 1 is excluded due to pagination limit of 2, while job 4 is excluded // due to its state: job2 := bundle.jobs[1] diff --git a/job_list_params.go b/job_list_params.go index b8d649d3..6599c1c0 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -179,7 +179,7 @@ func (p *JobListParams) toDBParams() (*dbadapter.JobListParams, error) { } if len(p.kinds) > 0 { - conditions = append(conditions, `"kind" = ANY(@kinds)`) + conditions = append(conditions, `"kind = ANY(@kinds)`) namedArgs["kinds"] = pq.Array(p.kinds) } if len(p.states) > 0 { From 9dfcc1cd6b2db45c93fb709dbfdc6530357a3c7f Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 17:35:52 -0800 Subject: [PATCH 3/8] Add tests --- internal/dbadapter/db_adapter_test.go | 48 +++++++++++++++++++++++++-- internal/dblist/job_list.go | 1 - 2 files changed, 45 insertions(+), 4 deletions(-) diff --git a/internal/dbadapter/db_adapter_test.go b/internal/dbadapter/db_adapter_test.go index 42905e7f..dc463776 100644 --- a/internal/dbadapter/db_adapter_test.go +++ b/internal/dbadapter/db_adapter_test.go @@ -831,7 +831,7 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { adapter := NewStandardAdapter(riverinternaltest.BaseServiceArchetype(t), testAdapterConfig(bundle.ex)) adapter.TimeNowUTC = func() time.Time { return bundle.baselineTime } - params := makeFakeJobInsertParams(1, &makeFakeJobInsertParamsOpts{Queue: ptrutil.Ptr("priority")}) + params := makeFakeJobInsertParams(1, &makeFakeJobInsertParamsOpts{Queue: ptrutil.Ptr("priority"), Kind: ptrutil.Ptr("different_kind"), ScheduledAt: ptrutil.Ptr(time.Now().Add(2 * time.Second))}) job1, err := adapter.JobInsert(ctx, params) require.NoError(t, err) @@ -843,7 +843,7 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { job3, err := adapter.JobInsert(ctx, params) require.NoError(t, err) - params = makeFakeJobInsertParams(4, &makeFakeJobInsertParamsOpts{State: ptrutil.Ptr(dbsqlc.JobStateRunning)}) + params = makeFakeJobInsertParams(4, &makeFakeJobInsertParamsOpts{State: ptrutil.Ptr(dbsqlc.JobStateRunning), ScheduledAt: ptrutil.Ptr(time.Now().Add(1 * time.Second))}) job4, err := adapter.JobInsert(ctx, params) require.NoError(t, err) @@ -960,6 +960,47 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) { require.Equal(t, []int64{job3.ID}, returnedIDs) }) }) + + t.Run("WithKindFilter", func(t *testing.T) { + t.Parallel() + + adapter, bundle := setupTx(t) + + params := JobListParams{ + Conditions: "kind = @kind", + LimitCount: 2, + NamedArgs: map[string]any{"kind": "different_kind"}, + OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}}, + } + + execTest(ctx, t, adapter, params, bundle.tx, func(jobs []*dbsqlc.RiverJob, err error) { + require.NoError(t, err) + + job1 := bundle.jobs[0] + returnedIDs := sliceutil.Map(jobs, func(j *dbsqlc.RiverJob) int64 { return j.ID }) + require.Equal(t, []int64{job1.ID}, returnedIDs) + }) + }) + + t.Run("WithSpecificTimeTypeSort", func(t *testing.T) { + t.Parallel() + + adapter, bundle := setupTx(t) + + params := JobListParams{ + LimitCount: 2, + OrderBy: []JobListOrderBy{{Expr: "scheduled_at", Order: SortOrderDesc}}, + } + + execTest(ctx, t, adapter, params, bundle.tx, func(jobs []*dbsqlc.RiverJob, err error) { + require.NoError(t, err) + // Job 1 is scheduled 2 seconds later, job 4 is scheduled 1 second later + job1 := bundle.jobs[0] + job4 := bundle.jobs[3] + returnedIDs := sliceutil.Map(jobs, func(j *dbsqlc.RiverJob) int64 { return j.ID }) + require.Equal(t, []int64{job1.ID, job4.ID}, returnedIDs) + }) + }) } func Test_StandardAdapter_JobRetryImmediately(t *testing.T) { @@ -1523,6 +1564,7 @@ func testAdapterConfig(ex dbutil.Executor) *StandardAdapterConfig { } type makeFakeJobInsertParamsOpts struct { + Kind *string Metadata []byte Queue *string ScheduledAt *time.Time @@ -1541,7 +1583,7 @@ func makeFakeJobInsertParams(i int, opts *makeFakeJobInsertParamsOpts) *JobInser return &JobInsertParams{ EncodedArgs: []byte(fmt.Sprintf(`{"job_num":%d}`, i)), - Kind: "fake_job", + Kind: ptrutil.ValOrDefault(opts.Kind, "fake_job"), MaxAttempts: rivercommon.MaxAttemptsDefault, Metadata: metadata, Priority: rivercommon.PriorityDefault, diff --git a/internal/dblist/job_list.go b/internal/dblist/job_list.go index 43154235..049bcded 100644 --- a/internal/dblist/job_list.go +++ b/internal/dblist/job_list.go @@ -72,7 +72,6 @@ func JobList(ctx context.Context, tx pgx.Tx, arg JobListParams) ([]*dbsqlc.River } var conditions []string - if arg.Conditions != "" { conditions = append(conditions, arg.Conditions) } From e27647f7f7ab34672cefd7313ff10c3784f49ace Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 18:08:47 -0800 Subject: [PATCH 4/8] Fix adding states to JobListParams --- job_list_params.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/job_list_params.go b/job_list_params.go index 6599c1c0..73d091f8 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -268,8 +268,8 @@ func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) // state. func (p *JobListParams) States(states ...rivertype.JobState) *JobListParams { result := p.copy() - result.states = make([]rivertype.JobState, 0, len(p.states)) - copy(result.states, p.states) + result.states = make([]rivertype.JobState, 0, len(states)) + copy(result.states, states) return result } From 2dcabf41402497df8613acfa9d0a0d40e2a8470b Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 18:21:28 -0800 Subject: [PATCH 5/8] Copy correctly --- job_list_params.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/job_list_params.go b/job_list_params.go index 73d091f8..67856539 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -141,8 +141,8 @@ func (p *JobListParams) copy() *JobListParams { queues: append([]string(nil), p.queues...), sortField: p.sortField, sortOrder: p.sortOrder, - states: p.states, - kinds: p.kinds, + states: append([]rivertype.JobState(nil), p.states...), + kinds: append([]string(nil), p.kinds...), } } From 35902624adb06c344e90c85360ce532c55e38ee9 Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 19:00:53 -0800 Subject: [PATCH 6/8] Fix copy behavior in JobParam queues, states and kinds --- job_list_params.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/job_list_params.go b/job_list_params.go index 67856539..a52581c3 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -250,7 +250,7 @@ func (p *JobListParams) Metadata(json string) *JobListParams { // given queues. func (p *JobListParams) Queues(queues ...string) *JobListParams { result := p.copy() - result.queues = make([]string, 0, len(queues)) + result.queues = make([]string, len(queues)) copy(result.queues, queues) return result } @@ -268,7 +268,7 @@ func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) // state. func (p *JobListParams) States(states ...rivertype.JobState) *JobListParams { result := p.copy() - result.states = make([]rivertype.JobState, 0, len(states)) + result.states = make([]rivertype.JobState, len(states)) copy(result.states, states) return result } @@ -277,7 +277,7 @@ func (p *JobListParams) States(states ...rivertype.JobState) *JobListParams { // kinds. func (p *JobListParams) Kinds(kinds ...string) *JobListParams { result := p.copy() - result.kinds = make([]string, 0, len(kinds)) + result.kinds = make([]string, len(kinds)) copy(result.kinds, kinds) return result } From c586c618996a2baf3fc2f3ea32a2a377d18a47df Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 19:21:41 -0800 Subject: [PATCH 7/8] Wayward double quote --- job_list_params.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/job_list_params.go b/job_list_params.go index a52581c3..b812fc1c 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -179,7 +179,7 @@ func (p *JobListParams) toDBParams() (*dbadapter.JobListParams, error) { } if len(p.kinds) > 0 { - conditions = append(conditions, `"kind = ANY(@kinds)`) + conditions = append(conditions, `kind = ANY(@kinds)`) namedArgs["kinds"] = pq.Array(p.kinds) } if len(p.states) > 0 { From 1e23a7fd386f02b5aef7ca9bd2e2332f4e13343b Mon Sep 17 00:00:00 2001 From: Jos Kraaijeveld Date: Sat, 24 Feb 2024 19:23:04 -0800 Subject: [PATCH 8/8] Consistent quote usage --- job_list_params.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/job_list_params.go b/job_list_params.go index b812fc1c..76fcaa81 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -183,7 +183,7 @@ func (p *JobListParams) toDBParams() (*dbadapter.JobListParams, error) { namedArgs["kinds"] = pq.Array(p.kinds) } if len(p.states) > 0 { - conditions = append(conditions, "state = ANY(@states)") + conditions = append(conditions, `state = ANY(@states)`) namedArgs["states"] = pq.Array(p.states) }