From 24be9bc79ccf88cb38f2690d22341ba46b286c04 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Fri, 25 Sep 2026 12:01:56 -0500 Subject: [PATCH] honor keep-forever retention in SQLite job cleaner The job cleaner signals a retention period of -1 ("keep forever") for a finalized state by passing `CancelledDoDelete`, `CompletedDoDelete`, or `DiscardedDoDelete` as false to `JobDeleteBefore`, along with a horizon of roughly now. The Postgres query checks these flags, but the SQLite query ignores them and filters only on the horizons. On SQLite, a config like `DiscardedJobRetentionPeriod: -1` with a finite completed retention therefore deletes every discarded job on the next cleaner run. Add the `*_do_delete` flags to the SQLite `JobDeleteBefore` query so each state's clause is skipped when its deletion is disabled, matching the Postgres query, and pass them through from the SQLite driver. --- CHANGELOG.md | 1 + riverdriver/riverdrivertest/job_delete.go | 33 +++++++++++++++++++ .../riversqlite/internal/dbsqlc/river_job.sql | 6 ++-- .../internal/dbsqlc/river_job.sql.go | 16 ++++++--- .../riversqlite/river_sqlite_driver.go | 3 ++ 5 files changed, 51 insertions(+), 8 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 22610a51..a208ca37 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed - Fixed SQLite job list pagination skipping or repeating jobs by formatting cursor timestamps consistently with stored timestamps. [PR #1374](https://github.com/riverqueue/river/pull/1374). +- Fixed SQLite drivers deleting jobs in a finalized state whose retention period was set to -1 (keep forever), like `Config.DiscardedJobRetentionPeriod: -1`, whenever another state's retention period was finite. [PR #1389](https://github.com/riverqueue/river/pull/1389). - Improved PostgreSQL job listing performance when filtering by one finalized state (`completed`, `cancelled`, or `discarded`) and sorting by finalized time, including in River UI. [PR #1374](https://github.com/riverqueue/river/pull/1374). - 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. Fixes [#1302](https://github.com/riverqueue/river/issues/1302). [PR #1373](https://github.com/riverqueue/river/pull/1373). diff --git a/riverdriver/riverdrivertest/job_delete.go b/riverdriver/riverdrivertest/job_delete.go index 8b0eeace..a00f44c9 100644 --- a/riverdriver/riverdrivertest/job_delete.go +++ b/riverdriver/riverdrivertest/job_delete.go @@ -196,6 +196,39 @@ func exerciseJobDelete[TTx any](ctx context.Context, t *testing.T, executorWithT require.NoError(t, err) }) + t.Run("DoDeleteFalse", func(t *testing.T) { + t.Parallel() + + exec, _ := setup(ctx, t) + + var ( + cancelledJob = testfactory.Job(ctx, t, exec, &testfactory.JobOpts{FinalizedAt: &beforeHorizon, State: new(rivertype.JobStateCancelled)}) + completedJob = testfactory.Job(ctx, t, exec, &testfactory.JobOpts{FinalizedAt: &beforeHorizon, State: new(rivertype.JobStateCompleted)}) + discardedJob = testfactory.Job(ctx, t, exec, &testfactory.JobOpts{FinalizedAt: &beforeHorizon, State: new(rivertype.JobStateDiscarded)}) + ) + + // Discarded jobs are kept even though they're older than their + // horizon because their deletion is disabled. + numDeleted, err := exec.JobDeleteBefore(ctx, &riverdriver.JobDeleteBeforeParams{ + CancelledDoDelete: true, + CancelledFinalizedAtHorizon: horizon, + CompletedDoDelete: true, + CompletedFinalizedAtHorizon: horizon, + DiscardedDoDelete: false, + DiscardedFinalizedAtHorizon: horizon, + Max: 1_000, + }) + require.NoError(t, err) + require.Equal(t, 2, numDeleted) + + _, err = exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: cancelledJob.ID}) + require.ErrorIs(t, err, rivertype.ErrNotFound) + _, err = exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: completedJob.ID}) + require.ErrorIs(t, err, rivertype.ErrNotFound) + _, err = exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: discardedJob.ID}) + require.NoError(t, err) + }) + t.Run("QueuesExcluded", func(t *testing.T) { t.Parallel() diff --git a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql index 528de8d8..1cf46567 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql @@ -98,9 +98,9 @@ WHERE SELECT id FROM /* TEMPLATE: schema */river_job WHERE - (state = 'cancelled' AND finalized_at < cast(@cancelled_finalized_at_horizon AS text)) OR - (state = 'completed' AND finalized_at < cast(@completed_finalized_at_horizon AS text)) OR - (state = 'discarded' AND finalized_at < cast(@discarded_finalized_at_horizon AS text)) + (state = 'cancelled' AND cast(@cancelled_do_delete AS boolean) AND finalized_at < cast(@cancelled_finalized_at_horizon AS text)) OR + (state = 'completed' AND cast(@completed_do_delete AS boolean) AND finalized_at < cast(@completed_finalized_at_horizon AS text)) OR + (state = 'discarded' AND cast(@discarded_do_delete AS boolean) AND finalized_at < cast(@discarded_finalized_at_horizon AS text)) ORDER BY id LIMIT @max ) diff --git a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go index 1bccd481..4848272a 100644 --- a/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go +++ b/riverdriver/riversqlite/internal/dbsqlc/river_job.sql.go @@ -219,11 +219,11 @@ WHERE SELECT id FROM /* TEMPLATE: schema */river_job WHERE - (state = 'cancelled' AND finalized_at < cast(?1 AS text)) OR - (state = 'completed' AND finalized_at < cast(?2 AS text)) OR - (state = 'discarded' AND finalized_at < cast(?3 AS text)) + (state = 'cancelled' AND cast(?1 AS boolean) AND finalized_at < cast(?2 AS text)) OR + (state = 'completed' AND cast(?3 AS boolean) AND finalized_at < cast(?4 AS text)) OR + (state = 'discarded' AND cast(?5 AS boolean) AND finalized_at < cast(?6 AS text)) ORDER BY id - LIMIT ?4 + LIMIT ?7 ) -- This is really awful, but unless the ` + "`" + `sqlc.slice` + "`" + ` appears as the very -- last parameter in the query things will fail if it includes more than one @@ -239,14 +239,17 @@ WHERE -- charts buggy, and there's little interest from the maintainers in fixing -- any of it. We already started using it though, so plough on. AND ( - cast(?5 AS boolean) + cast(?8 AS boolean) OR river_job.queue NOT IN (/*SLICE:queues_excluded*/?) ) ` type JobDeleteBeforeParams struct { + CancelledDoDelete bool CancelledFinalizedAtHorizon string + CompletedDoDelete bool CompletedFinalizedAtHorizon string + DiscardedDoDelete bool DiscardedFinalizedAtHorizon string Max int64 QueuesExcludedEmpty bool @@ -256,8 +259,11 @@ type JobDeleteBeforeParams struct { func (q *Queries) JobDeleteBefore(ctx context.Context, db DBTX, arg *JobDeleteBeforeParams) (sql.Result, error) { query := jobDeleteBefore var queryParams []interface{} + queryParams = append(queryParams, arg.CancelledDoDelete) queryParams = append(queryParams, arg.CancelledFinalizedAtHorizon) + queryParams = append(queryParams, arg.CompletedDoDelete) queryParams = append(queryParams, arg.CompletedFinalizedAtHorizon) + queryParams = append(queryParams, arg.DiscardedDoDelete) queryParams = append(queryParams, arg.DiscardedFinalizedAtHorizon) queryParams = append(queryParams, arg.Max) queryParams = append(queryParams, arg.QueuesExcludedEmpty) diff --git a/riverdriver/riversqlite/river_sqlite_driver.go b/riverdriver/riversqlite/river_sqlite_driver.go index a6b0264d..28889608 100644 --- a/riverdriver/riversqlite/river_sqlite_driver.go +++ b/riverdriver/riversqlite/river_sqlite_driver.go @@ -437,8 +437,11 @@ func (e *Executor) JobDeleteBefore(ctx context.Context, params *riverdriver.JobD } res, err := dbsqlc.New().JobDeleteBefore(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.JobDeleteBeforeParams{ + CancelledDoDelete: params.CancelledDoDelete, CancelledFinalizedAtHorizon: timeString(params.CancelledFinalizedAtHorizon), + CompletedDoDelete: params.CompletedDoDelete, CompletedFinalizedAtHorizon: timeString(params.CompletedFinalizedAtHorizon), + DiscardedDoDelete: params.DiscardedDoDelete, DiscardedFinalizedAtHorizon: timeString(params.DiscardedFinalizedAtHorizon), Max: int64(params.Max), QueuesExcluded: params.QueuesExcluded,