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

- 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).

Expand Down
33 changes: 33 additions & 0 deletions riverdriver/riverdrivertest/job_delete.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
6 changes: 3 additions & 3 deletions riverdriver/riversqlite/internal/dbsqlc/river_job.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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
)
Expand Down
16 changes: 11 additions & 5 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.

3 changes: 3 additions & 0 deletions riverdriver/riversqlite/river_sqlite_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading