Skip to content
Merged
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 @@ -48,6 +48,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Fixed the `Job appears to be stuck` log line reporting the client-level `JobTimeout` instead of the worker-level timeout when a worker overrides `Timeout`. [PR #1394](https://github.com/riverqueue/river/pull/1394).
- Fixed `river bench` inserting every benchmark job with a `num` arg of `0` instead of numbering jobs sequentially. [PR #1379](https://github.com/riverqueue/river/pull/1379).
- Fixed `rivermigrate` leaving `river_migration` rows behind after migrating a non-main migration line down through its version 1, which caused a later up migration of that line to skip version 1. With `MigrateTx`, rows for every removed version were left behind. [PR #1378](https://github.com/riverqueue/river/pull/1378).
- Fixed the job completer panicking when a job it was finalizing had its state changed concurrently, like being moved to `pending` out of band, or being rescued while the completer's update waited on the row lock (in which case PostgreSQL returns the job's pre-update `running` row). Such jobs are now skipped without emitting a completion event. [PR #1383](https://github.com/riverqueue/river/pull/1383).

## [0.47.0] - 2026-09-01

Expand Down
49 changes: 39 additions & 10 deletions internal/jobcompleter/job_completer.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,30 @@ type CompleterJobUpdated struct {
Reason riverdriver.JobSetStateReason
}

func completerJobUpdatedFromStateAndReason(job *rivertype.JobRow, stats *jobstats.JobStatistics, requestedReason riverdriver.JobSetStateReason) CompleterJobUpdated {
// completerJobUpdated maps a job row returned from JobSetStateIfRunningMany to
// an update event, logging and returning false if the row isn't in a state
// that an event should be emitted for (see
// completerJobUpdatedFromStateAndReason).
func completerJobUpdated(ctx context.Context, baseService *baseservice.BaseService, job *rivertype.JobRow, stats *jobstats.JobStatistics, requestedReason riverdriver.JobSetStateReason) (CompleterJobUpdated, bool) {
update, emit := completerJobUpdatedFromStateAndReason(job, stats, requestedReason)
if !emit {
baseService.Logger.DebugContext(ctx, baseService.Name+": Job wasn't finalized because its state was changed concurrently; skipping completion event",
slog.Int64("job_id", job.ID),
slog.String("job_state", string(job.State)),
)
}
return update, emit
}

// completerJobUpdatedFromStateAndReason maps a job row returned from
// JobSetStateIfRunningMany to an update event. JobSetStateIfRunningMany returns
// rows that it didn't update (because they were no longer running) in their
// current state, so a returned row may be in a state that isn't the result of
// a completion, like `pending` if it was moved there out of band, or `running`
// if Postgres returned the row's pre-statement version after a concurrent
// transaction changed its state (e.g. a rescue). Returns false for those rows,
// which shouldn't produce an event.
func completerJobUpdatedFromStateAndReason(job *rivertype.JobRow, stats *jobstats.JobStatistics, requestedReason riverdriver.JobSetStateReason) (CompleterJobUpdated, bool) {
var reason riverdriver.JobSetStateReason
switch job.State {
case rivertype.JobStateAvailable:
Expand Down Expand Up @@ -77,20 +100,20 @@ func completerJobUpdatedFromStateAndReason(job *rivertype.JobRow, stats *jobstat

case rivertype.JobStatePending, rivertype.JobStateRunning:
// Neither state represents a finalized job, so emitting a completion
// event would be misleading. Reaching this case indicates a River bug or
// that the job's state was changed out of band during finalization.
panic("completion subscriber received a job that wasn't finalized, river bug")
// event would be misleading.
return CompleterJobUpdated{}, false

default:
// linter exhaustive rule prevents this from being reached.
panic("completion subscriber received a job with an unknown state, river bug")
// An unknown state (e.g. one added by a newer migration) isn't
// something this completer finalized either.
return CompleterJobUpdated{}, false
}

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

type InlineCompleter struct {
Expand Down Expand Up @@ -146,7 +169,9 @@ func (c *InlineCompleter) JobSetStateIfRunning(ctx context.Context, stats *jobst
}

stats.CompleteDuration = c.Time.Now().Sub(start)
c.subscribeCh <- []CompleterJobUpdated{completerJobUpdatedFromStateAndReason(jobs[0], stats, params.Reason)}
if update, ok := completerJobUpdated(ctx, &c.BaseService, jobs[0], stats, params.Reason); ok {
c.subscribeCh <- []CompleterJobUpdated{update}
}

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

stats.CompleteDuration = c.Time.Now().Sub(start)
c.subscribeCh <- []CompleterJobUpdated{completerJobUpdatedFromStateAndReason(jobs[0], stats, params.Reason)}
if update, ok := completerJobUpdated(ctx, &c.BaseService, jobs[0], stats, params.Reason); ok {
c.subscribeCh <- []CompleterJobUpdated{update}
}

return nil
})
Expand Down Expand Up @@ -545,7 +572,9 @@ 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, completerJobUpdatedFromStateAndReason(jobRow, setState.Stats, setState.Params.Reason))
if update, ok := completerJobUpdated(ctx, &c.BaseService, jobRow, setState.Stats, setState.Params.Reason); ok {
events = append(events, update)
}
}

if len(events) > 0 {
Expand Down
81 changes: 71 additions & 10 deletions internal/jobcompleter/job_completer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,20 +75,24 @@ func TestCompleterJobUpdatedFromStateAndReason(t *testing.T) {
t.Parallel()

tests := []struct {
expectedOK bool
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},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonFailed, name: "AvailableFailed", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateAvailable},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonInterrupted, name: "AvailableInterrupted", requestedReason: riverdriver.JobSetStateReasonInterrupted, state: rivertype.JobStateAvailable},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonFailed, name: "AvailableRequestedCompleted", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateAvailable},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonSnoozed, name: "AvailableSnoozed", requestedReason: riverdriver.JobSetStateReasonSnoozed, state: rivertype.JobStateAvailable},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonCancelled, name: "Cancelled", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateCancelled},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonCompleted, name: "Completed", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateCompleted},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonFailed, name: "Discarded", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateDiscarded},
{expectedOK: false, name: "Pending", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStatePending},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonFailed, name: "Retryable", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateRetryable},
{expectedOK: false, name: "Running", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobStateRunning},
{expectedOK: true, expectedReason: riverdriver.JobSetStateReasonSnoozed, name: "Scheduled", requestedReason: riverdriver.JobSetStateReasonFailed, state: rivertype.JobStateScheduled},
{expectedOK: false, name: "UnknownState", requestedReason: riverdriver.JobSetStateReasonCompleted, state: rivertype.JobState("unknown_state")},
}

