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 @@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed

- Guard against empty job slice returned by `JobSetStateIfRunningMany` when a job has been deleted mid-run. [PR #1308](https://github.com/riverqueue/river/pull/1308).
- Fixed `JobRescuer` pagination so a full batch of running jobs with disabled or longer worker-specific timeouts can't prevent later stuck jobs from being rescued. [PR #1318](https://github.com/riverqueue/river/pull/1318).

## [0.40.0] - 2026-07-02

Expand Down
17 changes: 12 additions & 5 deletions internal/maintenance/job_rescuer.go
Original file line number Diff line number Diff line change
Expand Up @@ -191,11 +191,14 @@ type metadataWithCancelAttemptedAt struct {
}

func (s *JobRescuer) runOnce(ctx context.Context) (*rescuerRunOnceResult, error) {
var afterID int64

res := &rescuerRunOnceResult{}
stuckHorizon := time.Now().Add(-s.Config.RescueAfter)

for {
stuckHorizon := time.Now().Add(-s.Config.RescueAfter)
stuckJobs, err := s.getStuckJobs(ctx, stuckHorizon)
batchSize := s.batchSize()
stuckJobs, err := s.getStuckJobs(ctx, afterID, batchSize, stuckHorizon)
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
s.reducedBatchSizeBreaker.Trip()
Expand All @@ -207,6 +210,9 @@ func (s *JobRescuer) runOnce(ctx context.Context) (*rescuerRunOnceResult, error)
s.reducedBatchSizeBreaker.ResetIfNotOpen()

s.TestSignals.FetchedBatch.Signal(struct{}{})
if len(stuckJobs) > 0 {
afterID = stuckJobs[len(stuckJobs)-1].ID
}

now := time.Now().UTC()

Expand Down Expand Up @@ -277,7 +283,7 @@ func (s *JobRescuer) runOnce(ctx context.Context) (*rescuerRunOnceResult, error)

// Number of rows fetched was less than query `LIMIT` which means work is
// done for this round:
if len(stuckJobs) < s.batchSize() {
if len(stuckJobs) < batchSize {
break
}

Expand All @@ -287,12 +293,13 @@ func (s *JobRescuer) runOnce(ctx context.Context) (*rescuerRunOnceResult, error)
return res, nil
}

func (s *JobRescuer) getStuckJobs(ctx context.Context, stuckHorizon time.Time) ([]*rivertype.JobRow, error) {
func (s *JobRescuer) getStuckJobs(ctx context.Context, afterID int64, batchSize int, stuckHorizon time.Time) ([]*rivertype.JobRow, error) {
ctx, cancelFunc := context.WithTimeout(ctx, riversharedmaintenance.TimeoutDefault)
defer cancelFunc()

return s.Config.Pilot.JobGetStuck(ctx, s.exec, &riverdriver.JobGetStuckParams{
Max: s.batchSize(),
AfterID: afterID,
Max: batchSize,
Schema: s.Config.Schema,
StuckHorizon: stuckHorizon,
})
Expand Down
26 changes: 26 additions & 0 deletions internal/maintenance/job_rescuer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -292,6 +292,32 @@ func TestJobRescuer(t *testing.T) {
}
})

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

rescuer, bundle := setup(t)
rescuer.Config.Default = 3

noTimeoutJobs := make([]*rivertype.JobRow, rescuer.Config.Default+1)
for i := range noTimeoutJobs {
noTimeoutJobs[i] = testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Kind: ptrutil.Ptr(rescuerJobKindNoTimeout), State: ptrutil.Ptr(rivertype.JobStateRunning), AttemptedAt: ptrutil.Ptr(bundle.rescueHorizon.Add(-24 * time.Hour)), MaxAttempts: ptrutil.Ptr(5)})
}
jobToRescue := testfactory.Job(ctx, t, bundle.exec, &testfactory.JobOpts{Kind: ptrutil.Ptr(rescuerJobKind), State: ptrutil.Ptr(rivertype.JobStateRunning), AttemptedAt: ptrutil.Ptr(bundle.rescueHorizon.Add(-1 * time.Hour)), MaxAttempts: ptrutil.Ptr(5)})

_, err := rescuer.runOnce(ctx)
require.NoError(t, err)

for _, job := range noTimeoutJobs {
jobAfter, err := bundle.exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID, Schema: rescuer.Config.Schema})
require.NoError(t, err)
require.Equal(t, rivertype.JobStateRunning, jobAfter.State)
}

jobToRescueAfter, err := bundle.exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: jobToRescue.ID, Schema: rescuer.Config.Schema})
require.NoError(t, err)
require.Equal(t, rivertype.JobStateRetryable, jobToRescueAfter.State)
})

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

Expand Down
1 change: 1 addition & 0 deletions riverdriver/river_driver_interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -438,6 +438,7 @@ type JobGetByKindManyParams struct {
}

