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: 5 additions & 1 deletion rust/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,15 @@ Changes to River for Go are recorded in the [repository changelog](../CHANGELOG.

## [Unreleased]

### Changed

- **Breaking:** Renamed `WorkerRegistry` to `Workers`, `register` to `add`, and `register_fn` to `add_fn` to align worker registration with Go's naming. Both methods retain their `Result` return type and registration behavior. [PR #1469](https://github.com/riverqueue/river/pull/1469).

## [0.2.0] - 2026-10-06

### Changed

- Renamed `Client::start_with_graceful_shutdown` to `Client::start_with_graceful_stop` to match River's standard start/stop terminology. This is a rare breaking name change as the new Rust API stabilizes. [PR #1467](https://github.com/riverqueue/river/pull/1467).
- **Breaking:** Renamed `Client::start_with_graceful_shutdown` to `Client::start_with_graceful_stop` to match River's standard start/stop terminology. This is a rare breaking name change as the new Rust API stabilizes. [PR #1467](https://github.com/riverqueue/river/pull/1467).

## [0.1.0] - 2026-10-06

Expand Down
6 changes: 3 additions & 3 deletions rust/riverqueue-cli/src/bench.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ use std::{

use riverqueue::{
Client, EventKind, EventReceiver, EventRecvError, InsertOpts, Job, JobArgs, QueueConfig,
SubscribeConfig, WorkContext, WorkOutcome, Worker, WorkerRegistry, database::SchemaName,
SubscribeConfig, WorkContext, WorkOutcome, Worker, Workers, database::SchemaName,
};
use serde::{Deserialize, Serialize};
use sqlx::{
Expand Down Expand Up @@ -467,8 +467,8 @@ async fn run_benchmark(options: BenchOptions) -> Result<(), Box<dyn StdError + S
}

fn benchmark_client(pool: PgPool, max_workers: usize) -> Result<Client, riverqueue::Error> {
let mut workers = WorkerRegistry::new();
workers.register::<BenchmarkArgs, _>(BenchmarkWorker)?;
let mut workers = Workers::new();
workers.add::<BenchmarkArgs, _>(BenchmarkWorker)?;
Client::builder(pool)
.id("riverqueue-benchmark")
.workers(workers)
Expand Down
12 changes: 8 additions & 4 deletions rust/riverqueue/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ apply River's migrations, and start a client:
use riverqueue::migrate::PostgresMigrator;
use riverqueue::sqlx::PgPool;
use riverqueue::{
BoxError, Client, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, WorkerRegistry,
BoxError, Client, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, Workers,
};
use serde::{Deserialize, Serialize};

Expand All @@ -75,8 +75,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
// `riverqueue migrate-up` from `riverqueue-cli` at deploy time instead.
PostgresMigrator::new(pool.clone()).migrate_up().await?;

let mut workers = WorkerRegistry::new();
workers.register_fn(send_email)?;
let mut workers = Workers::new();
workers.add_fn(send_email)?;

let client = Client::builder(pool)
.workers(workers)
Expand Down Expand Up @@ -229,8 +229,12 @@ completes the job. After the configured stuck threshold, River can abort a
Tokio task that yields, which fails its attempt, but it can't stop CPU-bound
work or a blocking call already in progress.

Collect workers in [`Workers`], the equivalent of Go's `Workers` bundle.
Use [`Workers::add`] to add a [`Worker`] implementation, like Go's `AddWorker`.
Implement [`Worker`] when a kind needs a custom timeout or next-retry decision.
Use `WorkerRegistry::register_fn` for an async function or capturing closure.
Use [`Workers::add_fn`] for an async function or capturing closure; it combines
Go's `WorkFunc` adapter and registration in one call. Both methods return
`Result` so registration errors can be propagated with `?`.

## Events

Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/examples/basic_worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@ use std::error::Error;

use riverqueue::sqlx::PgPool;
use riverqueue::{
BoxError, Client, EventKind, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome,
WorkerRegistry, migrate::PostgresMigrator,
BoxError, Client, EventKind, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, Workers,
migrate::PostgresMigrator,
};
use serde::{Deserialize, Serialize};

Expand All @@ -30,8 +30,8 @@ async fn main() -> Result<(), Box<dyn Error>> {
let pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
// Apply River's schema before starting a client.
PostgresMigrator::new(pool.clone()).migrate_up().await?;
let mut workers = WorkerRegistry::new();
workers.register_fn(send_email)?;
let mut workers = Workers::new();
workers.add_fn(send_email)?;
let client = Client::builder(pool)
.workers(workers)
.queue("default", QueueConfig::new(10))
Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/examples/cancellation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,8 @@ use std::{error::Error, time::Duration};

use riverqueue::sqlx::PgPool;
use riverqueue::{
Client, Job, JobArgs, QueueConfig, WorkCancelled, WorkContext, WorkOutcome, Worker,
WorkerRegistry, migrate::PostgresMigrator,
Client, Job, JobArgs, QueueConfig, WorkCancelled, WorkContext, WorkOutcome, Worker, Workers,
migrate::PostgresMigrator,
};
use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -49,8 +49,8 @@ impl Worker<CancellableReport> for CancellableReportWorker {
async fn main() -> Result<(), Box<dyn Error>> {
let pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
PostgresMigrator::new(pool.clone()).migrate_up().await?;
let mut workers = WorkerRegistry::new();
workers.register(CancellableReportWorker)?;
let mut workers = Workers::new();
workers.add(CancellableReportWorker)?;
let client = Client::builder(pool)
.workers(workers)
.queue("default", QueueConfig::new(1))
Expand Down
6 changes: 3 additions & 3 deletions rust/riverqueue/examples/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use std::error::Error;
use riverqueue::sqlx::PgPool;
use riverqueue::{
Client, Event, EventKind, InsertOpts, Job, JobArgs, JobEventKind, QueueConfig, WorkContext,
WorkOutcome, WorkerRegistry, migrate::PostgresMigrator,
WorkOutcome, Workers, migrate::PostgresMigrator,
};
use serde::{Deserialize, Serialize};

Expand All @@ -38,8 +38,8 @@ async fn main() -> Result<(), Box<dyn Error>> {
let pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
PostgresMigrator::new(pool.clone()).migrate_up().await?;

let mut workers = WorkerRegistry::new();
workers.register_fn(charge_card)?;
let mut workers = Workers::new();
workers.add_fn(charge_card)?;
let client = Client::builder(pool)
.workers(workers)
.queue("default", QueueConfig::new(4))
Expand Down
7 changes: 3 additions & 4 deletions rust/riverqueue/examples/graceful_shutdown.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,7 @@ use std::{error::Error, time::Duration};

use riverqueue::sqlx::PgPool;
use riverqueue::{
BoxError, Client, Job, JobArgs, QueueConfig, WorkCancelled, WorkContext, WorkOutcome,
WorkerRegistry,
BoxError, Client, Job, JobArgs, QueueConfig, WorkCancelled, WorkContext, WorkOutcome, Workers,
};
use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -35,8 +34,8 @@ async fn generate_report(
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
let mut workers = WorkerRegistry::new();
workers.register_fn(generate_report)?;
let mut workers = Workers::new();
workers.add_fn(generate_report)?;
let client = Client::builder(pool)
.workers(workers)
.queue("default", QueueConfig::new(10))
Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/examples/mixed_go_rust.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,8 @@ use std::error::Error;

use riverqueue::sqlx::PgPool;
use riverqueue::{
Client, EventKind, InsertOpts, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome,
WorkerRegistry, migrate::PostgresMigrator,
Client, EventKind, InsertOpts, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, Workers,
migrate::PostgresMigrator,
};
use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -86,8 +86,8 @@ async fn main() -> Result<(), Box<dyn Error>> {
let pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
PostgresMigrator::new(pool.clone()).migrate_up().await?;

let mut workers = WorkerRegistry::new();
workers.register_fn(resize_image)?;
let mut workers = Workers::new();
workers.add_fn(resize_image)?;
let client = Client::builder(pool)
.workers(workers)
// Only Rust's queue: Rust never fetches the Go service's jobs.
Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/examples/periodic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use std::{error::Error, time::Duration};
use riverqueue::sqlx::PgPool;
use riverqueue::{
Client, CronSchedule, EventKind, IntervalSchedule, Job, JobArgs, PeriodicJob, PeriodicJobOpts,
QueueConfig, WorkContext, WorkOutcome, WorkerRegistry, migrate::PostgresMigrator,
QueueConfig, WorkContext, WorkOutcome, Workers, migrate::PostgresMigrator,
};
use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -45,9 +45,9 @@ async fn main() -> Result<(), Box<dyn Error>> {
let pool = PgPool::connect(&std::env::var("DATABASE_URL")?).await?;
PostgresMigrator::new(pool.clone()).migrate_up().await?;

let mut workers = WorkerRegistry::new();
workers.register_fn(refresh_cache)?;
workers.register_fn(nightly_report)?;
let mut workers = Workers::new();
workers.add_fn(refresh_cache)?;
workers.add_fn(nightly_report)?;
let client = Client::builder(pool)
.workers(workers)
.queue("default", QueueConfig::new(4))
Expand Down
6 changes: 3 additions & 3 deletions rust/riverqueue/examples/sqlite.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use std::{error::Error, str::FromStr, time::Duration};

use riverqueue::sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions};
use riverqueue::{
Client, EventKind, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, WorkerRegistry,
Client, EventKind, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, Workers,
migrate::SqliteMigrator,
};
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -49,8 +49,8 @@ async fn main() -> Result<(), Box<dyn Error>> {
// Apply River's schema before starting a client.
SqliteMigrator::new(pool.clone()).migrate_up().await?;

let mut workers = WorkerRegistry::new();
workers.register_fn(resize_image)?;
let mut workers = Workers::new();
workers.add_fn(resize_image)?;
let client = Client::builder(pool.clone())
.workers(workers)
.queue("default", QueueConfig::new(4))
Expand Down
6 changes: 3 additions & 3 deletions rust/riverqueue/examples/transactions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ use std::error::Error;

use riverqueue::sqlx::{self, PgPool};
use riverqueue::{
Client, EventKind, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, WorkerRegistry,
Client, EventKind, Job, JobArgs, QueueConfig, WorkContext, WorkOutcome, Workers,
migrate::PostgresMigrator,
};
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -63,8 +63,8 @@ async fn main() -> Result<(), Box<dyn Error>> {
.execute(&pool)
.await?;

let mut workers = WorkerRegistry::new();
workers.register_fn(confirm_order)?;
let mut workers = Workers::new();
workers.add_fn(confirm_order)?;
let client = Client::builder(pool.clone())
.workers(workers)
.queue("default", QueueConfig::new(4))
Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/src/client/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ use crate::database::Database;
use crate::periodic::{PeriodicJob, PeriodicJobs};
use crate::{
Client, Error, ErrorHandler, FETCH_COOLDOWN_MIN, FETCH_POLL_INTERVAL_DEFAULT, Hook,
InsertMiddleware, Plugin, QUEUE_NUM_WORKERS_MAX, RetryPolicy, WorkMiddleware, WorkerRegistry,
InsertMiddleware, Plugin, QUEUE_NUM_WORKERS_MAX, RetryPolicy, WorkMiddleware, Workers,
};

/// Default age at which running jobs are rescued (Go
Expand Down Expand Up @@ -362,7 +362,7 @@ pub struct ClientBuilder {
pub(super) retry_policy: Arc<dyn RetryPolicy>,
pub(super) soft_stop_timeout: Option<Duration>,
pub(super) work_middleware: Vec<Arc<dyn crate::extension::DynWorkMiddleware>>,
pub(crate) workers: WorkerRegistry,
pub(crate) workers: Workers,
}

impl std::fmt::Debug for ClientBuilder {
Expand Down Expand Up @@ -642,9 +642,9 @@ impl ClientBuilder {
self
}

/// Installs a typed worker registry.
/// Sets the client's workers.
#[must_use]
pub fn workers(mut self, workers: WorkerRegistry) -> Self {
pub fn workers(mut self, workers: Workers) -> Self {
self.workers = workers;
self
}
Expand Down
6 changes: 3 additions & 3 deletions rust/riverqueue/src/client/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ use crate::maintenance::LeadershipWakeup;
use crate::{
DefaultRetryPolicy, Error, Event, EventKind, EventReceiver, FETCH_COOLDOWN_DEFAULT,
JOB_STUCK_THRESHOLD_DEFAULT, JOB_TIMEOUT_DEFAULT, MAX_ATTEMPTS_DEFAULT, RetryPolicy,
SubscribeConfig, WorkerRegistry,
SubscribeConfig, Workers,
database::{ClientDatabase, Database, DatabasePool, DatabaseTransactionExecutor, IntoDatabase},
periodic::PeriodicJobs,
};
Expand Down Expand Up @@ -154,7 +154,7 @@ pub(crate) struct ClientInner {
soft_stop_timeout: Option<Duration>,
started: AtomicBool,
work_middleware: Vec<Arc<dyn crate::extension::DynWorkMiddleware>>,
pub(crate) workers: WorkerRegistry,
pub(crate) workers: Workers,
}

#[cfg(feature = "sqlite")]
Expand Down Expand Up @@ -322,7 +322,7 @@ impl Client {
retry_policy: Arc::new(DefaultRetryPolicy::default()),
soft_stop_timeout: None,
work_middleware: Vec::new(),
workers: WorkerRegistry::new(),
workers: Workers::new(),
}
}

Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/src/client/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -479,9 +479,9 @@ async fn subscription_forwarder_stops_when_the_receiver_drops() {
#[river(kind = "subscription_forwarder_test")]
struct ForwarderArgs {}

let mut workers = WorkerRegistry::new();
let mut workers = Workers::new();
workers
.register_fn(|_context: WorkContext, _job: Job<ForwarderArgs>| async {
.add_fn(|_context: WorkContext, _job: Job<ForwarderArgs>| async {
Ok::<_, std::convert::Infallible>(WorkOutcome::Complete)
})
.unwrap();
Expand Down Expand Up @@ -674,9 +674,9 @@ async fn readiness_survives_a_notification_listener_panic() {
.migrate_up()
.await
.unwrap();
let mut workers = WorkerRegistry::new();
let mut workers = Workers::new();
workers
.register_fn(|_context: WorkContext, _job: Job<ReadinessArgs>| async {
.add_fn(|_context: WorkContext, _job: Job<ReadinessArgs>| async {
Ok::<_, std::convert::Infallible>(WorkOutcome::Complete)
})
.unwrap();
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,7 @@ pub use sqlx;
/// The `tokio-util` version of [`WorkContext::cancellation_token`]'s
/// [`CancellationToken`](tokio_util::sync::CancellationToken).
pub use tokio_util;
pub use worker::{WorkContext, WorkOutcome, Worker, WorkerRegistry, WorkerTimeout};
pub use worker::{WorkContext, WorkOutcome, Worker, WorkerTimeout, Workers};

/// Default maximum number of attempts for a job.
pub const MAX_ATTEMPTS_DEFAULT: i16 = 25;
Expand Down
12 changes: 5 additions & 7 deletions rust/riverqueue/src/maintenance/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ use super::{
};
use crate::{
Client, Job, JobArgs, JobState, MaintenanceConfig, QueueConfig, SchemaName, UniqueOpts,
WorkContext, WorkOutcome, Worker, WorkerRegistry, WorkerTimeout,
WorkContext, WorkOutcome, Worker, WorkerTimeout, Workers,
database::{PostgresDatabase, PostgresReindexConfig},
};

Expand Down Expand Up @@ -80,13 +80,11 @@ impl Worker<ShortTimeoutArgs> for ShortTimeoutWorker {
}
}

fn workers() -> WorkerRegistry {
let mut workers = WorkerRegistry::new();
fn workers() -> Workers {
let mut workers = Workers::new();
workers.add::<NoTimeoutArgs, _>(NoTimeoutWorker).unwrap();
workers
.register::<NoTimeoutArgs, _>(NoTimeoutWorker)
.unwrap();
workers
.register::<ShortTimeoutArgs, _>(ShortTimeoutWorker)
.add::<ShortTimeoutArgs, _>(ShortTimeoutWorker)
.unwrap();
workers
}
Expand Down
Loading
Loading