for _, test := range tests {
Expand All @@ -98,7 +102,12 @@ func TestCompleterJobUpdatedFromStateAndReason(t *testing.T) {
job := &rivertype.JobRow{State: test.state}
stats := &jobstats.JobStatistics{}

update := completerJobUpdatedFromStateAndReason(job, stats, test.requestedReason)
update, ok := completerJobUpdatedFromStateAndReason(job, stats, test.requestedReason)
require.Equal(t, test.expectedOK, ok)
if !test.expectedOK {
require.Equal(t, CompleterJobUpdated{}, update)
return
}
require.Same(t, job, update.Job)
require.Same(t, stats, update.JobStats)
require.Equal(t, test.expectedReason, update.Reason)
Expand Down Expand Up @@ -1116,6 +1125,58 @@ func testCompleter[TCompleter JobCompleter](
require.Equal(t, riverdriver.JobSetStateReasonSnoozed, job3Update.Reason)
})

// JobSetStateIfRunningMany returns rows it didn't update in their current
// state. A job moved to `pending` out of band while its completion is in
// flight, or a stale `running` version that Postgres returns when a
// concurrent transaction (like a rescue) changes the job's state first,
// must not produce a completion event or crash the completer.
t.Run("NonFinalizedJobsSkipped", func(t *testing.T) {
t.Parallel()

completer, bundle := setup(t)

var (
job1 = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Schema: bundle.schema, State: new(rivertype.JobStateRunning)})
job2 = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Schema: bundle.schema, State: new(rivertype.JobStatePending)})
job3 = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Schema: bundle.schema, State: new(rivertype.JobStateRunning)})
)

// Simulate the concurrent state change race for job3 by returning its
// row as it was before the statement.
execMock := NewPartialExecutorMock(bundle.exec)
execMock.JobSetStateIfRunningManyFunc = func(ctx context.Context, params *riverdriver.JobSetStateIfRunningManyParams) ([]*rivertype.JobRow, error) {
jobRows, err := bundle.exec.JobSetStateIfRunningMany(ctx, params)
if err != nil {
return nil, err
}
for i, jobRow := range jobRows {
if jobRow.ID == job3.ID {
jobRows[i] = job3
}
}
return jobRows, nil
}
setExec(completer, execMock)

require.NoError(t, completer.JobSetStateIfRunning(ctx, &jobstats.JobStatistics{}, riverdriver.JobSetStateCompleted(job1.ID, time.Now(), nil)))
require.NoError(t, completer.JobSetStateIfRunning(ctx, &jobstats.JobStatistics{}, riverdriver.JobSetStateCompleted(job2.ID, time.Now(), nil)))
require.NoError(t, completer.JobSetStateIfRunning(ctx, &jobstats.JobStatistics{}, riverdriver.JobSetStateCompleted(job3.ID, time.Now(), nil)))

completer.Stop()

// Stop closes the subscribe channel, so this collects every update.
var jobUpdates []CompleterJobUpdated
for updates := range bundle.subscribeCh {
jobUpdates = append(jobUpdates, updates...)
}
require.Len(t, jobUpdates, 1)
require.Equal(t, job1.ID, jobUpdates[0].Job.ID)
require.Equal(t, riverdriver.JobSetStateReasonCompleted, jobUpdates[0].Reason)

requireState(t, bundle, job1.ID, rivertype.JobStateCompleted)
requireState(t, bundle, job2.ID, rivertype.JobStatePending)
})

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

Expand Down
Loading