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
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

⚠️ This release contains a new database migration, version 8, but it only affects SQLite:

- If you're on Postgres, you can ignore it with no adverse affects.
- If you're on SQLite, it rebuilds `river_job` to add an `AUTOINCREMENT` keyword to the primary key, preventing a possible edge case where generated job IDs could be reused after deletion. It's not necessary to run the migration for River to work, but it's a good idea to get it in when convenient. [PR #1390](https://github.com/riverqueue/river/pull/1390).

### Added

- Added `Config.FetchOnlyKnownKinds` to restrict job fetching to registered worker kinds, including aliases. Clients with different workers can share a queue while leaving unknown jobs available without consuming attempts. Disabled by default; leader election and stuck-job rescue behavior are unchanged. [PR #1396](https://github.com/riverqueue/river/pull/1396).
Expand Down Expand Up @@ -34,6 +39,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- 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).
- Fixed SQLite reusing the ID of a deleted job when that job held the largest ID, which could cause an ID observed earlier to refer to an unrelated job later. [PR #1390](https://github.com/riverqueue/river/pull/1390).
- Fixed the `Job appears to be stuck` log line reporting the client-level `JobTimeout` instead of the worker-level timeout when a worker overrides `Timeout`. [PR #1394](https://github.com/riverqueue/river/pull/1394).

## [0.47.0] - 2026-09-01
Expand Down
2 changes: 1 addition & 1 deletion riverdriver/river_driver_interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -993,7 +993,7 @@ func MigrationLineMainTruncateTables(version int) []string {
return []string{"river_job", "river_leader", "river_queue"}
case 5, 6:
return []string{"river_job", "river_leader", "river_queue", "river_client", "river_client_queue"}
case 0, 7:
case 0, 7, 8:
return []string{"river_job", "river_leader", "river_queue", "river_notification"}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
-- No-op. PostgreSQL sequences already prevent automatically generated job IDs
-- from being reused.
SELECT 1;
Comment on lines +2 to +3

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we include #1012 and #875 in this release so these aren't just no-ops?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think we actually got all the SQL cleanup back then and with nothing new coming in since.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I should add to that I think this is kind of better because we can make the claim that if you're on Postgres, 008 is a pure no-op, like you don't need to do anything at all.

Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
-- No-op. PostgreSQL sequences already prevent automatically generated job IDs
-- from being reused.
SELECT 1;
14 changes: 14 additions & 0 deletions riverdriver/riverdrivertest/job_insert.go
Original file line number Diff line number Diff line change
Expand Up @@ -664,6 +664,20 @@ func exerciseJobInsert[TTx any](ctx context.Context, t *testing.T,
t.Run("JobInsertFull", func(t *testing.T) {
t.Parallel()

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

exec, _ := setup(ctx, t)

job := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{})

_, err := exec.JobDelete(ctx, &riverdriver.JobDeleteParams{ID: job.ID})
require.NoError(t, err)

jobAfter := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{})
require.Greater(t, jobAfter.ID, job.ID)
})

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

Expand Down
100 changes: 99 additions & 1 deletion riverdriver/riverdrivertest/migration.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,8 @@ func exerciseMigration[TTx any](ctx context.Context, t *testing.T,
driver.GetMigrationTruncateTables(riverdriver.MigrationLineMain, 6))
require.Equal(t, expectedLatestTables,
driver.GetMigrationTruncateTables(riverdriver.MigrationLineMain, 7))
require.Equal(t, expectedLatestTables,
driver.GetMigrationTruncateTables(riverdriver.MigrationLineMain, 8))
require.Equal(t, expectedLatestTables,
driver.GetMigrationTruncateTables(riverdriver.MigrationLineMain, 0))
})
Expand Down Expand Up @@ -144,7 +146,7 @@ func exerciseMigration[TTx any](ctx context.Context, t *testing.T,
}
})

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

