Skip to content

Commit 5341d91

Browse files
committed
retry sequence jobs into pending state
Fixes #971.
1 parent 24a8c2f commit 5341d91

6 files changed

Lines changed: 40 additions & 5 deletions

File tree

riverdriver/riverdatabasesql/internal/dbsqlc/river_job.sql.go

Lines changed: 4 additions & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

riverdriver/riverdrivertest/riverdrivertest.go

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2302,6 +2302,32 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,
23022302
require.Error(t, err)
23032303
require.ErrorIs(t, err, rivertype.ErrNotFound)
23042304
})
2305+
2306+
t.Run("SetsStatePendingWhenSeqKeyPresent", func(t *testing.T) {
2307+
t.Parallel()
2308+
2309+
exec, bundle := setup(ctx, t)
2310+
2311+
now := time.Now().UTC()
2312+
2313+
job := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{
2314+
ScheduledAt: ptrutil.Ptr(now.Add(1 * time.Hour)),
2315+
State: ptrutil.Ptr(rivertype.JobStateRetryable),
2316+
Metadata: []byte(`{"seq_key":"foo"}`),
2317+
})
2318+
2319+
jobAfter, err := exec.JobRetry(ctx, &riverdriver.JobRetryParams{
2320+
ID: job.ID,
2321+
Now: &now,
2322+
})
2323+
require.NoError(t, err)
2324+
require.Equal(t, rivertype.JobStatePending, jobAfter.State)
2325+
require.WithinDuration(t, now, jobAfter.ScheduledAt, bundle.driver.TimePrecision())
2326+
2327+
jobUpdated, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID})
2328+
require.NoError(t, err)
2329+
require.Equal(t, rivertype.JobStatePending, jobUpdated.State)
2330+
})
23052331
})
23062332

23072333
t.Run("JobSchedule", func(t *testing.T) {

riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -418,7 +418,10 @@ WITH job_to_update AS (
418418
updated_job AS (
419419
UPDATE /* TEMPLATE: schema */river_job
420420
SET
421-
state = 'available',
421+
state = CASE WHEN river_job.metadata ? 'seq_key'
422+
THEN 'pending'::/* TEMPLATE: schema */river_job_state
423+
ELSE 'available'::/* TEMPLATE: schema */river_job_state
424+
END,
422425
max_attempts = CASE WHEN attempt = max_attempts THEN max_attempts + 1 ELSE max_attempts END,
423426
finalized_at = NULL,
424427
scheduled_at = coalesce(sqlc.narg('now')::timestamptz, now())

riverdriver/riverpgxv5/internal/dbsqlc/river_job.sql.go

Lines changed: 4 additions & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

riverdriver/riversqlite/internal/dbsqlc/river_job.sql

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -320,7 +320,7 @@ WHERE id = @id;
320320
-- name: JobRetry :one
321321
UPDATE /* TEMPLATE: schema */river_job
322322
SET
323-
state = 'available',
323+
state = CASE WHEN river_job.metadata -> 'seq_key' IS NOT NULL THEN 'pending' ELSE 'available' END,
324324
max_attempts = CASE WHEN attempt = max_attempts THEN max_attempts + 1 ELSE max_attempts END,
325325
finalized_at = NULL,
326326
scheduled_at = coalesce(cast(sqlc.narg('now') AS text), datetime('now', 'subsec'))

riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

0 commit comments

Comments
 (0)