From 20b7323cfdbb6e83fc3dad081ed84a1d09c07958 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Mon, 7 Sep 2026 19:31:41 -0500 Subject: [PATCH] guard job rescue against stale snapshots A worker can complete or release a job after the rescuer fetches it but before the rescue update executes. The job can also be claimed again, leaving it running under a new worker when the stale rescue arrives. Require jobs to still be `running` with `attempted_at` before the original rescue horizon in both PostgreSQL and SQLite updates. Forward the horizon through every driver so stale rescues preserve completed jobs and fresh attempts, including their errors, metadata, and timestamps. Add shared driver coverage for completion, immediate retry, and worker interruption between fetch and rescue, plus strict horizon boundaries and mixed batches containing eligible jobs. Document the fix in the changelog. Fixes #1302. --- CHANGELOG.md | 4 + .../internal/dbsqlc/river_job.sql.go | 24 +- .../river_database_sql_driver.go | 11 +- riverdriver/riverdrivertest/job_update.go | 266 +++++++++++++++--- .../riverpgxv5/internal/dbsqlc/river_job.sql | 4 +- .../internal/dbsqlc/river_job.sql.go | 24 +- riverdriver/riverpgxv5/river_pgx_v5_driver.go | 11 +- .../riversqlite/internal/dbsqlc/river_job.sql | 4 +- .../internal/dbsqlc/river_job.sql.go | 14 +- .../riversqlite/river_sqlite_driver.go | 11 +- 10 files changed, 292 insertions(+), 81 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b875f8a..6e6445a0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- Fixed `JobRescuer` overwriting jobs that complete, leave the running state, or are claimed again by another worker after being fetched for rescue, preserving their state, errors, metadata, and timestamps across PostgreSQL and SQLite drivers. [Issue #1302](https://github.com/riverqueue/river/issues/1302). + ## [0.47.0] - 2026-09-01 ### Added diff --git a/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go b/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go index a197f7bb..024a45a7 100644 --- a/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go @@ -1231,26 +1231,30 @@ SET state = updated_job.state FROM ( SELECT - unnest($1::bigint[]) AS id, - unnest($2::jsonb[]) AS error, - nullif(unnest($3::timestamptz[]), '0001-01-01 00:00:00 +0000') AS finalized_at, - unnest($4::timestamptz[]) AS scheduled_at, - unnest($5::text[])::/* TEMPLATE: schema */river_job_state AS state + unnest($2::bigint[]) AS id, + unnest($3::jsonb[]) AS error, + nullif(unnest($4::timestamptz[]), '0001-01-01 00:00:00 +0000') AS finalized_at, + unnest($5::timestamptz[]) AS scheduled_at, + unnest($6::text[])::/* TEMPLATE: schema */river_job_state AS state ) AS updated_job WHERE river_job.id = updated_job.id + AND river_job.state = 'running' + AND river_job.attempted_at < $1::timestamptz ` type JobRescueManyParams struct { - ID []int64 - Error []string - FinalizedAt []time.Time - ScheduledAt []time.Time - State []string + StuckHorizon time.Time + ID []int64 + Error []string + FinalizedAt []time.Time + ScheduledAt []time.Time + State []string } // Run by the rescuer to queue for retry or discard depending on job state. func (q *Queries) JobRescueMany(ctx context.Context, db DBTX, arg *JobRescueManyParams) error { _, err := db.ExecContext(ctx, jobRescueMany, + arg.StuckHorizon, pq.Array(arg.ID), pq.Array(arg.Error), pq.Array(arg.FinalizedAt), diff --git a/riverdriver/riverdatabasesql/river_database_sql_driver.go b/riverdriver/riverdatabasesql/river_database_sql_driver.go index b672c77d..8677d0fd 100644 --- a/riverdriver/riverdatabasesql/river_database_sql_driver.go +++ b/riverdriver/riverdatabasesql/river_database_sql_driver.go @@ -630,11 +630,12 @@ func (e *Executor) JobList(ctx context.Context, params *riverdriver.JobListParam func (e *Executor) JobRescueMany(ctx context.Context, params *riverdriver.JobRescueManyParams) (*struct{}, error) { if err := dbsqlc.New().JobRescueMany(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.JobRescueManyParams{ - ID: params.ID, - Error: sliceutil.Map(params.Error, func(e []byte) string { return string(e) }), - FinalizedAt: sliceutil.Map(params.FinalizedAt, func(t *time.Time) time.Time { return ptrutil.ValOrDefault(t, time.Time{}) }), - ScheduledAt: params.ScheduledAt, - State: params.State, + ID: params.ID, + Error: sliceutil.Map(params.Error, func(e []byte) string { return string(e) }), + FinalizedAt: sliceutil.Map(params.FinalizedAt, func(t *time.Time) time.Time { return ptrutil.ValOrDefault(t, time.Time{}) }), + ScheduledAt: params.ScheduledAt, + State: params.State, + StuckHorizon: params.StuckHorizon, }); err != nil { return nil, interpretError(err) } diff --git a/riverdriver/riverdrivertest/job_update.go b/riverdriver/riverdrivertest/job_update.go index a9c4cea3..2512ff31 100644 --- a/riverdriver/riverdrivertest/job_update.go +++ b/riverdriver/riverdrivertest/job_update.go @@ -36,6 +36,41 @@ func exerciseJobUpdate[TTx any](ctx context.Context, t *testing.T, executorWithT } } + setStateManyParams := func(params ...*riverdriver.JobSetStateIfRunningParams) *riverdriver.JobSetStateIfRunningManyParams { + batchParams := &riverdriver.JobSetStateIfRunningManyParams{} + for _, param := range params { + var ( + attempt *int + errData []byte + finalizedAt *time.Time + scheduledAt *time.Time + ) + if param.Attempt != nil { + attempt = param.Attempt + } + if param.ErrData != nil { + errData = param.ErrData + } + if param.FinalizedAt != nil { + finalizedAt = param.FinalizedAt + } + if param.ScheduledAt != nil { + scheduledAt = param.ScheduledAt + } + + batchParams.ID = append(batchParams.ID, param.ID) + batchParams.Attempt = append(batchParams.Attempt, attempt) + batchParams.ErrData = append(batchParams.ErrData, errData) + batchParams.FinalizedAt = append(batchParams.FinalizedAt, finalizedAt) + batchParams.MetadataDoMerge = append(batchParams.MetadataDoMerge, param.MetadataDoMerge) + batchParams.MetadataUpdates = append(batchParams.MetadataUpdates, param.MetadataUpdates) + batchParams.ScheduledAt = append(batchParams.ScheduledAt, scheduledAt) + batchParams.State = append(batchParams.State, param.State) + } + + return batchParams + } + // Deliberately includes sub-millisecond precision so tests can verify that // drivers normalize timestamps to their declared precision. precisionTestTime := time.Date(2025, 4, 30, 13, 26, 39, 123400000, time.UTC) @@ -153,12 +188,14 @@ func exerciseJobUpdate[TTx any](ctx context.Context, t *testing.T, executorWithT now := precisionTestTime job1 := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{ - Metadata: []byte(`{"river:rescue_count": 5, "something": "else"}`), - State: new(rivertype.JobStateRunning), + AttemptedAt: new(now.Add(-time.Hour)), + Metadata: []byte(`{"river:rescue_count": 5, "something": "else"}`), + State: new(rivertype.JobStateRunning), }) job2 := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{ - Metadata: []byte(`{}`), - State: new(rivertype.JobStateRunning), + AttemptedAt: new(now.Add(-time.Hour)), + Metadata: []byte(`{}`), + State: new(rivertype.JobStateRunning), }) _, err := exec.JobRescueMany(ctx, &riverdriver.JobRescueManyParams{ @@ -183,6 +220,7 @@ func exerciseJobUpdate[TTx any](ctx context.Context, t *testing.T, executorWithT string(rivertype.JobStateAvailable), string(rivertype.JobStateDiscarded), }, + StuckHorizon: now, }) require.NoError(t, err) @@ -203,6 +241,191 @@ func exerciseJobUpdate[TTx any](ctx context.Context, t *testing.T, executorWithT require.JSONEq(t, `{"river:rescue_count": 1}`, string(updatedJob2.Metadata)) }) + t.Run("JobRescueMany_CompletedAfterFetch", func(t *testing.T) { + t.Parallel() + + for _, rescueState := range []rivertype.JobState{ + rivertype.JobStateCancelled, + rivertype.JobStateDiscarded, + rivertype.JobStateRetryable, + } { + t.Run(string(rescueState), func(t *testing.T) { + t.Parallel() + + exec, bundle := setup(ctx, t) + + now := precisionTestTime.Truncate(bundle.driver.TimePrecision()) + jobOpts := &testfactory.JobOpts{ + AttemptedAt: new(now.Add(-2 * time.Hour)), + Metadata: []byte(`{"river:rescue_count": 5, "something": "else"}`), + ScheduledAt: new(now.Add(-2 * time.Hour)), + State: new(rivertype.JobStateRunning), + } + job := testfactory.Job(ctx, t, exec, jobOpts) + stillRunningJob := testfactory.Job(ctx, t, exec, jobOpts) + + stuckJobs, err := exec.JobGetStuck(ctx, &riverdriver.JobGetStuckParams{ + Max: 10, + StuckHorizon: now.Add(-time.Hour), + }) + require.NoError(t, err) + require.Len(t, stuckJobs, 2) + require.Equal(t, job.ID, stuckJobs[0].ID) + require.Equal(t, stillRunningJob.ID, stuckJobs[1].ID) + + // Reproduce the interleaving deterministically: the worker completes + // after the rescuer fetches the job, but before its rescue write. + completedJobs, err := exec.JobSetStateIfRunningMany(ctx, setStateManyParams( + riverdriver.JobSetStateCompleted(job.ID, now, []byte(`{"worker": "finished"}`)), + )) + require.NoError(t, err) + require.Len(t, completedJobs, 1) + require.Equal(t, rivertype.JobStateCompleted, completedJobs[0].State) + + completedJob, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID}) + require.NoError(t, err) + + rescueAt := now.Add(time.Minute) + var finalizedAt *time.Time + if rescueState != rivertype.JobStateRetryable { + finalizedAt = &rescueAt + } + + _, err = exec.JobRescueMany(ctx, &riverdriver.JobRescueManyParams{ + ID: []int64{stuckJobs[0].ID, stuckJobs[1].ID}, + Error: [][]byte{[]byte(`{"error": "stale rescue"}`), []byte(`{"error": "stuck job rescued"}`)}, + FinalizedAt: []*time.Time{finalizedAt, finalizedAt}, + ScheduledAt: []time.Time{rescueAt, rescueAt}, + State: []string{string(rescueState), string(rescueState)}, + StuckHorizon: now.Add(-time.Hour), + }) + require.NoError(t, err) + + rescuedJob, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: stillRunningJob.ID}) + require.NoError(t, err) + require.Equal(t, rescueState, rescuedJob.State) + require.Equal(t, finalizedAt, rescuedJob.FinalizedAt) + require.Equal(t, rescueAt, rescuedJob.ScheduledAt) + require.Len(t, rescuedJob.Errors, 1) + require.Equal(t, "stuck job rescued", rescuedJob.Errors[0].Error) + require.JSONEq(t, `{"river:rescue_count": 6, "something": "else"}`, string(rescuedJob.Metadata)) + + jobAfterRescue, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID}) + require.NoError(t, err) + // Check the entire row, including errors, metadata, and timestamps. + require.Equal(t, completedJob, jobAfterRescue) + }) + } + }) + + t.Run("JobRescueMany_ReclaimedAfterFetch", func(t *testing.T) { + t.Parallel() + + for _, release := range []struct { + name string + setStateFunc func(int64, time.Time) *riverdriver.JobSetStateIfRunningParams + }{ + {"Failed", func(id int64, now time.Time) *riverdriver.JobSetStateIfRunningParams { + return riverdriver.JobSetStateErrorAvailable(id, now, []byte(`{"error":"worker failed"}`), nil) + }}, + {"Interrupted", func(id int64, now time.Time) *riverdriver.JobSetStateIfRunningParams { + return riverdriver.JobSetStateInterrupted(id, now, 0, nil) + }}, + } { + t.Run(release.name, func(t *testing.T) { + t.Parallel() + + exec, bundle := setup(ctx, t) + + now := precisionTestTime.Truncate(bundle.driver.TimePrecision()) + horizon := now.Add(-time.Hour) + job := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{ + Attempt: new(1), + AttemptedAt: new(now.Add(-2 * time.Hour)), + Metadata: []byte(`{"river:rescue_count":5}`), + ScheduledAt: new(now.Add(-2 * time.Hour)), + State: new(rivertype.JobStateRunning), + }) + + stuckJobs, err := exec.JobGetStuck(ctx, &riverdriver.JobGetStuckParams{Max: 10, StuckHorizon: horizon}) + require.NoError(t, err) + require.Len(t, stuckJobs, 1) + require.Equal(t, job.ID, stuckJobs[0].ID) + + // The old worker releases the job and a new worker claims it before + // the rescuer writes its stale snapshot. State alone still matches. + releasedJobs, err := exec.JobSetStateIfRunningMany(ctx, setStateManyParams(release.setStateFunc(job.ID, now))) + require.NoError(t, err) + require.Len(t, releasedJobs, 1) + require.Equal(t, rivertype.JobStateAvailable, releasedJobs[0].State) + + claimedJobs, err := exec.JobGetAvailable(ctx, &riverdriver.JobGetAvailableParams{ + ClientID: "new-worker", + MaxAttemptedBy: 10, + MaxToLock: 1, + Now: &now, + Queue: job.Queue, + }) + require.NoError(t, err) + require.Len(t, claimedJobs, 1) + require.Equal(t, job.ID, claimedJobs[0].ID) + require.Equal(t, rivertype.JobStateRunning, claimedJobs[0].State) + require.True(t, claimedJobs[0].AttemptedAt.After(horizon)) + + _, err = exec.JobRescueMany(ctx, &riverdriver.JobRescueManyParams{ + ID: []int64{stuckJobs[0].ID}, + Error: [][]byte{[]byte(`{"error":"stale rescue"}`)}, + FinalizedAt: []*time.Time{nil}, + ScheduledAt: []time.Time{now.Add(time.Minute)}, + State: []string{string(rivertype.JobStateRetryable)}, + StuckHorizon: horizon, + }) + require.NoError(t, err) + + jobAfter, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID}) + require.NoError(t, err) + require.Equal(t, claimedJobs[0], jobAfter) + }) + } + }) + + t.Run("JobRescueMany_StuckHorizon", func(t *testing.T) { + t.Parallel() + + exec, bundle := setup(ctx, t) + + horizon := precisionTestTime.Truncate(bundle.driver.TimePrecision()) + params := &riverdriver.JobRescueManyParams{StuckHorizon: horizon} + jobs := make([]*rivertype.JobRow, 0, 3) + for _, offset := range []time.Duration{-bundle.driver.TimePrecision(), 0, bundle.driver.TimePrecision()} { + job := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{ + AttemptedAt: new(horizon.Add(offset)), + State: new(rivertype.JobStateRunning), + }) + jobs = append(jobs, job) + params.ID = append(params.ID, job.ID) + params.Error = append(params.Error, []byte(`{"error":"stuck job rescued"}`)) + params.FinalizedAt = append(params.FinalizedAt, nil) + params.ScheduledAt = append(params.ScheduledAt, horizon) + params.State = append(params.State, string(rivertype.JobStateRetryable)) + } + + _, err := exec.JobRescueMany(ctx, params) + require.NoError(t, err) + + for i, job := range jobs { + jobAfter, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID}) + require.NoError(t, err) + if i == 0 { + require.Equal(t, rivertype.JobStateRetryable, jobAfter.State) + require.Len(t, jobAfter.Errors, 1) + } else { + // As in JobGetStuck, jobs at or after the horizon are ineligible. + require.Equal(t, job, jobAfter) + } + } + }) + t.Run("JobRetry", func(t *testing.T) { t.Parallel() @@ -547,41 +770,6 @@ func exerciseJobUpdate[TTx any](ctx context.Context, t *testing.T, executorWithT return errPayload } - setStateManyParams := func(params ...*riverdriver.JobSetStateIfRunningParams) *riverdriver.JobSetStateIfRunningManyParams { - batchParams := &riverdriver.JobSetStateIfRunningManyParams{} - for _, param := range params { - var ( - attempt *int - errData []byte - finalizedAt *time.Time - scheduledAt *time.Time - ) - if param.Attempt != nil { - attempt = param.Attempt - } - if param.ErrData != nil { - errData = param.ErrData - } - if param.FinalizedAt != nil { - finalizedAt = param.FinalizedAt - } - if param.ScheduledAt != nil { - scheduledAt = param.ScheduledAt - } - - batchParams.ID = append(batchParams.ID, param.ID) - batchParams.Attempt = append(batchParams.Attempt, attempt) - batchParams.ErrData = append(batchParams.ErrData, errData) - batchParams.FinalizedAt = append(batchParams.FinalizedAt, finalizedAt) - batchParams.MetadataDoMerge = append(batchParams.MetadataDoMerge, param.MetadataDoMerge) - batchParams.MetadataUpdates = append(batchParams.MetadataUpdates, param.MetadataUpdates) - batchParams.ScheduledAt = append(batchParams.ScheduledAt, scheduledAt) - batchParams.State = append(batchParams.State, param.State) - } - - return batchParams - } - t.Run("JobSetStateIfRunningMany_JobSetStateCompleted", func(t *testing.T) { t.Parallel() diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql index 40098ed0..26e023fa 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql @@ -501,7 +501,9 @@ FROM ( unnest(@scheduled_at::timestamptz[]) AS scheduled_at, unnest(@state::text[])::/* TEMPLATE: schema */river_job_state AS state ) AS updated_job -WHERE river_job.id = updated_job.id; +WHERE river_job.id = updated_job.id + AND river_job.state = 'running' + AND river_job.attempted_at < @stuck_horizon::timestamptz; -- name: JobRetry :one WITH job_to_update AS ( diff --git a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go index 7eb082c5..62b07c2b 100644 --- a/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go @@ -1198,26 +1198,30 @@ SET state = updated_job.state FROM ( SELECT - unnest($1::bigint[]) AS id, - unnest($2::jsonb[]) AS error, - nullif(unnest($3::timestamptz[]), '0001-01-01 00:00:00 +0000') AS finalized_at, - unnest($4::timestamptz[]) AS scheduled_at, - unnest($5::text[])::/* TEMPLATE: schema */river_job_state AS state + unnest($2::bigint[]) AS id, + unnest($3::jsonb[]) AS error, + nullif(unnest($4::timestamptz[]), '0001-01-01 00:00:00 +0000') AS finalized_at, + unnest($5::timestamptz[]) AS scheduled_at, + unnest($6::text[])::/* TEMPLATE: schema */river_job_state AS state ) AS updated_job WHERE river_job.id = updated_job.id + AND river_job.state = 'running' + AND river_job.attempted_at < $1::timestamptz ` type JobRescueManyParams struct { - ID []int64 - Error [][]byte - FinalizedAt []time.Time - ScheduledAt []time.Time - State []string + StuckHorizon time.Time + ID []int64 + Error [][]byte + FinalizedAt []time.Time + ScheduledAt []time.Time + State []string } // Run by the rescuer to queue for retry or discard depending on job state. func (q *Queries) JobRescueMany(ctx context.Context, db DBTX, arg *JobRescueManyParams) error { _, err := db.Exec(ctx, jobRescueMany, + arg.StuckHorizon, arg.ID, arg.Error, arg.FinalizedAt, diff --git a/riverdriver/riverpgxv5/river_pgx_v5_driver.go b/riverdriver/riverpgxv5/river_pgx_v5_driver.go index 34954fac..62601230 100644 --- a/riverdriver/riverpgxv5/river_pgx_v5_driver.go +++ b/riverdriver/riverpgxv5/river_pgx_v5_driver.go @@ -585,11 +585,12 @@ func (e *Executor) JobList(ctx context.Context, params *riverdriver.JobListParam func (e *Executor) JobRescueMany(ctx context.Context, params *riverdriver.JobRescueManyParams) (*struct{}, error) { err := dbsqlc.New().JobRescueMany(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.JobRescueManyParams{ - ID: params.ID, - Error: params.Error, - FinalizedAt: sliceutil.Map(params.FinalizedAt, func(t *time.Time) time.Time { return ptrutil.ValOrDefault(t, time.Time{}) }), - ScheduledAt: params.ScheduledAt, - State: params.State, + ID: params.ID, + Error: params.Error, + FinalizedAt: sliceutil.Map(params.FinalizedAt, func(t *time.Time) time.Time { return ptrutil.ValOrDefault(t, time.Time{}) }), + ScheduledAt: params.ScheduledAt, + State: params.State, + StuckHorizon: params.StuckHorizon, }) if err != nil { return nil, interpretError(err) diff --git a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql index dd284e7e..528de8d8 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql @@ -503,7 +503,9 @@ SET ) + 1 ), state = @state -WHERE id = @id; +WHERE id = @id + AND state = 'running' + AND attempted_at < cast(@stuck_horizon AS text); -- Differs by necessity from other drivers because SQLite doesn't support -- `UPDATE` inside CTEs so we can't retry if running but select otherwise. diff --git a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go index 68701243..1bccd481 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go @@ -1259,14 +1259,17 @@ SET ), state = ?4 WHERE id = ?5 + AND state = 'running' + AND attempted_at < cast(?6 AS text) ` type JobRescueParams struct { - Error interface{} - FinalizedAt *string - ScheduledAt string - State string - ID int64 + Error interface{} + FinalizedAt *string + ScheduledAt string + State string + ID int64 + StuckHorizon string } // Rescue a job. @@ -1286,6 +1289,7 @@ func (q *Queries) JobRescue(ctx context.Context, db DBTX, arg *JobRescueParams) arg.ScheduledAt, arg.State, arg.ID, + arg.StuckHorizon, ) return err } diff --git a/riverdriver/riversqlite/river_sqlite_driver.go b/riverdriver/riversqlite/river_sqlite_driver.go index c778396d..89e2fe81 100644 --- a/riverdriver/riversqlite/river_sqlite_driver.go +++ b/riverdriver/riversqlite/river_sqlite_driver.go @@ -722,11 +722,12 @@ func (e *Executor) JobRescueMany(ctx context.Context, params *riverdriver.JobRes // Should be a batch rescue, but that's currently impossible with SQLite/sqlc. https://github.com/sqlc-dev/sqlc/issues/3802 for i := range params.ID { if err := dbsqlc.New().JobRescue(ctx, dbtx, &dbsqlc.JobRescueParams{ - ID: params.ID[i], - Error: params.Error[i], - FinalizedAt: timeStringNullable(params.FinalizedAt[i]), - ScheduledAt: timeString(params.ScheduledAt[i]), - State: params.State[i], + ID: params.ID[i], + Error: params.Error[i], + FinalizedAt: timeStringNullable(params.FinalizedAt[i]), + ScheduledAt: timeString(params.ScheduledAt[i]), + State: params.State[i], + StuckHorizon: timeString(params.StuckHorizon), }); err != nil { return interpretError(err) }