Skip to content
Open
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed

- Write timestamps in SQLite to always include three digits after the second like `.000`. Previously, they may have been truncated down to just `.0` in the case of trailing zeroes. [PR #1349](https://github.com/riverqueue/river/pull/1349)
- Job completion events now reflect the job's persisted outcome when it differs from the transition requested by the worker. For example, a job completed with `JobCompleteTx` before its worker returns an error emits `job_completed`, and a remotely cancelled job whose worker errors or snoozes emits `job_cancelled`. [PR #1350](https://github.com/riverqueue/river/pull/1350)

## [0.43.0] - 2026-08-05

Expand Down
93 changes: 90 additions & 3 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import (
"github.com/riverqueue/river/rivershared/riversharedtest"
"github.com/riverqueue/river/rivershared/startstoptest"
"github.com/riverqueue/river/rivershared/testfactory"
"github.com/riverqueue/river/rivershared/testsignal"
"github.com/riverqueue/river/rivershared/util/dbutil"
"github.com/riverqueue/river/rivershared/util/ptrutil"
"github.com/riverqueue/river/rivershared/util/randutil"
Expand All @@ -49,6 +50,18 @@ import (
"github.com/riverqueue/river/rivertype"
)

type clientJobCancelTestSignals struct {
ContinueWork testsignal.TestSignal[struct{}]
JobStarted testsignal.TestSignal[int64]
}

func (s *clientJobCancelTestSignals) Init(tb testing.TB) {
tb.Helper()

s.ContinueWork.Init(tb)
s.JobStarted.Init(tb)
}

type invalidKindArgs struct{}

func (a invalidKindArgs) Kind() string { return "this kind is invalid" }
Expand Down Expand Up @@ -1210,6 +1223,69 @@ func Test_Client_Common(t *testing.T) {
require.WithinDuration(t, time.Now(), *jobAfterCancel.FinalizedAt, 2*time.Second)
}

cancelRunningJobBeforeWorkReturnsTestHelper := func(t *testing.T) {
t.Helper()

config, bundle := setupConfig(t)
config.RetryPolicy = &retrypolicytest.RetryPolicySlow{}

client, err := NewClient(NewDriverWithoutListenNotify(bundle.dbPool), config)
require.NoError(t, err)

var signals clientJobCancelTestSignals
signals.Init(t)

type JobArgs struct {
testutil.JobArgsReflectKind[JobArgs]

Result string `json:"result"`
}

AddWorker(client.config.Workers, WorkFunc(func(ctx context.Context, job *Job[JobArgs]) error {
signals.JobStarted.Signal(job.ID)
<-signals.ContinueWork.WaitC()

switch job.Args.Result {
case "error":
return errors.New("retry later")
case "snooze":
return JobSnooze(time.Hour)
default:
return fmt.Errorf("unknown test result: %s", job.Args.Result)
}
}))

subscribeChan := subscribe(t, client)
startClient(ctx, t, client)
t.Cleanup(func() { signals.ContinueWork.Signal(struct{}{}) })

for _, result := range []string{"error", "snooze"} {
insertRes, err := client.Insert(ctx, &JobArgs{Result: result}, nil)
require.NoError(t, err)
require.Equal(t, insertRes.Job.ID, signals.JobStarted.WaitOrTimeout())

var jobAfterCancel *rivertype.JobRow
err = pgx.BeginFunc(ctx, bundle.dbPool, func(tx pgx.Tx) error {
var err error
jobAfterCancel, err = client.JobCancelTx(ctx, tx, insertRes.Job.ID)
return err
})
require.NoError(t, err)
require.Equal(t, rivertype.JobStateRunning, jobAfterCancel.State)

signals.ContinueWork.Signal(struct{}{})

event := riversharedtest.WaitOrTimeout(t, subscribeChan)
require.Equal(t, EventKindJobCancelled, event.Kind)
require.Equal(t, insertRes.Job.ID, event.Job.ID)
require.Equal(t, rivertype.JobStateCancelled, event.Job.State)

reloadedJob, err := client.JobGet(ctx, insertRes.Job.ID)
require.NoError(t, err)
require.Equal(t, rivertype.JobStateCancelled, reloadedJob.State)
}
}

t.Run("CancelRunningJob", func(t *testing.T) {
t.Parallel()

Expand All @@ -1218,6 +1294,12 @@ func Test_Client_Common(t *testing.T) {
})
})

t.Run("CancelRunningJobBeforeWorkReturns", func(t *testing.T) {
t.Parallel()

cancelRunningJobBeforeWorkReturnsTestHelper(t)
})

t.Run("CancelRunningJobWithLongPollInterval", func(t *testing.T) {
t.Parallel()

Expand Down Expand Up @@ -8077,7 +8159,7 @@ func Test_Client_JobCompletion(t *testing.T) {
require.NotNil(t, reloadedJob.FinalizedAt)
})

t.Run("JobThatIsCompletedManuallyIsNotTouchedByCompleter", func(t *testing.T) {
t.Run("JobThatIsCompletedManuallyBeforeReturningErrIsNotTouchedByCompleter", func(t *testing.T) {
t.Parallel()

client, bundle := setup(t, newTestConfig(t, ""))
Expand All @@ -8096,17 +8178,22 @@ func Test_Client_JobCompletion(t *testing.T) {
updatedJob, err = JobCompleteTx[*riverpgxv5.Driver](ctx, tx, job)
require.NoError(t, err)

return tx.Commit(ctx)
if err := tx.Commit(ctx); err != nil {
return err
}

return errors.New("error after committed completion")
}))

insertRes, err := client.Insert(ctx, JobArgs{}, nil)
require.NoError(t, err)

event := riversharedtest.WaitOrTimeout(t, bundle.subscribeChan)
require.Equal(t, EventKindJobCompleted, event.Kind)
require.Equal(t, insertRes.Job.ID, event.Job.ID)
require.Equal(t, rivertype.JobStateCompleted, event.Job.State)
require.Equal(t, rivertype.JobStateCompleted, updatedJob.State)
require.NotNil(t, updatedJob)
require.Equal(t, rivertype.JobStateCompleted, updatedJob.State)
require.NotNil(t, event.Job.FinalizedAt)
require.NotNil(t, updatedJob.FinalizedAt)

Expand Down
60 changes: 45 additions & 15 deletions internal/jobcompleter/job_completer.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,48 @@ type CompleterJobUpdated struct {
Reason riverdriver.JobSetStateReason
}

func newCompleterJobUpdated(job *rivertype.JobRow, stats *jobstats.JobStatistics, requestedReason riverdriver.JobSetStateReason) CompleterJobUpdated {
var reason riverdriver.JobSetStateReason
switch job.State {
case rivertype.JobStateAvailable:
// Available is ambiguous: failed jobs being retried immediately, short
// snoozes, and jobs interrupted by client shutdown all use it. Preserve
// the requested reason when it describes one of those transitions.
switch requestedReason {
case riverdriver.JobSetStateReasonFailed, riverdriver.JobSetStateReasonInterrupted, riverdriver.JobSetStateReasonSnoozed:
reason = requestedReason
case riverdriver.JobSetStateReasonCancelled, riverdriver.JobSetStateReasonCompleted:
reason = riverdriver.JobSetStateReasonFailed
default:
panic("completion subscriber received an unknown reason, river bug")
}
case rivertype.JobStateCancelled:
reason = riverdriver.JobSetStateReasonCancelled

case rivertype.JobStateCompleted:
reason = riverdriver.JobSetStateReasonCompleted

case rivertype.JobStateDiscarded, rivertype.JobStateRetryable:
reason = riverdriver.JobSetStateReasonFailed

case rivertype.JobStateScheduled:
reason = riverdriver.JobSetStateReasonSnoozed

case rivertype.JobStatePending, rivertype.JobStateRunning:
panic("completion subscriber received a job that wasn't finalized, river bug")

default:
// linter exhaustive rule prevents this from being reached.
panic("completion subscriber received a job with an unknown state, river bug")
}

return CompleterJobUpdated{
Job: job,
JobStats: stats,
Reason: reason,
}
}

type InlineCompleter struct {
baseservice.BaseService
startstop.BaseStartStop
Expand Down Expand Up @@ -101,11 +143,7 @@ func (c *InlineCompleter) JobSetStateIfRunning(ctx context.Context, stats *jobst
}

stats.CompleteDuration = c.Time.Now().Sub(start)
c.subscribeCh <- []CompleterJobUpdated{{
Job: jobs[0],
JobStats: stats,
Reason: params.Reason,
}}
c.subscribeCh <- []CompleterJobUpdated{newCompleterJobUpdated(jobs[0], stats, params.Reason)}

return nil
}
Expand Down Expand Up @@ -216,11 +254,7 @@ func (c *AsyncCompleter) JobSetStateIfRunning(ctx context.Context, stats *jobsta
}

stats.CompleteDuration = c.Time.Now().Sub(start)
c.subscribeCh <- []CompleterJobUpdated{{
Job: jobs[0],
JobStats: stats,
Reason: params.Reason,
}}
c.subscribeCh <- []CompleterJobUpdated{newCompleterJobUpdated(jobs[0], stats, params.Reason)}

return nil
})
Expand Down Expand Up @@ -508,11 +542,7 @@ func (c *BatchCompleter) handleBatch(ctx context.Context) error {
for _, jobRow := range jobRows {
setState := setStateBatch[jobRow.ID]
setState.Stats.CompleteDuration = completeTime.Sub(setState.StartTime)
events = append(events, CompleterJobUpdated{
Job: jobRow,
JobStats: setState.Stats,
Reason: setState.Params.Reason,
})
events = append(events, newCompleterJobUpdated(jobRow, setState.Stats, setState.Params.Reason))
}

if len(events) > 0 {
Expand Down
35 changes: 35 additions & 0 deletions internal/jobcompleter/job_completer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,41 @@ func (m *partialExecutorTxMock) JobSetStateIfRunningMany(ctx context.Context, pa
return m.partial.JobSetStateIfRunningMany(ctx, params)
}

func TestNewCompleterJobUpdated(t *testing.T) {
t.Parallel()

tests := []struct {
expectedReason riverdriver.JobSetStateReason
name string
requestedReason riverdriver.JobSetStateReason
state rivertype.JobState
}{
{expectedReason: riverdriver.JobSetStateReasonFailed, name: "AvailableFailed", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateAvailable},
{expectedReason: riverdriver.JobSetStateReasonInterrupted, name: "AvailableInterrupted", requestedReason: riverdriver.JobSetStateReasonInterrupted, state: rivertype.JobStateAvailable},
{expectedReason: riverdriver.JobSetStateReasonFailed, name: "AvailableRequestedCompleted", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateAvailable},
{expectedReason: riverdriver.JobSetStateReasonSnoozed, name: "AvailableSnoozed", requestedReason: riverdriver.JobSetStateReasonSnoozed, state: rivertype.JobStateAvailable},
{expectedReason: riverdriver.JobSetStateReasonCancelled, name: "Cancelled", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateCancelled},
{expectedReason: riverdriver.JobSetStateReasonCompleted, name: "Completed", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateCompleted},
{expectedReason: riverdriver.JobSetStateReasonFailed, name: "Discarded", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateDiscarded},
{expectedReason: riverdriver.JobSetStateReasonFailed, name: "Retryable", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateRetryable},
{expectedReason: riverdriver.JobSetStateReasonSnoozed, name: "Scheduled", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateScheduled},
}

for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
t.Parallel()

job := &rivertype.JobRow{State: test.state}
stats := &jobstats.JobStatistics{}

update := newCompleterJobUpdated(job, stats, test.requestedReason)
require.Same(t, job, update.Job)
require.Same(t, stats, update.JobStats)
require.Equal(t, test.expectedReason, update.Reason)
})
}
}

func TestInlineJobCompleter_Complete(t *testing.T) {
t.Parallel()

Expand Down
5 changes: 5 additions & 0 deletions job_complete_tx.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,11 @@ import (
// // handle error
// }
//
// After successfully committing the transaction, the worker should normally
// return nil. If it returns an error instead, the committed completion takes
// precedence. the job remains completed and a resulting subscribe event has
// kind EventKindJobCompleted. The error has no effect.
//
// 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 != rivertype.JobStateRunning {
Expand Down
8 changes: 4 additions & 4 deletions riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading