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 @@ -16,6 +16,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- 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).
- 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).
- Fixed SQLite notification listeners delivering notifications from before a subscription or from an unsubscribe gap. Notification reads now fetch subscribed topics in bounded batches, and cleanup deletes expired notifications in batches of 10,000 rows (reduced to 1,000 after repeated timeouts), with pauses between batches to reduce write lock contention. [PR #1381](https://github.com/riverqueue/river/pull/1381).

## [0.47.0] - 2026-09-01

Expand Down
68 changes: 58 additions & 10 deletions internal/maintenance/sqlite_notification_cleaner.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,12 @@ import (

"github.com/riverqueue/river/riverdriver"
"github.com/riverqueue/river/rivershared/baseservice"
"github.com/riverqueue/river/rivershared/circuitbreaker"
"github.com/riverqueue/river/rivershared/riversharedmaintenance"
"github.com/riverqueue/river/rivershared/startstop"
"github.com/riverqueue/river/rivershared/testsignal"
"github.com/riverqueue/river/rivershared/util/randutil"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivershared/util/timeutil"
Expand All @@ -24,14 +27,16 @@ const (

// SQLiteNotificationCleanerTestSignals are internal signals used exclusively in tests.
type SQLiteNotificationCleanerTestSignals struct {
DeletedBatch testsignal.TestSignal[struct{}] // notifies when runOnce finishes a pass
DeletedBatch testsignal.TestSignal[struct{}] // notifies when a delete batch finishes
}

func (ts *SQLiteNotificationCleanerTestSignals) Init(tb testutil.TestingTB) {
ts.DeletedBatch.Init(tb)
}

type SQLiteNotificationCleanerConfig struct {
riversharedmaintenance.BatchSizes

// Interval is the amount of time to wait between cleaner runs.
Interval time.Duration

Expand All @@ -47,6 +52,8 @@ type SQLiteNotificationCleanerConfig struct {
}

func (c *SQLiteNotificationCleanerConfig) mustValidate() *SQLiteNotificationCleanerConfig {
c.MustValidate()

if c.Interval <= 0 {
panic("SQLiteNotificationCleanerConfig.Interval must be above zero")
}
Expand All @@ -72,18 +79,25 @@ type SQLiteNotificationCleaner struct {
TestSignals SQLiteNotificationCleanerTestSignals

exec riverdriver.Executor

// After repeated timeouts, keep using smaller batches until restart.
reducedBatchSizeBreaker *circuitbreaker.CircuitBreaker
}

// NewSQLiteNotificationCleaner returns a SQLite notification cleaner.
func NewSQLiteNotificationCleaner(archetype *baseservice.Archetype, config *SQLiteNotificationCleanerConfig, exec riverdriver.Executor) *SQLiteNotificationCleaner {
batchSizes := config.WithDefaults()

return baseservice.Init(archetype, &SQLiteNotificationCleaner{
Config: (&SQLiteNotificationCleanerConfig{
BatchSizes: batchSizes,
Interval: cmp.Or(config.Interval, SQLiteNotificationCleanerIntervalDefault),
RetentionPeriod: cmp.Or(config.RetentionPeriod, SQLiteNotificationCleanerRetentionPeriodDefault),
Schema: config.Schema,
Timeout: cmp.Or(config.Timeout, riversharedmaintenance.TimeoutDefault),
}).mustValidate(),
exec: exec,
exec: exec,
reducedBatchSizeBreaker: riversharedmaintenance.ReducedBatchSizeBreaker(batchSizes),
})
}

Expand Down Expand Up @@ -129,24 +143,58 @@ func (s *SQLiteNotificationCleaner) Start(ctx context.Context) error { //nolint:
return nil
}

func (s *SQLiteNotificationCleaner) batchSize() int {
if s.reducedBatchSizeBreaker.Open() {
return s.Config.Reduced
}
return s.Config.Default
}

type sqliteNotificationCleanerRunOnceResult struct {
NumNotificationsDeleted int
}

func (s *SQLiteNotificationCleaner) runOnce(ctx context.Context) (*sqliteNotificationCleanerRunOnceResult, error) {
return timeoututil.WithTimeoutV(ctx, s.Config.Timeout, s.Name+".runOnce", func(ctx context.Context) (*sqliteNotificationCleanerRunOnceResult, error) {
numDeleted, err := s.exec.NotificationDeleteBefore(ctx, &riverdriver.NotificationDeleteBeforeParams{
CreatedAtHorizon: time.Now().Add(-s.Config.RetentionPeriod),
Schema: s.Config.Schema,
res := &sqliteNotificationCleanerRunOnceResult{}
// Keep a fixed horizon so new expirations don't extend a cleanup pass.
createdAtHorizon := time.Now().Add(-s.Config.RetentionPeriod)

for {
if err := ctx.Err(); err != nil {
return nil, err
}

numDeleted, err := timeoututil.WithTimeoutV(ctx, s.Config.Timeout, s.Name+".runOnce", func(ctx context.Context) (int, error) {
numDeleted, err := s.exec.NotificationDeleteBefore(ctx, &riverdriver.NotificationDeleteBeforeParams{
CreatedAtHorizon: createdAtHorizon,
Max: s.batchSize(),
Schema: s.Config.Schema,
})
if err != nil {
return 0, err
}

s.reducedBatchSizeBreaker.ResetIfNotOpen()

return numDeleted, nil
})
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
s.reducedBatchSizeBreaker.Trip()
}

return nil, err
}

s.TestSignals.DeletedBatch.Signal(struct{}{})
res.NumNotificationsDeleted += numDeleted

return &sqliteNotificationCleanerRunOnceResult{
NumNotificationsDeleted: numDeleted,
}, nil
})
if numDeleted < s.batchSize() {
return res, nil
}

// Each delete commits independently. Yield SQLite's writer lock before
// the next batch so job inserts and updates can make progress.
serviceutil.CancellableSleep(ctx, randutil.DurationBetween(riversharedmaintenance.BatchBackoffMin, riversharedmaintenance.BatchBackoffMax))
}
}
152 changes: 152 additions & 0 deletions internal/maintenance/sqlite_notification_cleaner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package maintenance

import (
"context"
"errors"
"testing"
"time"

Expand All @@ -10,10 +11,21 @@ import (
"github.com/riverqueue/river/riverdbtest"
"github.com/riverqueue/river/riverdriver"
"github.com/riverqueue/river/riverdriver/riverpgxv5"
"github.com/riverqueue/river/rivershared/riversharedmaintenance"
"github.com/riverqueue/river/rivershared/riversharedtest"
"github.com/riverqueue/river/rivershared/startstoptest"
)

type sqliteNotificationCleanerExecutor struct {
riverdriver.Executor

notificationDeleteBeforeFunc func(context.Context, *riverdriver.NotificationDeleteBeforeParams) (int, error)
}

func (e *sqliteNotificationCleanerExecutor) NotificationDeleteBefore(ctx context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error) {
return e.notificationDeleteBeforeFunc(ctx, params)
}

func TestSQLiteNotificationCleaner(t *testing.T) {
t.Parallel()

Expand Down Expand Up @@ -59,6 +71,36 @@ func TestSQLiteNotificationCleaner(t *testing.T) {
return count
}

t.Run("CancelsBetweenBatches", func(t *testing.T) {
t.Parallel()

cleaner, bundle := setup(t)
cleaner.Config.Default = 2
cleaner.TestSignals.Init(t)

require.NoError(t, bundle.exec.Exec(ctx, `
INSERT INTO river_notification (created_at, payload, topic)
SELECT $1, 'old_payload', 'topic' FROM generate_series(1, 5)
`, time.Now().Add(-2*time.Hour)))

ctx, cancelFunc := context.WithCancel(ctx)
defer cancelFunc()
cleaner.exec = &sqliteNotificationCleanerExecutor{
Executor: bundle.exec,
notificationDeleteBeforeFunc: func(ctx context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error) {
numDeleted, err := bundle.exec.NotificationDeleteBefore(ctx, params)
cancelFunc()
return numDeleted, err
},
}

_, err := cleaner.runOnce(ctx)
require.ErrorIs(t, err, context.Canceled)
cleaner.TestSignals.DeletedBatch.WaitOrTimeout()
cleaner.TestSignals.DeletedBatch.RequireEmpty()
require.Equal(t, 3, notificationCount(t, bundle.exec))
})

t.Run("Defaults", func(t *testing.T) {
t.Parallel()

Expand All @@ -68,6 +110,8 @@ func TestSQLiteNotificationCleaner(t *testing.T) {
nil,
)

require.Equal(t, riversharedmaintenance.BatchSizeDefault, cleaner.Config.Default)
require.Equal(t, riversharedmaintenance.BatchSizeReduced, cleaner.Config.Reduced)
require.Equal(t, SQLiteNotificationCleanerIntervalDefault, cleaner.Config.Interval)
require.Equal(t, SQLiteNotificationCleanerRetentionPeriodDefault, cleaner.Config.RetentionPeriod)
})
Expand All @@ -94,6 +138,114 @@ func TestSQLiteNotificationCleaner(t *testing.T) {
require.Equal(t, 1, notificationCount(t, bundle.exec))
})

t.Run("DeletesMultipleBatches", func(t *testing.T) {
t.Parallel()

cleaner, bundle := setup(t)
cleaner.Config.Default = 2
cleaner.TestSignals.Init(t)

require.NoError(t, bundle.exec.Exec(ctx, `
INSERT INTO river_notification (created_at, payload, topic)
SELECT $1, 'old_payload', 'topic' FROM generate_series(1, 5)
`, time.Now().Add(-2*time.Hour)))
require.NoError(t, bundle.exec.Exec(ctx, `
INSERT INTO river_notification (payload, topic) VALUES ('new_payload', 'topic')
`))

res, err := cleaner.runOnce(ctx)
require.NoError(t, err)
require.Equal(t, 5, res.NumNotificationsDeleted)
for range 3 { // Two full batches followed by a partial batch.
cleaner.TestSignals.DeletedBatch.WaitOrTimeout()
}
cleaner.TestSignals.DeletedBatch.RequireEmpty()
require.Equal(t, 1, notificationCount(t, bundle.exec))
})

t.Run("ReducedBatchSizeBreakerIgnoresOtherErrors", func(t *testing.T) {
t.Parallel()

cleaner, bundle := setup(t)

for _, queryErr := range []error{context.Canceled, errors.New("notification delete failed")} {
cleaner.exec = &sqliteNotificationCleanerExecutor{
Executor: bundle.exec,
notificationDeleteBeforeFunc: func(context.Context, *riverdriver.NotificationDeleteBeforeParams) (int, error) {
return 0, queryErr
},
}

for range cleaner.reducedBatchSizeBreaker.Limit() {
_, err := cleaner.runOnce(ctx)
require.ErrorIs(t, err, queryErr)
}
require.Equal(t, riversharedmaintenance.BatchSizeDefault, cleaner.batchSize())
}
})

t.Run("ReducedBatchSizeBreakerResetsOnSuccess", func(t *testing.T) {
t.Parallel()

cleaner, bundle := setup(t)
var queryErr error
cleaner.exec = &sqliteNotificationCleanerExecutor{
Executor: bundle.exec,
notificationDeleteBeforeFunc: func(_ context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error) {
require.Equal(t, riversharedmaintenance.BatchSizeDefault, params.Max)
return 0, queryErr
},
}

for range 2 {
queryErr = context.DeadlineExceeded
for range cleaner.reducedBatchSizeBreaker.Limit() - 1 {
_, err := cleaner.runOnce(ctx)
require.ErrorIs(t, err, context.DeadlineExceeded)
require.Equal(t, riversharedmaintenance.BatchSizeDefault, cleaner.batchSize())
}

queryErr = nil
_, err := cleaner.runOnce(ctx)
require.NoError(t, err)
require.Equal(t, riversharedmaintenance.BatchSizeDefault, cleaner.batchSize())
}
})

t.Run("ReducedBatchSizeBreakerTrips", func(t *testing.T) {
t.Parallel()

cleaner, bundle := setup(t)
expectedMax := riversharedmaintenance.BatchSizeDefault
queryErr := context.DeadlineExceeded
cleaner.exec = &sqliteNotificationCleanerExecutor{
Executor: bundle.exec,
notificationDeleteBeforeFunc: func(_ context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error) {
require.Equal(t, expectedMax, params.Max)
return 0, queryErr
},
}

for range cleaner.reducedBatchSizeBreaker.Limit() - 1 {
_, err := cleaner.runOnce(ctx)
require.ErrorIs(t, err, context.DeadlineExceeded)
require.Equal(t, riversharedmaintenance.BatchSizeDefault, cleaner.batchSize())
}

_, err := cleaner.runOnce(ctx)
require.ErrorIs(t, err, context.DeadlineExceeded)
require.Equal(t, riversharedmaintenance.BatchSizeReduced, cleaner.batchSize())

// Once tripped, successful deletes keep using the reduced batch size.
expectedMax = riversharedmaintenance.BatchSizeReduced
queryErr = nil
for range 2 {
_, err := cleaner.runOnce(ctx)
require.NoError(t, err)
require.Equal(t, riversharedmaintenance.BatchSizeReduced, cleaner.batchSize())
}
})

t.Run("StartStopStress", func(t *testing.T) {
t.Parallel()

Expand Down
5 changes: 3 additions & 2 deletions riverdriver/river_driver_interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -280,8 +280,8 @@ type Executor interface {
// the `line` column was added to the migrations table.
MigrationInsertManyAssumingMain(ctx context.Context, params *MigrationInsertManyAssumingMainParams) ([]*Migration, error)

// NotificationDeleteBefore deletes notifications before a certain time
// horizon.
// NotificationDeleteBefore deletes up to Max notifications before a certain
// time horizon, oldest first.
//
// A "notification" in this context refers to a row in `river_notification`
// which is a special table implemented in some databases (e.g. SQLite) that
Expand Down Expand Up @@ -845,6 +845,7 @@ type NotifyManyParams struct {

type NotificationDeleteBeforeParams struct {
CreatedAtHorizon time.Time
Max int
Schema string
}

Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion riverdriver/riverdatabasesql/river_database_sql_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -921,7 +921,10 @@ func (e *Executor) MigrationInsertManyAssumingMain(ctx context.Context, params *
}

func (e *Executor) NotificationDeleteBefore(ctx context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error) {
numDeleted, err := dbsqlc.New().NotificationDeleteBefore(schemaTemplateParam(ctx, params.Schema), e.dbtx, params.CreatedAtHorizon)
numDeleted, err := dbsqlc.New().NotificationDeleteBefore(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.NotificationDeleteBeforeParams{
CreatedAtHorizon: params.CreatedAtHorizon,
Max: int64(params.Max),
})
return int(numDeleted), interpretError(err)
}

Expand Down
Loading
Loading