Skip to content

Add River for Rust - #1436

Closed
bgentry wants to merge 33 commits into
bg/conformance-harnessfrom
bg/rust-port
Closed

bgentry wants to merge 33 commits into
bg/conformance-harnessfrom
bg/rust-port

Conversation

@bgentry

@bgentry bgentry commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

Teams running River in Go sometimes have services in Rust that need to enqueue or work the same jobs. Right now their choices are to write raw SQL against River's tables, which is easy to get subtly wrong (unique keys, attempt errors, notifications), or to put a Go service in front of the queue. This PR adds a Rust implementation of River that shares the schema and job protocol with River Go on PostgreSQL and SQLite. A Rust service can insert jobs that Go workers run, work jobs Go inserts, and take part in leader election and maintenance alongside Go clients in the same database.

This PR is stacked on the conformance suite PR (#1435) and should be reviewed after it. The cross-language harness from that PR is how this implementation proves it's interoperable: the Rust crate added here is its default candidate.

All the code lives in a new rust/ workspace. The crates are versioned and released together, starting at 0.48.0-alpha.1 to match River Go 0.48.

  • riverqueue is the client. It's built on Tokio and takes a caller-owned SQLx pool, without generics over the database or a driver trait. It has typed async workers that get a CancellationToken, and request builders for insertion (transactional, batch, unique, scheduled, pending), job and queue management, and transactional completion. The runtime fetches, completes in bounded concurrent batches, retries, snoozes, cancels, and rescues jobs the way Go does. It also runs leader election and the maintenance services: rescuer, cleaners, scheduler, reindexer, and periodic jobs with Go-compatible cron parsing. Hooks, middleware, error handlers, and retry policies follow Go's extension points.
  • riverqueue-macros provides #[derive(JobArgs)], with compile-time errors for invalid kinds and unique options.
  • riverqueue-migrate applies River's PostgreSQL and SQLite migration lines and shares the river_migration history with Go, so either language can migrate a database the other uses. Its migrations are mirrors of Go's files. syncrustmigrations generates them, and CI fails if they drift, so Rust never gets its own schema history.
  • riverqueue-test has insertion assertions and work-once helpers, similar to rivertest.
  • riverqueue-cli installs a riverqueue binary with migrate-* commands and bench.
  • riverqueue-conformance (unpublished) serves the conformance adapter contract (adapter_version 22) with riverqueue. conformance/adapter/candidates/rust.json makes it the harness's default candidate and peer.

A few behaviors are deliberate and worth checking during review:

  • Requests given the caller's transaction run directly in it and open no savepoint, which matches River Go after Avoid opening subtransactions when already in a transaction #1420. If a request fails after River has written to the transaction, for example because insert middleware throws, that write stays in the transaction for the caller to roll back. On PostgreSQL, a failed statement aborts the transaction. caller_transactions.rs covers both. There's no Rust equivalent of RIVER_USE_LEGACY_SUBTRANSACTIONS; callers who need partial rollback can open their own savepoint.
  • Values that other engines read must be stored exactly. Unique key hash inputs are encoded the way Go's encoding/json encodes them, so 1.0 hashes like Go's 1, and unique columns, timestamps, and attempt errors match Go's format. Other JSON, such as args, metadata, and notification payloads, only needs to mean the same thing as Go's, so the crate has no raw-text APIs whose only purpose is to reproduce Go's exact bytes. Unit tests check unique keys and cron schedules against the Go-generated conformance fixtures, and the generator copies them into the crate because published tests can't read files outside it.
  • YugabyteDB is detected per database from the server's version string and yb_enable_listen_notify. Without xmax, unique inserts use a river:unique_nonce. Without LISTEN/NOTIFY, clients send no notifications and poll for cancellations, like Go.

On CI, the Rust targets stay out of make lint and make test, so Go contributors don't need a Rust toolchain. A separate Rust workflow runs lints (including PostgreSQL-only and SQLite-only builds), docs, tests on each supported Rust version, cargo deny, packaging, and semver checks. It also runs SQLite conformance, and for PostgreSQL 14 through 18 it runs the PostgreSQL tests, mixed and insert-only conformance, and a ten-minute soak. A manual release-candidate workflow gates on performance and a one-hour soak, and a weekly workflow runs a six-hour soak.

Here's how to review it:

Commits Notes
add the Rust migration crate, add the Rust JobArgs derive macro Small and self-contained. A good warm-up.
add River's Rust client This is the main commit, about 35k lines. It can't be split into smaller commits that each build, because the client, storage, pilot, and maintenance modules all depend on ClientInner and on each other. I'd suggest reading it in this order: lib.rs and README.md, then client/ (builder, insert, producer, completer), then database*/storage*, then maintenance/, then unique.rs and encoding*. docs/mixed-deployments.md explains what has to match between Go and Rust.
test the JobArgs derive… through test Rust extension points Integration tests grouped by layer: macro, storage and parity with Go's rows, runtime, database faults and YugabyteDB, extension points and caller transactions.
add Rust examples, add Rust test helpers, add the riverqueue command Separate crates and examples.
run the Rust client as a conformance candidate The Rust adapter, about 4k lines, mostly a mapping from contract methods to client calls.
run the Rust crates and conformance in CI, keep Rust and workflow dependencies updated, document the Rust workspace CI, Dependabot, the workspace README, and the changelog.

bgentry added 29 commits October 5, 2026 11:35
River's database protocol is shared by any implementation that reads or
writes River's tables, but nothing describes it outside the Go code.
Another implementation has no way to check that it inserts, claims,
completes, and notifies the way Go does.

Add a language-neutral description of that protocol under
`conformance/`. `manifest.json` declares matched implementation
versions, the migration line, and the protocol capabilities, with a
recorded decision for any capability that isn't complete. The migration
inventories record the canonical PostgreSQL and SQLite migration files
with their hashes.

`adapter/contract.json` specifies a versioned JSON-RPC 2.0 process
adapter: every method, its parameters and result schema, and the error
codes, plus normalized job and queue shapes. Profiles name the subsets an
adapter may serve (`postgres-full-v1`, `portable-storage-v1`,
`sqlite-runtime-v1`, and `insert-only-v1`), and a candidate schema
describes how a harness builds and starts an implementation's adapter.
The adapter README explains the process, connection, and environment
rules that the schemas can't express.
Add `riverconformanceadapter`, which exposes River Go through the
conformance adapter contract so a harness can drive Go and another
implementation against the same database and compare what each one
writes and observes.

The adapter serves every method in the PostgreSQL and SQLite profiles:
migrations, insertion including unique and transactional batches, job
and queue CRUD, list filters and cursors, worker runs with configurable
outcomes, barriers, extensions, periodic and resumable jobs, leadership,
maintenance, subscriptions, and fault injection. Results use the
contract's normalized job and queue shapes, and failures map to its error
codes.

It's its own module in the workspace, like River's driver modules, so
its SQLite driver and other dependencies stay out of River's main
module.
Add a Go test harness, behind the `riverconformance` build tag, that
starts the Go reference adapter and a candidate adapter against one
disposable PostgreSQL database and checks that each implementation can
read and work what the other writes.

The harness manages adapter processes with bounded request and exit
waits, validates every request and response against the adapter
contract, and loads the candidate from a descriptor in
`RIVER_CONFORMANCE_CANDIDATE_FILE` or `RIVER_CONFORMANCE_CANDIDATE`, so
nothing in it names a candidate language. A registry binds each scenario
to the one test that owns it; each scenario runs as its own subtest and
is credited only by its own assertions, and the owner fails unless every
scenario it owns ran. `TestCompatibilityArtifacts` checks the manifest,
contract, profiles, migration inventories, and scenario catalog against
each other and their schemas. With `RIVER_CONFORMANCE_REQUIRED=1`, a
missing database URL or candidate fails instead of skipping.

The first scenarios cover handshakes, migrations in both directions and
in custom schemas, insertion and work across engines, uniqueness, batch
and transactional insertion, job and queue CRUD, list filters and
cursors, row round trips, 64-bit job IDs, transaction visibility,
cancel and retry races, worker outcomes, completion batching, and
cross-engine rescue after a killed process.

`make lint` now also lints the harness, and `make test/conformance` runs
it.
Add mixed scenarios for what happens while and after a job runs:
barriers, panics and their attempt traces, transactional completion,
snoozes, terminal completion racing an external update, hook and
middleware order, resumable jobs, dynamic queue reconfiguration,
periodic jobs, error handler cancellation, refetched attempts,
timeouts, claim order, scheduler unique conflicts, retries of exhausted
jobs, clock boundaries, stuck job detection, and completion under pool
pressure.

Each runs in both directions where it can, with Go inserting or
observing and the candidate working, and the reverse.
Add scenarios for fleets whose clients don't all register the same
workers: a worker renamed with a kind alias, clients that fetch only the
kinds they know, the rescuer discarding a job with an unknown kind, and
the error recorded when a client works a kind it has no worker for.
Add scenarios in which one engine's notification has to reach the other:
insert wakeups, including only after a transaction commits, queue pause
and resume, queue subscription events, and remote cancellation, whether
the worker listens, polls, cooperates with the cancellation, or receives
it while claiming the job. Notification payloads are compared field by
field.
Add scenarios that move leadership between engines, by request,
graceful handoff, or the leader's death, and that run without leader
election in both directions. Chaos scenarios disconnect listeners, drop
notifications so workers must fall back to polling, compete for jobs
with `SKIP LOCKED`, hard-abort jobs that ignore cancellation, kill and
restart a worker so its jobs are rescued, and replace each engine's
process in turn during a rolling deployment.
Start a second pair of adapters whose connections use a schema that
makes PostgreSQL look like YugabyteDB, without `LISTEN`/`NOTIFY`. Both
engines must detect it, write unique jobs with a nonce instead of
relying on `xmax`, send no notifications, and poll for cancellations of
running jobs.
Add `TestMaintenanceConformance`, which runs maintenance services on one
engine against rows the other wrote: the rescuer's stale selection and a
full batch of unexpired jobs, job cleaner retention, the queue cleaner
keeping active queues, and the reindexer skipping leftover artifacts.
It also checks leader renewal while maintenance is slow, a new term for
the same client ID, due periodic jobs being inserted available,
migrations in a mixed-case schema, and queue name validation and control
of unknown queues.
Add `TestResilienceConformance`. The harness runs a second Go adapter
and the candidate behind their own TCP proxy, so it can make the
database unavailable to one worker while the reference keeps working,
then checks that the worker reconnects. With direct SQL it also injects
transient completion failures, holds row locks during completion, and
inserts rows that can't be decoded. It also checks how a hard shutdown
classifies jobs that were stopping, and that a job whose cancellation
was attempted during shutdown ends up cancelled.
Add opt-in release performance gates that compare the candidate's
enqueue, worker, and mixed throughput and p95 latency with the reference,
using bounds from the candidate's descriptor, and a mixed soak that also
bounds the connection pool. Both use the same deterministic 10 ms worker
in every engine.

A soak fails at startup when its duration plus time to finish doesn't fit
in the `go test` timeout, so a long soak can't be cut off as a timeout
without saying why.
Add `TestMixedSQLiteConformance`, the `portable-storage-v1` profile.
Both adapters share one temporary SQLite file with WAL and a busy
timeout and check handshakes, migrations, job CRUD, rows written by
either engine, cross-engine insertion and uniqueness, unique column
bytes, batch atomicity, transactions, timestamp rounding and ordering,
and 64-bit job IDs in requests and cursors.
Add `TestMixedSQLiteRuntimeConformance`, the `sqlite-runtime-v1` worker
and queue profile, which works jobs on SQLite across engines: competing
workers, claim order and `attempted_by`, queue CRUD and pause,
notification wakeups and payloads, remote cancellation, leadership and
failover, periodic and scheduler work, poll-only recovery, resumable
jobs, extensions and subscriptions, kinds, and lifecycle shutdown.

`TestResilienceSQLiteConformance` adds completion under a held writer
lock, Go integer ranges, and columns holding invalid JSON.
Add the `insert-only-v1` profile for implementations that insert jobs
but don't work them. `TestInsertOnlyConformance` checks the handshake,
that the reference works what the candidate inserts, typed batches,
transactional insertion, and that inserts notify a reference worker.
Add tiers that start the reference and at least two candidates against
one database at once, so a fault can't degrade into a pairwise test.
They fill one worker slot in every engine, move leadership through every
runtime, terminate each engine's connections, run work, notification,
and cancellation between every ordered pair, kill each candidate in turn
so another implementation takes over leadership and rescues its job, and
interchange list and resumable cursors. SQLite storage and runtime
checks run between every pair of candidates. Opt-in performance and soak
tiers compare release builds running together.

Peers come from `RIVER_CONFORMANCE_PEER` or `RIVER_CONFORMANCE_PEER_FILE`.
Add `retrypolicy.DelayBounds`, which returns the smallest and largest
delay the default policy schedules for a given error count, including
jitter and the cap at the maximum `time.Duration`. A golden generator can
then record the allowed range for each error count instead of
reproducing the policy's arithmetic. The jitter fraction becomes a
constant shared by both.
Implementations must agree with Go on values that are easy to get
subtly wrong: unique keys, retry delays, notification payloads, cron
schedules, and reserved metadata keys. They also need to know which of
Go's features they're expected to match.

`generateconformance` computes golden fixtures from Go itself: unique
key hashes for a matrix of options and argument encodings; protocol
values such as job states, notification topics and payloads, attempt
error encoding, retry delay ranges, and the reserved metadata keys the
feature inventory lists; and cron schedules and snooze counters.
Scenarios check every engine against them on PostgreSQL, SQLite, and
the insert-only profile.

`generatefeatureinventory` extracts Go's public API, configuration,
driver and extension interfaces, metadata keys, notification topics and
payloads, and `rivertype` fields from source, and requires each to be
classified: protocol-visible items name the scenarios that cover them or
record a gap, and every named scenario must have an owner in the harness
registry. It renders `feature-matrix.md` from the result.

`make generate` refreshes both, and a CI job runs their checks so a Go
change that affects the protocol can't land without updating them.
Describe what the conformance directory contains, how scenarios and
their owners are checked, how to run each tier against a candidate, what
the simulated scenarios stand in for, and what River CI runs.
Start a Rust workspace under `rust/` with `riverqueue-migrate`, which
applies, lists, previews, and validates River's PostgreSQL and SQLite
migration lines from Rust. It shares the `river_migration` history with
River Go, so either language can migrate a database the other uses, and
it accepts any quotable PostgreSQL schema name.

The crate can't read files outside its package, so it carries mirrors of
Go's canonical migrations. `syncrustmigrations` writes the mirrors and
the hashes in the conformance migration inventories from Go's migration
directories, and `make verify/rust-migrations`, which the conformance
artifacts CI job now runs, fails when they drift, so Rust never gains an
independent schema history.
Add `riverqueue-macros` with `#[derive(JobArgs)]`, which requires a
stable `#[river(kind = "...")]` and can declare kind aliases, the default
queue, max attempts, priority, pending state, tags, and default unique
options, including the fields that make up a `by_args` key. Invalid
attributes fail at compile time with spans that point at them.
Applications receive the macro through `riverqueue`.
Add `riverqueue`, a Rust and Tokio implementation of River that shares
the database schema and job protocol with River Go on PostgreSQL and
SQLite, so Rust and Go clients can insert and work jobs in the same
database.

`Client` takes a caller-owned SQLx pool and isn't generic over the
database or a driver trait. Workers are typed and async and get a
`CancellationToken`; request builders cover insertion (including
transactional, batch, unique, scheduled, and pending jobs), job and
queue management, and transactional completion from inside a worker.
Like River Go's `*Tx` methods, a request given the caller's transaction
runs directly in it without a savepoint, so the caller rolls back on an
error. The runtime fetches, completes in bounded concurrent batches,
retries, snoozes, cancels, and rescues jobs the way Go does, publishes
events only after their database update commits, and runs leader
election and maintenance services: the rescuer, cleaners, scheduler,
reindexer, and periodic jobs with Go-compatible cron parsing. Hooks,
middleware, error handlers, retry policies, and a hidden extension
module for lockstep add-on crates mirror Go's extension points.

Persisted values follow Go exactly where another engine reads them:
unique keys, metadata, attempt errors, notifications, and timestamps.
Unit tests check unique keys and cron schedules against the conformance
fixtures, which the fixture generator now also copies into the crate
because published tests can't read files outside it.
Now that `riverqueue` exists, test the derive macro end to end: derived
kinds, aliases, insert options, and unique keys, plus compile-fail cases
for each attribute error, checked with `trybuild`.
Add integration tests for the storage layer: the backend contract, job
and queue CRUD, list filters and cursors, exact metadata, the protocol
fixtures, and parity with the rows River Go writes on each database. A
test also fails if a dependency enables `serde_json` features that would
change its behavior for the rest of an application.

PostgreSQL tests build only with `--cfg river_postgres_tests` and run in
a freshly migrated schema with a unique name, so test binaries can share
one disposable database. They fail rather than skip when
`RIVER_RUST_DATABASE_URL` is unset.
Add integration tests for running clients: start and stop lifecycles,
runtime configuration, producer lifetimes, stuck jobs, fetching only
known kinds, poll-only cancellation, running without leader election,
insert notifications, the job and queue handles, and the full SQLite
runtime.
Add integration tests that run clients against malformed rows and
database faults on PostgreSQL and SQLite: proxied connections that are
dropped and refused while clients run, after which producers,
completion, notifications, and leadership must recover. Clients must
also detect a PostgreSQL server that looks like YugabyteDB and fall back
to polling.
Add integration tests for what add-on crates build on: work middleware,
hooks, and error handlers and their ordering; extension services and the
client's stop order; producer sessions and the checks River applies to
their claims; peer attempts that River owns until each outcome persists;
prepared insertion of stored jobs; filtered deletion of finalized jobs;
and batched insert interception.

Tests also cover requests run in a caller's transaction: they open no
savepoint, so an extension step, insert middleware, or decode failure
after River's write leaves that write in the transaction for the caller
to roll back, a failed statement aborts a PostgreSQL transaction, and
every write carries the caller's transaction ID.
Add runnable examples for a basic worker, graceful shutdown,
cancellation, transactional completion, unique and periodic jobs, event
subscriptions, custom schemas, SQLite, and a deployment where Go and
Rust clients share one database.
Add `riverqueue-test`: assertions that check which jobs a test's code
inserted, with optional expected properties and variants that read
through an open transaction, and helpers that run a worker once with or
without a database.
Add `riverqueue-cli`, which installs a `riverqueue` binary that migrates
PostgreSQL and SQLite databases like `river migrate-*` and benchmarks
worker throughput and end-to-end latency like `river bench`.
Add `riverqueue-conformance`, an unpublished crate that serves the
conformance adapter contract with `riverqueue`, and a candidate
descriptor that builds and starts it. The harness now uses it as the
default candidate and peer, and its artifact checks require at least
one checked descriptor. A test in the crate keeps the manifest's Rust
version in step with the workspace.
Add Makefile targets that lint the Rust workspace, including
PostgreSQL-only and SQLite-only builds, build documentation and doctests
for each backend, test it with and without PostgreSQL, audit
dependencies with `cargo deny`, package the crates, check semver against
the last Rust release, and run the benchmark. Rust targets stay out of
`make lint` and `make test`, so Go contributors don't need a Rust
toolchain. `SQLC` can now override the `sqlc` binary.

A Rust workflow runs those checks, the unit and SQLite tests on each
supported Rust version, SQLite conformance, and for PostgreSQL 14
through 18 the PostgreSQL tests, mixed and insert-only conformance, and
a ten-minute soak, with advisory performance gates. A manual release
candidate workflow gates on performance and a one-hour soak, and a
weekly workflow runs a six-hour soak.
Have Dependabot update the Rust workspace's dependencies weekly, grouping
minor and patch updates after a seven-day cooldown, and the GitHub
Actions the workflows use.
Add a workspace README listing the crates and how to run their checks,
benchmark, and PostgreSQL tests, and a changelog for the crates, which
are versioned and released together.
@bgentry
bgentry force-pushed the bg/conformance-harness branch 2 times, most recently from 97f1e4a to 96c7b6c Compare October 5, 2026 17:19
@brandur brandur mentioned this pull request Oct 5, 2026
@bgentry
bgentry force-pushed the bg/conformance-harness branch from 96c7b6c to 3a61e95 Compare October 5, 2026 19:15
@bgentry

bgentry commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

Superseded by #1442, which merged the Rust port without the conformance suite. Conformance is being slimmed down and will come back as separate, smaller PRs.

@bgentry bgentry closed this Oct 5, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant