diff --git a/CHANGELOG.md b/CHANGELOG.md index 62501c68..cca06887 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Changed - **Breaking change:** JobList/JobListTx now support querying Jobs by a list of Job Kinds and States. Also allows for filtering by specific timestamp values. Thank you Jos Kraaijeveld (@thatjos)! 🙏🏻 [PR #236](https://github.com/riverqueue/river/pull/236). +- **Breaking change:** `river.JobState*` type aliases have been removed. All job state constants should be accessed through `rivertype.JobState*` instead. [PR #300](https://github.com/riverqueue/river/pull/300). ## [0.3.0] - 2024-04-15 diff --git a/client.go b/client.go index 63dd9479..0e5dac68 100644 --- a/client.go +++ b/client.go @@ -904,13 +904,13 @@ func (c *Client[TTx]) distributeJob(job *rivertype.JobRow, stats *JobStatistics) var event *Event switch job.State { - case JobStateCancelled: + case rivertype.JobStateCancelled: event = &Event{Kind: EventKindJobCancelled, Job: job, JobStats: stats} - case JobStateCompleted: + case rivertype.JobStateCompleted: event = &Event{Kind: EventKindJobCompleted, Job: job, JobStats: stats} - case JobStateScheduled: + case rivertype.JobStateScheduled: event = &Event{Kind: EventKindJobSnoozed, Job: job, JobStats: stats} - case JobStateAvailable, JobStateDiscarded, JobStateRetryable, JobStateRunning: + case rivertype.JobStateAvailable, rivertype.JobStateDiscarded, rivertype.JobStateRetryable, rivertype.JobStateRunning: event = &Event{Kind: EventKindJobFailed, Job: job, JobStats: stats} default: // linter exhaustive rule prevents this from being reached diff --git a/client_test.go b/client_test.go index 59d7bb5a..5d6324ac 100644 --- a/client_test.go +++ b/client_test.go @@ -268,7 +268,7 @@ func Test_Client(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, subscribeChan) require.Equal(t, EventKindJobCancelled, event.Kind) - require.Equal(t, JobStateCancelled, event.Job.State) + require.Equal(t, rivertype.JobStateCancelled, event.Job.State) require.WithinDuration(t, time.Now(), *event.Job.FinalizedAt, 2*time.Second) updatedJob, err := client.JobGet(ctx, insertedJob.ID) @@ -298,7 +298,7 @@ func Test_Client(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, subscribeChan) require.Equal(t, EventKindJobSnoozed, event.Kind) - require.Equal(t, JobStateScheduled, event.Job.State) + require.Equal(t, rivertype.JobStateScheduled, event.Job.State) require.WithinDuration(t, time.Now().Add(15*time.Minute), event.Job.ScheduledAt, 2*time.Second) updatedJob, err := client.JobGet(ctx, insertedJob.ID) @@ -347,7 +347,7 @@ func Test_Client(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, subscribeChan) require.Equal(t, EventKindJobCancelled, event.Kind) - require.Equal(t, JobStateCancelled, event.Job.State) + require.Equal(t, rivertype.JobStateCancelled, event.Job.State) require.WithinDuration(t, time.Now(), *event.Job.FinalizedAt, 2*time.Second) jobAfterCancel, err := client.JobGet(ctx, insertedJob.ID) @@ -485,7 +485,7 @@ func Test_Client(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, subscribeChan) require.Equal(t, EventKindJobCompleted, event.Kind) require.Equal(t, job.ID, event.Job.ID) - require.Equal(t, JobStateCompleted, event.Job.State) + require.Equal(t, rivertype.JobStateCompleted, event.Job.State) }) t.Run("StartStopStress", func(t *testing.T) { @@ -702,7 +702,7 @@ func Test_Client_Stop(t *testing.T) { require.NoError(t, client.Stop(ctx)) - res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning)) require.NoError(t, err) require.Empty(t, res.Jobs, "expected no jobs to be left running") }) @@ -1487,12 +1487,12 @@ func Test_Client_JobList(t *testing.T) { job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)}) job3 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateRunning)}) - res, err := client.JobList(ctx, NewJobListParams().States(JobStateAvailable)) + res, err := client.JobList(ctx, NewJobListParams().States(rivertype.JobStateAvailable)) require.NoError(t, err) // jobs ordered by ScheduledAt ASC by default require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - res, err = client.JobList(ctx, NewJobListParams().States(JobStateRunning)) + res, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning)) require.NoError(t, err) require.Equal(t, []int64{job3.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) }) @@ -1504,14 +1504,14 @@ func Test_Client_JobList(t *testing.T) { now := time.Now().UTC() - states := map[rivertype.JobState]rivertype.JobState{ - JobStateAvailable: rivertype.JobStateAvailable, - JobStateRetryable: rivertype.JobStateRetryable, - JobStateScheduled: rivertype.JobStateScheduled, + states := []rivertype.JobState{ + rivertype.JobStateAvailable, + rivertype.JobStateRetryable, + rivertype.JobStateScheduled, } - for state, dbState := range states { - 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))}) + for _, state := range states { + job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(state), ScheduledAt: &now}) + job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(state), ScheduledAt: ptrutil.Ptr(now.Add(-5 * time.Second))}) res, err := client.JobList(ctx, NewJobListParams().States(state)) require.NoError(t, err) @@ -1530,14 +1530,14 @@ func Test_Client_JobList(t *testing.T) { now := time.Now().UTC() - states := map[rivertype.JobState]rivertype.JobState{ - JobStateCancelled: rivertype.JobStateCancelled, - JobStateCompleted: rivertype.JobStateCompleted, - JobStateDiscarded: rivertype.JobStateDiscarded, + states := []rivertype.JobState{ + rivertype.JobStateCancelled, + rivertype.JobStateCompleted, + rivertype.JobStateDiscarded, } - for state, dbState := range states { - 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))}) + for _, state := range states { + job1 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(state), FinalizedAt: ptrutil.Ptr(now.Add(-10 * time.Second))}) + job2 := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{State: ptrutil.Ptr(state), FinalizedAt: ptrutil.Ptr(now.Add(-15 * time.Second))}) res, err := client.JobList(ctx, NewJobListParams().States(state)) require.NoError(t, err) @@ -1558,11 +1558,11 @@ 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))}) - res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning)) require.NoError(t, err) require.Equal(t, []int64{job2.ID, job1.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) - res, err = client.JobList(ctx, NewJobListParams().States(JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc)) + res, err = client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning).OrderBy(JobListOrderByTime, SortOrderDesc)) require.NoError(t, err) // Sort order was explicitly reversed: require.Equal(t, []int64{job1.ID, job2.ID}, sliceutil.Map(res.Jobs, func(job *rivertype.JobRow) int64 { return job.ID })) @@ -1646,7 +1646,7 @@ func Test_Client_JobList(t *testing.T) { ctx, cancel := context.WithCancel(ctx) cancel() // cancel immediately - res, err := client.JobList(ctx, NewJobListParams().States(JobStateRunning)) + res, err := client.JobList(ctx, NewJobListParams().States(rivertype.JobStateRunning)) require.ErrorIs(t, context.Canceled, err) require.Nil(t, res) }) @@ -2350,28 +2350,28 @@ func Test_Client_Subscribe(t *testing.T) { eventCompleted1 := eventsByName["completed1"] require.Equal(t, EventKindJobCompleted, eventCompleted1.Kind) require.Equal(t, jobCompleted1.ID, eventCompleted1.Job.ID) - require.Equal(t, JobStateCompleted, eventCompleted1.Job.State) + require.Equal(t, rivertype.JobStateCompleted, eventCompleted1.Job.State) } { eventCompleted2 := eventsByName["completed2"] require.Equal(t, EventKindJobCompleted, eventCompleted2.Kind) require.Equal(t, jobCompleted2.ID, eventCompleted2.Job.ID) - require.Equal(t, JobStateCompleted, eventCompleted2.Job.State) + require.Equal(t, rivertype.JobStateCompleted, eventCompleted2.Job.State) } { eventFailed1 := eventsByName["failed1"] require.Equal(t, EventKindJobFailed, eventFailed1.Kind) require.Equal(t, jobFailed1.ID, eventFailed1.Job.ID) - require.Equal(t, JobStateRetryable, eventFailed1.Job.State) + require.Equal(t, rivertype.JobStateRetryable, eventFailed1.Job.State) } { eventFailed2 := eventsByName["failed2"] require.Equal(t, EventKindJobFailed, eventFailed2.Kind) require.Equal(t, jobFailed2.ID, eventFailed2.Job.ID) - require.Equal(t, JobStateRetryable, eventFailed2.Job.State) + require.Equal(t, rivertype.JobStateRetryable, eventFailed2.Job.State) } }) @@ -2412,7 +2412,7 @@ func Test_Client_Subscribe(t *testing.T) { eventCompleted := eventsByName["completed1"] require.Equal(t, EventKindJobCompleted, eventCompleted.Kind) require.Equal(t, jobCompleted.ID, eventCompleted.Job.ID) - require.Equal(t, JobStateCompleted, eventCompleted.Job.State) + require.Equal(t, rivertype.JobStateCompleted, eventCompleted.Job.State) _, ok := eventsByName["failed1"] // filtered out require.False(t, ok) @@ -2458,7 +2458,7 @@ func Test_Client_Subscribe(t *testing.T) { eventFailed := eventsByName["failed1"] require.Equal(t, EventKindJobFailed, eventFailed.Kind) require.Equal(t, jobFailed.ID, eventFailed.Job.ID) - require.Equal(t, JobStateRetryable, eventFailed.Job.State) + require.Equal(t, rivertype.JobStateRetryable, eventFailed.Job.State) }) t.Run("PanicOnUnknownKind", func(t *testing.T) { @@ -2567,28 +2567,28 @@ func Test_Client_SubscribeConfig(t *testing.T) { eventCompleted1 := eventsByName["completed1"] require.Equal(t, EventKindJobCompleted, eventCompleted1.Kind) require.Equal(t, jobCompleted1.ID, eventCompleted1.Job.ID) - require.Equal(t, JobStateCompleted, eventCompleted1.Job.State) + require.Equal(t, rivertype.JobStateCompleted, eventCompleted1.Job.State) } { eventCompleted2 := eventsByName["completed2"] require.Equal(t, EventKindJobCompleted, eventCompleted2.Kind) require.Equal(t, jobCompleted2.ID, eventCompleted2.Job.ID) - require.Equal(t, JobStateCompleted, eventCompleted2.Job.State) + require.Equal(t, rivertype.JobStateCompleted, eventCompleted2.Job.State) } { eventFailed1 := eventsByName["failed1"] require.Equal(t, EventKindJobFailed, eventFailed1.Kind) require.Equal(t, jobFailed1.ID, eventFailed1.Job.ID) - require.Equal(t, JobStateRetryable, eventFailed1.Job.State) + require.Equal(t, rivertype.JobStateRetryable, eventFailed1.Job.State) } { eventFailed2 := eventsByName["failed2"] require.Equal(t, EventKindJobFailed, eventFailed2.Kind) require.Equal(t, jobFailed2.ID, eventFailed2.Job.ID) - require.Equal(t, JobStateRetryable, eventFailed2.Job.State) + require.Equal(t, rivertype.JobStateRetryable, eventFailed2.Job.State) } }) @@ -2799,7 +2799,7 @@ func Test_Client_JobCompletion(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, bundle.SubscribeChan) require.Equal(job.ID, event.Job.ID) - require.Equal(JobStateCompleted, event.Job.State) + require.Equal(rivertype.JobStateCompleted, event.Job.State) reloadedJob, err := client.JobGet(ctx, job.ID) require.NoError(err) @@ -2834,7 +2834,7 @@ func Test_Client_JobCompletion(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, bundle.SubscribeChan) require.Equal(job.ID, event.Job.ID) - require.Equal(JobStateCompleted, event.Job.State) + require.Equal(rivertype.JobStateCompleted, event.Job.State) reloadedJob, err := client.JobGet(ctx, job.ID) require.NoError(err) @@ -2858,7 +2858,7 @@ func Test_Client_JobCompletion(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, bundle.SubscribeChan) require.Equal(job.ID, event.Job.ID) - require.Equal(JobStateRetryable, event.Job.State) + require.Equal(rivertype.JobStateRetryable, event.Job.State) reloadedJob, err := client.JobGet(ctx, job.ID) require.NoError(err) @@ -2883,7 +2883,7 @@ func Test_Client_JobCompletion(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, bundle.SubscribeChan) require.Equal(job.ID, event.Job.ID) - require.Equal(JobStateCancelled, event.Job.State) + require.Equal(rivertype.JobStateCancelled, event.Job.State) reloadedJob, err := client.JobGet(ctx, job.ID) require.NoError(err) @@ -2924,7 +2924,7 @@ func Test_Client_JobCompletion(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, bundle.SubscribeChan) require.Equal(job.ID, event.Job.ID) - require.Equal(JobStateDiscarded, event.Job.State) + require.Equal(rivertype.JobStateDiscarded, event.Job.State) reloadedJob, err := client.JobGet(ctx, job.ID) require.NoError(err) @@ -2961,8 +2961,8 @@ func Test_Client_JobCompletion(t *testing.T) { event := riverinternaltest.WaitOrTimeout(t, bundle.SubscribeChan) require.Equal(job.ID, event.Job.ID) - require.Equal(JobStateCompleted, event.Job.State) - require.Equal(JobStateCompleted, updatedJob.State) + require.Equal(rivertype.JobStateCompleted, event.Job.State) + require.Equal(rivertype.JobStateCompleted, updatedJob.State) require.NotNil(updatedJob) require.NotNil(event.Job.FinalizedAt) require.NotNil(updatedJob.FinalizedAt) @@ -3014,7 +3014,7 @@ func Test_Client_UnknownJobKindErrorsTheJob(t *testing.T) { require.Equal("RandomWorkerNameThatIsNeverRegistered", insertedJob.Kind) require.Len(event.Job.Errors, 1) require.Equal((&UnknownJobKindError{Kind: "RandomWorkerNameThatIsNeverRegistered"}).Error(), event.Job.Errors[0].Error) - require.Equal(JobStateRetryable, event.Job.State) + require.Equal(rivertype.JobStateRetryable, event.Job.State) // Ensure that ScheduledAt was updated with next run time: require.True(event.Job.ScheduledAt.After(insertedJob.ScheduledAt)) // It's the 1st attempt that failed. Attempt won't be incremented again until @@ -3768,7 +3768,7 @@ func TestInsert(t *testing.T) { require.WithinDuration(opts.ScheduledAt.UTC(), insertedJob.ScheduledAt, time.Microsecond) require.Equal([]string{"tag1", "tag2"}, insertedJob.Tags) // derived state: - require.Equal(JobStateScheduled, insertedJob.State) + require.Equal(rivertype.JobStateScheduled, insertedJob.State) require.Equal("noOp", insertedJob.Kind) // default state: // require.Equal([]byte("{}"), insertedJob.metadata) @@ -3785,7 +3785,7 @@ func TestInsert(t *testing.T) { // specified by inputs: requireEqualArgs(t, args, insertedJob.EncodedArgs) // derived state: - require.Equal(JobStateAvailable, insertedJob.State) + require.Equal(rivertype.JobStateAvailable, insertedJob.State) require.Equal("noOp", insertedJob.Kind) // default state: require.Equal(QueueDefault, insertedJob.Queue) @@ -3877,7 +3877,7 @@ func TestUniqueOpts(t *testing.T) { uniqueOpts := UniqueOpts{ ByPeriod: 24 * time.Hour, - ByState: []rivertype.JobState{JobStateAvailable, JobStateCompleted}, + ByState: []rivertype.JobState{rivertype.JobStateAvailable, rivertype.JobStateCompleted}, } job0, err := client.Insert(ctx, noOpArgs{}, &InsertOpts{ diff --git a/job.go b/job.go index 69cfdc7d..aad8b46f 100644 --- a/job.go +++ b/job.go @@ -4,17 +4,6 @@ import ( "github.com/riverqueue/river/rivertype" ) -const ( - JobStateAvailable = rivertype.JobStateAvailable - JobStateCancelled = rivertype.JobStateCancelled - JobStateCompleted = rivertype.JobStateCompleted - JobStateDiscarded = rivertype.JobStateDiscarded - JobStateRetryable = rivertype.JobStateRetryable - JobStateRunning = rivertype.JobStateRunning - JobStateScheduled = rivertype.JobStateScheduled -) - -// Job represents a single unit of work, holding both the arguments and // information for a job with args of type T. type Job[T JobArgs] struct { *rivertype.JobRow diff --git a/job_complete_tx.go b/job_complete_tx.go index bc5bfeca..d7eef9f6 100644 --- a/job_complete_tx.go +++ b/job_complete_tx.go @@ -7,6 +7,7 @@ import ( "time" "github.com/riverqueue/river/riverdriver" + "github.com/riverqueue/river/rivertype" ) // JobCompleteTx marks the job as completed as part of transaction tx. If tx is @@ -23,7 +24,7 @@ import ( // // Returns the updated, completed job. func JobCompleteTx[TDriver riverdriver.Driver[TTx], TTx any, TArgs JobArgs](ctx context.Context, tx TTx, job *Job[TArgs]) (*Job[TArgs], error) { - if job.State != JobStateRunning { + if job.State != rivertype.JobStateRunning { return nil, errors.New("job must be running") } diff --git a/job_complete_tx_test.go b/job_complete_tx_test.go index 4828e6b3..74a8f896 100644 --- a/job_complete_tx_test.go +++ b/job_complete_tx_test.go @@ -13,6 +13,7 @@ import ( "github.com/riverqueue/river/internal/util/ptrutil" "github.com/riverqueue/river/riverdriver" "github.com/riverqueue/river/riverdriver/riverpgxv5" + "github.com/riverqueue/river/rivertype" ) func TestJobCompleteTx(t *testing.T) { @@ -46,17 +47,17 @@ func TestJobCompleteTx(t *testing.T) { bundle := setup(t) job := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{ - State: ptrutil.Ptr(JobStateRunning), + State: ptrutil.Ptr(rivertype.JobStateRunning), }) completedJob, err := JobCompleteTx[*riverpgxv5.Driver](ctx, bundle.tx, &Job[JobArgs]{JobRow: job}) require.NoError(t, err) - require.Equal(t, JobStateCompleted, completedJob.State) + require.Equal(t, rivertype.JobStateCompleted, completedJob.State) require.WithinDuration(t, time.Now(), *completedJob.FinalizedAt, 2*time.Second) updatedJob, err := bundle.exec.JobGetByID(ctx, job.ID) require.NoError(t, err) - require.Equal(t, JobStateCompleted, updatedJob.State) + require.Equal(t, rivertype.JobStateCompleted, updatedJob.State) }) t.Run("ErrorIfNotRunning", func(t *testing.T) { diff --git a/job_list_params.go b/job_list_params.go index 2881be74..77b08c0a 100644 --- a/job_list_params.go +++ b/job_list_params.go @@ -149,13 +149,13 @@ func NewJobListParams() *JobListParams { sortField: JobListOrderByTime, sortOrder: SortOrderAsc, states: []rivertype.JobState{ - JobStateAvailable, - JobStateCancelled, - JobStateCompleted, - JobStateDiscarded, - JobStateRetryable, - JobStateRunning, - JobStateScheduled, + rivertype.JobStateAvailable, + rivertype.JobStateCancelled, + rivertype.JobStateCompleted, + rivertype.JobStateDiscarded, + rivertype.JobStateRetryable, + rivertype.JobStateRunning, + rivertype.JobStateScheduled, }, } } @@ -195,7 +195,7 @@ func (p *JobListParams) toDBParams() (*dblist.JobListParams, error) { for _, state := range p.states { //nolint:exhaustive switch state { - case JobStateCancelled, JobStateCompleted, JobStateDiscarded: + case rivertype.JobStateCancelled, rivertype.JobStateCompleted, rivertype.JobStateDiscarded: default: currentNonFinalizedStates = append(currentNonFinalizedStates, state) } @@ -317,9 +317,9 @@ func (p *JobListParams) OrderBy(field JobListOrderByField, direction SortOrder) result.sortField = field if !p.overrodeState { result.states = []rivertype.JobState{ - JobStateCancelled, - JobStateCompleted, - JobStateDiscarded, + rivertype.JobStateCancelled, + rivertype.JobStateCompleted, + rivertype.JobStateDiscarded, } } default: diff --git a/job_test.go b/job_test.go index 1af4bf10..b7576289 100644 --- a/job_test.go +++ b/job_test.go @@ -16,5 +16,5 @@ func TestJobUniqueOpts_isEmpty(t *testing.T) { require.False(t, (&UniqueOpts{ByArgs: true}).isEmpty()) require.False(t, (&UniqueOpts{ByPeriod: 1 * time.Nanosecond}).isEmpty()) require.False(t, (&UniqueOpts{ByQueue: true}).isEmpty()) - require.False(t, (&UniqueOpts{ByState: []rivertype.JobState{JobStateAvailable}}).isEmpty()) + require.False(t, (&UniqueOpts{ByState: []rivertype.JobState{rivertype.JobStateAvailable}}).isEmpty()) } diff --git a/rivertest/rivertest_test.go b/rivertest/rivertest_test.go index c906ca6e..55077f48 100644 --- a/rivertest/rivertest_test.go +++ b/rivertest/rivertest_test.go @@ -14,6 +14,7 @@ import ( "github.com/riverqueue/river" "github.com/riverqueue/river/internal/riverinternaltest" "github.com/riverqueue/river/riverdriver/riverpgxv5" + "github.com/riverqueue/river/rivertype" ) // Gives us a nice, stable time to test against. @@ -251,7 +252,7 @@ func TestRequireInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }) require.True(t, mockT.Failed) @@ -268,7 +269,7 @@ func TestRequireInsertedTx(t *testing.T) { Priority: 3, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }) require.True(t, mockT.Failed) @@ -285,7 +286,7 @@ func TestRequireInsertedTx(t *testing.T) { Priority: 2, Queue: "wrong-queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }) require.True(t, mockT.Failed) @@ -302,7 +303,7 @@ func TestRequireInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime.Add(3*time.Minute + 23*time.Second + 123*time.Microsecond), - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }) require.True(t, mockT.Failed) @@ -319,7 +320,7 @@ func TestRequireInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateCancelled, + State: rivertype.JobStateCancelled, Tags: []string{"tag1"}, }) require.True(t, mockT.Failed) @@ -336,7 +337,7 @@ func TestRequireInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag2"}, }) require.True(t, mockT.Failed) @@ -674,7 +675,7 @@ func TestRequireManyInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }, }, @@ -696,7 +697,7 @@ func TestRequireManyInsertedTx(t *testing.T) { Priority: 3, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }, }, @@ -718,7 +719,7 @@ func TestRequireManyInsertedTx(t *testing.T) { Priority: 2, Queue: "wrong-queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }, }, @@ -740,7 +741,7 @@ func TestRequireManyInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime.Add(3*time.Minute + 23*time.Second + 123*time.Microsecond), - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag1"}, }, }, @@ -761,7 +762,7 @@ func TestRequireManyInsertedTx(t *testing.T) { MaxAttempts: 78, Priority: 2, Queue: "another_queue", - State: river.JobStateCancelled, + State: rivertype.JobStateCancelled, ScheduledAt: testTime, Tags: []string{"tag1"}, }, @@ -784,7 +785,7 @@ func TestRequireManyInsertedTx(t *testing.T) { Priority: 2, Queue: "another_queue", ScheduledAt: testTime, - State: river.JobStateScheduled, + State: rivertype.JobStateScheduled, Tags: []string{"tag2"}, }, },