driver, schema := driverWithSchema(ctx, t, &riverdbtest.TestSchemaOpts{
Expand All @@ -171,6 +173,39 @@ func exerciseMigration[TTx any](ctx context.Context, t *testing.T,
require.NotZero(t, job.ID)
})

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

driver, schema := driverWithSchema(ctx, t, &riverdbtest.TestSchemaOpts{
DisableReuse: true,
LineTargetVersions: map[string]int{
riverdriver.MigrationLineMain: 7,
},
})
exec := driver.GetExecutor()

job := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{Schema: schema})

migrator, err := rivermigrate.New(driver, &rivermigrate.Config{
Line: riverdriver.MigrationLineMain,
Logger: riversharedtest.Logger(t),
Schema: schema,
})
require.NoError(t, err)

_, err = migrator.Migrate(ctx, rivermigrate.DirectionUp, nil)
require.NoError(t, err)

job, err = exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID, Schema: schema})
require.NoError(t, err)

_, err = exec.JobDelete(ctx, &riverdriver.JobDeleteParams{ID: job.ID, Schema: schema})
require.NoError(t, err)

jobAfter := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{Schema: schema})
require.Greater(t, jobAfter.ID, job.ID)
})

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

Expand Down Expand Up @@ -215,6 +250,69 @@ func exerciseMigration[TTx any](ctx context.Context, t *testing.T,
require.NotZero(t, queue.UpdatedAt)
})

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

for _, testCase := range []struct {
name string
sql string
}{
{"LegacyWorkflow", `CREATE INDEX river_job_workflow_scheduling ON river_job (state)`},
{"Pro", `CREATE TABLE river_job_sequence (id integer PRIMARY KEY, key text)`},
{"WorkflowV2", `CREATE TABLE river_workflow (id text PRIMARY KEY)`},
} {
t.Run(testCase.name, func(t *testing.T) {
t.Parallel()

for _, direction := range []rivermigrate.Direction{rivermigrate.DirectionDown, rivermigrate.DirectionUp} {
t.Run(string(direction), func(t *testing.T) {
t.Parallel()

version := 7
if direction == rivermigrate.DirectionDown {
version = 8
}
driver, schema := driverWithSchema(ctx, t, &riverdbtest.TestSchemaOpts{
DisableReuse: true,
LineTargetVersions: map[string]int{
riverdriver.MigrationLineMain: version,
},
Lines: []string{riverdriver.MigrationLineMain},
})
if driver.DatabaseName() != riverdriver.DatabaseNameSQLite {
t.Skip("SQLite table rebuild")
}
exec := driver.GetExecutor()
job := testfactory.Job(ctx, t, exec, &testfactory.JobOpts{Schema: schema})

// Deliberately omit Pro migration records, as when applying SQL
// through an external migration tool.
require.NoError(t, exec.Exec(ctx, testCase.sql))
require.NoError(t, exec.Exec(ctx, `ALTER TABLE river_job ADD COLUMN partition_key text`))
migrator, err := rivermigrate.New(driver, &rivermigrate.Config{Logger: riversharedtest.Logger(t), Schema: schema})
require.NoError(t, err)

_, err = migrator.Migrate(ctx, direction, &rivermigrate.MigrateOpts{MaxSteps: 1})
require.ErrorContains(t, err, "River SQLite migration 008 cannot run while River Pro schema is installed")

jobAfter, err := exec.JobGetByID(ctx, &riverdriver.JobGetByIDParams{ID: job.ID, Schema: schema})
require.NoError(t, err)
require.Equal(t, job, jobAfter)
exists, err := exec.ColumnExists(ctx, &riverdriver.ColumnExistsParams{Column: "partition_key", Schema: schema, Table: "river_job"})
require.NoError(t, err)
require.True(t, exists)
exists, err = exec.IndexExists(ctx, &riverdriver.IndexExistsParams{Index: "river_job_kind", Schema: schema})
require.NoError(t, err)
require.True(t, exists)
migrations, err := exec.MigrationGetByLine(ctx, &riverdriver.MigrationGetByLineParams{Line: riverdriver.MigrationLineMain, Schema: schema})
require.NoError(t, err)
require.Len(t, migrations, version)
})
}
})
}
})

