From e2678bfb36254ab26234f3993237a73bc18d5c77 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Thu, 24 Sep 2026 12:33:24 -0500 Subject: [PATCH] skip non-finalized rows in the job completer `JobSetStateIfRunningMany` returns rows it didn't update in their current state, and the completers map every returned row to a completion event. Mapping panics for `pending` and `running` rows, and because the panic happens in a completer goroutine it takes down the whole process. Both states are reachable. A job moved to `pending` out of band while its completion is in flight comes back as `pending`. On PostgreSQL, if another transaction (like the rescuer) changes a running job's state while the completer's `UPDATE` waits on the row lock, the update re-checks the row and skips it, but the statement's `SELECT` over non-updated rows still reads the snapshot taken before the lock wait and returns the job's old `running` version. Have the mapping report whether a row should produce an event, and skip rows in `pending`, `running`, or an unrecognized state in the inline, async, and batch completers, logging at debug level instead. A new completer subtest covers a `pending` row and a stale `running` row for all three completers. --- CHANGELOG.md | 1 + internal/jobcompleter/job_completer.go | 49 ++++++++++--- internal/jobcompleter/job_completer_test.go | 81 ++++++++++++++++++--- 3 files changed, 111 insertions(+), 20 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0917e8b7d..bc47dd3ba 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/internal/jobcompleter/job_completer.go b/internal/jobcompleter/job_completer.go index 1f6bfae16..44e17a8f8 100644 --- a/internal/jobcompleter/job_completer.go +++ b/internal/jobcompleter/job_completer.go @@ -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: @@ -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 { @@ -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 } @@ -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 }) @@ -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 { diff --git a/internal/jobcompleter/job_completer_test.go b/internal/jobcompleter/job_completer_test.go index 749d57d39..86bb8291a 100644 --- a/internal/jobcompleter/job_completer_test.go +++ b/internal/jobcompleter/job_completer_test.go @@ -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 { @@ -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) @@ -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()