type JobGetStuckParams struct {
AfterID int64
Max int
Schema string
StuckHorizon time.Time
Expand Down
8 changes: 5 additions & 3 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.

1 change: 1 addition & 0 deletions riverdriver/riverdatabasesql/river_database_sql_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -343,6 +343,7 @@ func (e *Executor) JobGetByKindMany(ctx context.Context, params *riverdriver.Job

func (e *Executor) JobGetStuck(ctx context.Context, params *riverdriver.JobGetStuckParams) ([]*rivertype.JobRow, error) {
jobs, err := dbsqlc.New().JobGetStuck(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.JobGetStuckParams{
AfterID: params.AfterID,
Max: int32(min(params.Max, math.MaxInt32)), //nolint:gosec
StuckHorizon: params.StuckHorizon,
})
Expand Down
13 changes: 11 additions & 2 deletions riverdriver/riverdrivertest/job_read.go
Original file line number Diff line number Diff line change
Expand Up @@ -528,8 +528,8 @@ func exerciseJobRead[TTx any](ctx context.Context, t *testing.T, executorWithTx

t.Logf("stuckJob1 full = %s", spew.Sdump(stuckJob1))

// Not returned because we put a maximum of two.
_ = testfactory.Job(ctx, t, exec, &testfactory.JobOpts{AttemptedAt: &beforeHorizon, State: ptrutil.Ptr(rivertype.JobStateRunning)})
// Not returned on the first page because we put a maximum of two.
stuckJob3 := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{AttemptedAt: &beforeHorizon, State: ptrutil.Ptr(rivertype.JobStateRunning)})

// Not stuck because not in running state.
_ = testfactory.Job(ctx, t, exec, &testfactory.JobOpts{State: ptrutil.Ptr(rivertype.JobStateAvailable)})
Expand All @@ -545,6 +545,15 @@ func exerciseJobRead[TTx any](ctx context.Context, t *testing.T, executorWithTx
require.NoError(t, err)
require.Equal(t, []int64{stuckJob1.ID, stuckJob2.ID},
sliceutil.Map(stuckJobs, func(j *rivertype.JobRow) int64 { return j.ID }))

stuckJobs, err = exec.JobGetStuck(ctx, &riverdriver.JobGetStuckParams{
AfterID: stuckJob2.ID,
Max: 2,
StuckHorizon: horizon,
})
require.NoError(t, err)
require.Equal(t, []int64{stuckJob3.ID},
sliceutil.Map(stuckJobs, func(j *rivertype.JobRow) int64 { return j.ID }))
})

t.Run("JobKindList", func(t *testing.T) {
Expand Down
3 changes: 2 additions & 1 deletion riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql
Original file line number Diff line number Diff line change
Expand Up @@ -258,6 +258,7 @@ ORDER BY id;
SELECT *
FROM /* TEMPLATE: schema */river_job
WHERE state = 'running'
AND id > @after_id::bigint
AND attempted_at < @stuck_horizon::timestamptz
ORDER BY id
LIMIT @max;
Expand Down Expand Up @@ -726,4 +727,4 @@ SET
metadata = CASE WHEN @metadata_do_update::boolean THEN @metadata::jsonb ELSE metadata END,
state = CASE WHEN @state_do_update::boolean THEN @state::/* TEMPLATE: schema */river_job_state ELSE state END
WHERE id = @id
RETURNING *;
RETURNING *;
8 changes: 5 additions & 3 deletions riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go

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

1 change: 1 addition & 0 deletions riverdriver/riverpgxv5/river_pgx_v5_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -347,6 +347,7 @@ func (e *Executor) JobGetByKindMany(ctx context.Context, params *riverdriver.Job

func (e *Executor) JobGetStuck(ctx context.Context, params *riverdriver.JobGetStuckParams) ([]*rivertype.JobRow, error) {
jobs, err := dbsqlc.New().JobGetStuck(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.JobGetStuckParams{
AfterID: params.AfterID,
Max: int32(min(params.Max, math.MaxInt32)), //nolint:gosec
StuckHorizon: params.StuckHorizon,
})
Expand Down
1 change: 1 addition & 0 deletions riverdriver/riversqlite/internal/dbsqlc/river_job.sql
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,7 @@ ORDER BY id;
SELECT *
FROM /* TEMPLATE: schema */river_job
WHERE state = 'running'
AND id > @after_id
AND attempted_at < cast(@stuck_horizon AS text)
ORDER BY id
LIMIT @max;
Expand Down
8 changes: 5 additions & 3 deletions riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go

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

1 change: 1 addition & 0 deletions riverdriver/riversqlite/river_sqlite_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -517,6 +517,7 @@ func (e *Executor) JobGetByKindMany(ctx context.Context, params *riverdriver.Job

func (e *Executor) JobGetStuck(ctx context.Context, params *riverdriver.JobGetStuckParams) ([]*rivertype.JobRow, error) {
jobs, err := dbsqlc.New().JobGetStuck(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.JobGetStuckParams{
AfterID: params.AfterID,
Max: int64(params.Max),
StuckHorizon: timeString(params.StuckHorizon),
})
Expand Down
Loading