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,