type testBundle struct {
driver riverdriver.Driver[TTx]
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
-- No-op. PostgreSQL sequences already prevent automatically generated job IDs
-- from being reused.
SELECT 1;
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
-- No-op. PostgreSQL sequences already prevent automatically generated job IDs
-- from being reused.
SELECT 1;
3 changes: 2 additions & 1 deletion riverdriver/riversqlite/internal/dbsqlc/river_job.sql
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
CREATE TABLE river_job (
id integer PRIMARY KEY, -- SQLite makes this autoincrementing automatically
-- AUTOINCREMENT prevents SQLite from reusing the IDs of deleted jobs.
id integer PRIMARY KEY AUTOINCREMENT,
args jsonb NOT NULL DEFAULT (jsonb('{}')),
attempt integer NOT NULL DEFAULT 0,
attempted_at timestamp,
Expand Down
4 changes: 2 additions & 2 deletions riverdriver/riversqlite/migration/main/006_bulk_unique.up.sql
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
DROP TABLE /* TEMPLATE: schema */river_job;

CREATE TABLE /* TEMPLATE: schema */river_job (
id integer PRIMARY KEY, -- SQLite makes this autoincrementing automatically
id integer PRIMARY KEY, -- SQLite aliases this to ROWID, which may reuse deleted IDs.
args blob NOT NULL DEFAULT '{}',
attempt integer NOT NULL DEFAULT 0,
attempted_at timestamp,
Expand Down Expand Up @@ -60,4 +60,4 @@ CREATE UNIQUE INDEX /* TEMPLATE: schema */river_job_unique_idx ON river_job (uni
WHEN 'running' THEN unique_states & (1 << 6)
WHEN 'scheduled' THEN unique_states & (1 << 7)
ELSE 0
END >= 1;
END >= 1;
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ DROP INDEX /* TEMPLATE: schema */river_job_unique_idx;
ALTER TABLE /* TEMPLATE: schema */river_job RENAME TO river_job_old;

CREATE TABLE /* TEMPLATE: schema */river_job (
id integer PRIMARY KEY, -- SQLite makes this autoincrementing automatically
id integer PRIMARY KEY, -- SQLite aliases this to ROWID, which may reuse deleted IDs.
args blob NOT NULL DEFAULT '{}',
attempt integer NOT NULL DEFAULT 0,
attempted_at timestamp,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ DROP INDEX /* TEMPLATE: schema */river_job_unique_idx;
ALTER TABLE /* TEMPLATE: schema */river_job RENAME TO river_job_old;

CREATE TABLE /* TEMPLATE: schema */river_job (
id integer PRIMARY KEY, -- SQLite makes this autoincrementing automatically
id integer PRIMARY KEY, -- SQLite aliases this to ROWID, which may reuse deleted IDs.
args blob NOT NULL DEFAULT (jsonb('{}')),
attempt integer NOT NULL DEFAULT 0,
attempted_at timestamp,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
-- Rebuild river_job to restore SQLite's default ROWID allocation behavior.

-- Rebuilding river_job would discard schema installed by River Pro. Check
-- schema objects instead of migration records to also catch manually applied
-- Pro migrations and the legacy workflow migration line.
CREATE TEMP TABLE river_job_pro_schema_guard (
id integer NOT NULL
);

CREATE TEMP TRIGGER river_job_pro_schema_guard_enforce
BEFORE INSERT ON river_job_pro_schema_guard
WHEN EXISTS (
SELECT 1
FROM /* TEMPLATE: schema */sqlite_master
WHERE name IN ('river_job_sequence', 'river_job_workflow_scheduling', 'river_workflow')
)
BEGIN
SELECT RAISE(ABORT, 'River SQLite migration 008 cannot run while River Pro schema is installed');
END;

INSERT INTO river_job_pro_schema_guard (id) VALUES (1);

DROP TRIGGER river_job_pro_schema_guard_enforce;
DROP TABLE river_job_pro_schema_guard;

DROP INDEX /* TEMPLATE: schema */river_job_kind;
DROP INDEX /* TEMPLATE: schema */river_job_state_and_finalized_at_index;
DROP INDEX /* TEMPLATE: schema */river_job_prioritized_fetching_index;
DROP INDEX /* TEMPLATE: schema */river_job_unique_idx;

ALTER TABLE /* TEMPLATE: schema */river_job RENAME TO river_job_old;

CREATE TABLE /* TEMPLATE: schema */river_job (
id integer PRIMARY KEY,
args blob NOT NULL DEFAULT (jsonb('{}')),
attempt integer NOT NULL DEFAULT 0,
attempted_at timestamp,
attempted_by blob, -- json
created_at timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP,
errors blob, -- json
finalized_at timestamp,
kind text NOT NULL,
max_attempts integer NOT NULL DEFAULT 25,
metadata blob NOT NULL DEFAULT (jsonb('{}')),
priority integer NOT NULL DEFAULT 1,
queue text NOT NULL DEFAULT 'default',
state text NOT NULL DEFAULT 'available',
scheduled_at timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP,
tags blob NOT NULL DEFAULT (jsonb('[]')),
unique_key blob,
unique_states integer,
CONSTRAINT finalized_or_finalized_at_null CHECK (
(finalized_at IS NULL AND state NOT IN ('cancelled', 'completed', 'discarded')) OR
(finalized_at IS NOT NULL AND state IN ('cancelled', 'completed', 'discarded'))
),
CONSTRAINT priority_in_range CHECK (priority >= 1 AND priority <= 4),
CONSTRAINT queue_length CHECK (length(queue) > 0 AND length(queue) < 128),
CONSTRAINT kind_length CHECK (length(kind) > 0 AND length(kind) < 128),
CONSTRAINT state_valid CHECK (state IN ('available', 'cancelled', 'completed', 'discarded', 'pending', 'retryable', 'running', 'scheduled'))
);

INSERT INTO /* TEMPLATE: schema */river_job (
id,
args,
attempt,
attempted_at,
attempted_by,
created_at,
errors,
finalized_at,
kind,
max_attempts,
metadata,
priority,
queue,
state,
scheduled_at,
tags,
unique_key,
unique_states
)
SELECT
id,
args,
attempt,
attempted_at,
attempted_by,
created_at,
errors,
finalized_at,
kind,
max_attempts,
metadata,
priority,
queue,
state,
scheduled_at,
tags,
unique_key,
unique_states
FROM /* TEMPLATE: schema */river_job_old;

DROP TABLE /* TEMPLATE: schema */river_job_old;

CREATE INDEX /* TEMPLATE: schema */river_job_kind ON river_job (kind);
CREATE INDEX /* TEMPLATE: schema */river_job_state_and_finalized_at_index ON river_job (state, finalized_at) WHERE finalized_at IS NOT NULL;
CREATE INDEX /* TEMPLATE: schema */river_job_prioritized_fetching_index ON river_job (state, queue, priority, scheduled_at, id);
CREATE UNIQUE INDEX /* TEMPLATE: schema */river_job_unique_idx ON river_job (unique_key)
WHERE unique_key IS NOT NULL
AND unique_states IS NOT NULL
AND CASE state
WHEN 'available' THEN unique_states & (1 << 0)
WHEN 'cancelled' THEN unique_states & (1 << 1)
WHEN 'completed' THEN unique_states & (1 << 2)
WHEN 'discarded' THEN unique_states & (1 << 3)
WHEN 'pending' THEN unique_states & (1 << 4)
WHEN 'retryable' THEN unique_states & (1 << 5)
WHEN 'running' THEN unique_states & (1 << 6)
WHEN 'scheduled' THEN unique_states & (1 << 7)
ELSE 0
END >= 1;
Loading
Loading