Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 13 additions & 13 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 }))
})
Expand All @@ -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 }))
}
Expand All @@ -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 }))
}
Expand All @@ -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 }))
Expand Down Expand Up @@ -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 }))
})
Expand All @@ -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 }))
Expand All @@ -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)
})
Expand Down
3 changes: 0 additions & 3 deletions internal/dbadapter/db_adapter.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
55 changes: 48 additions & 7 deletions internal/dbadapter/db_adapter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -830,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)

Expand All @@ -842,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)

Expand Down Expand Up @@ -881,12 +882,12 @@ func Test_StandardAdapter_JobList_and_JobListTx(t *testing.T) {
params := JobListParams{
LimitCount: 2,
OrderBy: []JobListOrderBy{{Expr: "id", Order: SortOrderDesc}},
State: rivertype.JobStateAvailable,
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]
Expand All @@ -907,7 +908,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) {
Expand All @@ -929,7 +929,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) {
Expand Down Expand Up @@ -961,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) {
Expand Down Expand Up @@ -1524,6 +1564,7 @@ func testAdapterConfig(ex dbutil.Executor) *StandardAdapterConfig {
}

type makeFakeJobInsertParamsOpts struct {
Kind *string
Metadata []byte
Queue *string
ScheduledAt *time.Time
Expand All @@ -1542,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,
Expand Down
5 changes: 0 additions & 5 deletions internal/dblist/job_list.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@ type JobListOrderBy struct {
}

type JobListParams struct {
State dbsqlc.JobState
Priorities []int16
Conditions string
OrderBy []JobListOrderBy
Expand Down Expand Up @@ -73,10 +72,6 @@ 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)
}
Expand Down
3 changes: 0 additions & 3 deletions internal/dblist/job_list_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand All @@ -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}},
})
Expand All @@ -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}},
})
Expand Down
Loading