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
7 changes: 6 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -206,9 +206,14 @@ check/js/package: build/js
check/rust/dependencies: ## Audit Rust advisories, licenses, bans, and sources
cd rust && cargo deny check

# Cargo caches temporary registry dependencies by path and version. A fresh
# build directory prevents stale sources and binaries after same-version edits.
# Verified archives still go to the normal target/package directory.
.PHONY: check/rust/package
check/rust/package: ## Build and verify publishable crate archives without publishing
cd rust && cargo package --workspace --allow-dirty --locked
cd rust && package_build_dir=$$(mktemp -d) && \
trap 'rm -rf "$$package_build_dir"' EXIT && \
CARGO_BUILD_BUILD_DIR="$$package_build_dir" cargo package --workspace --allow-dirty --locked

# The baseline is the latest published riverqueue-v* tag, and
# cargo-semver-checks infers the allowed change from the version bump. It
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue-cli/src/bench.rs
Original file line number Diff line number Diff line change
Expand Up @@ -440,7 +440,7 @@ async fn run_benchmark(options: BenchOptions) -> Result<(), Box<dyn StdError + S
if let Some(producer) = producer {
producer.await.map_err(|error| join_error(&error))??;
}
run.shutdown().await?;
run.stop().await?;
event_cancel.cancel();
event_task.await.map_err(|error| join_error(&error))??;
run_result?;
Expand Down
8 changes: 4 additions & 4 deletions rust/riverqueue/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,8 +102,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
```

Apply migrations before any client starts, and start clients inside a Tokio
runtime. `Client::start` returns a `RunHandle`: await `wait`, `shutdown` (a
soft stop that lets running jobs finish), or `shutdown_now` (which cancels
runtime. `Client::start` returns a `RunHandle`: await `wait`, `stop` (a
soft stop that lets running jobs finish), or `stop_and_cancel` (which cancels
them). `RunHandle::stopper` returns a cloneable `Stopper` for stopping the
client from another task, such as a signal handler. The handle controls the
running client: dropping every `Client` clone doesn't stop it, dropping the
Expand Down Expand Up @@ -224,8 +224,8 @@ async fn build_report(
}
```

During a client's hard stop (`RunHandle::shutdown_now` or
`Stopper::stop_now`), a job whose worker returns `WorkCancelled`, anywhere in
During a client's hard stop (`RunHandle::stop_and_cancel` or
`Stopper::stop_and_cancel`), a job whose worker returns `WorkCancelled`, anywhere in
its error's source chain, becomes available again without using up its
attempt. Any other error is recorded and consumes the attempt, and `Ok`
completes the job. After the configured stuck threshold, River can abort a
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/docs/mixed-deployments.md
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,7 @@ with `SQLITE_BUSY` when a Rust process commits in between.
## Rolling back

Rolling back doesn't touch the schema. Stop Rust clients gracefully with
`RunHandle::shutdown` and let the Go clients continue. Jobs that Rust inserted
`RunHandle::stop` and let the Go clients continue. Jobs that Rust inserted
are ordinary River rows that Go workers can run, and any job a stopped Rust
client left running is recovered by the rescuer. Only migrate down as a
separately planned operation once no deployed client needs the newer schema.
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/basic_worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
.await?;
while completed.recv().await?.as_job().map(|event| event.job.id) != Some(inserted.id()) {}

run.shutdown().await?;
run.stop().await?;
Ok(())
}
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/cancellation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
let job = client.insert(CancellableReport { report_id: 42 }).await?;

client.jobs().cancel(job.id()).await?;
run.shutdown().await?;
run.stop().await?;
Ok(())
}
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
seen += 1;
}

run.shutdown().await?;
run.stop().await?;
Ok(())
}
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/graceful_shutdown.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ async fn main() -> Result<(), Box<dyn Error>> {
stopper.stop();
}
if tokio::signal::ctrl_c().await.is_ok() {
stopper.stop_now();
stopper.stop_and_cancel();
}
});

Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/mixed_go_rust.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
while completed.recv().await?.as_job().map(|event| event.job.id) != Some(inserted.id()) {}
println!("send_receipt is waiting in the default queue for the Go service");

run.shutdown().await?;
run.stop().await?;
Ok(())
}
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/periodic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
break;
}
}
run.shutdown().await?;
run.stop().await?;
Ok(())
}
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/sqlite.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ async fn main() -> Result<(), Box<dyn Error>> {
while completed.recv().await?.as_job().map(|event| event.job.id) != Some(inserted.id()) {}
println!("job {} completed", inserted.id());

run.shutdown().await?;
run.stop().await?;
pool.close().await;
std::fs::remove_dir_all(&directory).ok();
Ok(())
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/examples/transactions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,6 @@ async fn main() -> Result<(), Box<dyn Error>> {
.await?;
assert!(confirmed);

run.shutdown().await?;
run.stop().await?;
Ok(())
}
4 changes: 2 additions & 2 deletions rust/riverqueue/src/client/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -568,11 +568,11 @@ impl ClientBuilder {
/// The timeout must be positive.
///
/// The client starts this timer when fetching stops, however the stop was
/// requested: [`RunHandle::shutdown`](crate::RunHandle::shutdown),
/// requested: [`RunHandle::stop`](crate::RunHandle::stop),
/// [`Stopper::stop`](crate::Stopper::stop), or the signal passed to
/// [`Client::start_with_graceful_shutdown`]. Jobs still running when it
/// expires are cancelled as if by
/// [`Stopper::stop_now`](crate::Stopper::stop_now).
/// [`Stopper::stop_and_cancel`](crate::Stopper::stop_and_cancel).
#[must_use]
pub fn soft_stop_timeout(mut self, timeout: Duration) -> Self {
self.soft_stop_timeout = Some(timeout);
Expand Down
42 changes: 21 additions & 21 deletions rust/riverqueue/src/client/run.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ impl Client {
/// client stops fetching jobs and lets running jobs finish, and the
/// builder's `soft_stop_timeout` escalates to cancelling them when set.
/// Hard stops remain available through [`RunHandle::stopper`] and
/// [`RunHandle::shutdown_now`]. The client drops `signal` without
/// [`RunHandle::stop_and_cancel`]. The client drops `signal` without
/// awaiting it further once it stops for any other reason.
///
/// This mirrors the graceful shutdown hooks of Tokio servers such as
Expand Down Expand Up @@ -452,8 +452,8 @@ impl Supervisor {
/// soft stop to a hard stop after `soft_stop_timeout`.
///
/// The escalation belongs to the client rather than to a caller awaiting
/// [`RunHandle::shutdown`], so it applies however the stop was requested and
/// dropping a shutdown future never changes it. This matches Go's client,
/// [`RunHandle::stop`], so it applies however the stop was requested and
/// dropping a stop future never changes it. This matches Go's client,
/// which starts its soft stop timer when fetching stops.
async fn watch_stop(
fetch_cancel: CancellationToken,
Expand Down Expand Up @@ -495,7 +495,7 @@ async fn watch_stop(
/// immediately; observe completion through the handle.
///
/// Requests are idempotent and ordered by severity: calling [`Stopper::stop`]
/// after [`Stopper::stop_now`] does not undo the hard stop, and requests made
/// after [`Stopper::stop_and_cancel`] does not undo the hard stop, and requests made
/// after the client stopped do nothing. A stopper only affects the run it came
/// from, not a later restart of the same [`Client`].
///
Expand All @@ -511,7 +511,7 @@ async fn watch_stop(
/// stopper.stop();
/// let _ = tokio::signal::ctrl_c().await;
/// // A second Ctrl-C cancels jobs that are still running.
/// stopper.stop_now();
/// stopper.stop_and_cancel();
/// });
/// run.wait().await
/// # }
Expand All @@ -530,7 +530,7 @@ impl Stopper {
/// stop at once, while each queue's producer keeps reporting its
/// running jobs until they finish. When the builder's `soft_stop_timeout`
/// is set, jobs still running after that timeout are cancelled as if by
/// [`Stopper::stop_now`].
/// [`Stopper::stop_and_cancel`].
pub fn stop(&self) {
self.fetch_cancel.cancel();
}
Expand All @@ -550,16 +550,16 @@ impl Stopper {
/// attempt. A job whose cancellation was requested with
/// [`Jobs::cancel`](crate::Jobs::cancel) is cancelled rather than made
/// available.
pub fn stop_now(&self) {
pub fn stop_and_cancel(&self) {
self.fetch_cancel.cancel();
self.work_cancel.cancel();
}
}

/// Controls one running client instance.
///
/// [`RunHandle::wait`], [`RunHandle::shutdown`], and
/// [`RunHandle::shutdown_now`] take `&mut self`, can be called repeatedly, and
/// [`RunHandle::wait`], [`RunHandle::stop`], and
/// [`RunHandle::stop_and_cancel`] take `&mut self`, can be called repeatedly, and
/// are cancel safe: dropping one of their futures, for example from
/// `tokio::time::timeout` or `tokio::select!`, leaves the client and the
/// handle as they were, apart from any stop the method already requested. To
Expand All @@ -570,9 +570,9 @@ impl Stopper {
/// The client's result is reported to the first call that observes it
/// stopping; later calls return `Ok(())`.
///
/// Dropping the handle requests a hard stop, like [`Stopper::stop_now`], but
/// cannot wait for in-flight work to be recorded. Use [`RunHandle::shutdown`]
/// or [`RunHandle::shutdown_now`] when shutdown must finish before returning,
/// Dropping the handle requests a hard stop, like [`Stopper::stop_and_cancel`], but
/// cannot wait for in-flight work to be recorded. Use [`RunHandle::stop`]
/// or [`RunHandle::stop_and_cancel`] when shutdown must finish before returning,
/// or [`RunHandle::detach`] to deliberately leave the client running.
///
/// # Examples
Expand All @@ -585,11 +585,11 @@ impl Stopper {
///
/// let mut run = client.start()?;
/// // ... serve until the application stops ...
/// if tokio::time::timeout(Duration::from_secs(30), run.shutdown())
/// if tokio::time::timeout(Duration::from_secs(30), run.stop())
/// .await
/// .is_err()
/// {
/// run.shutdown_now().await?;
/// run.stop_and_cancel().await?;
/// }
/// # Ok(())
/// # }
Expand Down Expand Up @@ -631,7 +631,7 @@ impl RunHandle {
/// fails, or the process exits. Nothing observes its result, and jobs
/// running when the process exits are left `running` for the rescuer.
/// Most applications should keep the handle and await
/// [`RunHandle::shutdown`] instead.
/// [`RunHandle::stop`] instead.
pub fn detach(mut self) {
// Without a join handle, dropping the handle requests no stop.
self.join.take();
Expand All @@ -647,20 +647,20 @@ impl RunHandle {
/// This method is cancel safe. Dropping the future after its first poll
/// leaves the soft stop in progress, including any `soft_stop_timeout`
/// escalation, and never escalates to a hard stop by itself. The handle
/// remains usable: call [`RunHandle::shutdown_now`] to cancel running jobs
/// remains usable: call [`RunHandle::stop_and_cancel`] to cancel running jobs
/// or [`RunHandle::wait`] to keep waiting.
///
/// # Errors
///
/// Returns the error that stopped the client, as [`RunHandle::wait`] does.
pub async fn shutdown(&mut self) -> Result<(), Error> {
pub async fn stop(&mut self) -> Result<(), Error> {
self.stopper.stop();
self.wait().await
}

/// Requests a hard stop and waits for the client to stop.
///
/// This is [`Stopper::stop_now`] followed by [`RunHandle::wait`]. The stop
/// This is [`Stopper::stop_and_cancel`] followed by [`RunHandle::wait`]. The stop
/// is requested when the future is first polled.
///
/// # Cancel safety
Expand All @@ -671,8 +671,8 @@ impl RunHandle {
/// # Errors
///
/// Returns the error that stopped the client, as [`RunHandle::wait`] does.
pub async fn shutdown_now(&mut self) -> Result<(), Error> {
self.stopper.stop_now();
pub async fn stop_and_cancel(&mut self) -> Result<(), Error> {
self.stopper.stop_and_cancel();
self.wait().await
}

Expand Down Expand Up @@ -742,7 +742,7 @@ impl RunHandle {
impl Drop for RunHandle {
fn drop(&mut self) {
if self.join.is_some() {
self.stopper.stop_now();
self.stopper.stop_and_cancel();
}
}
}
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/src/client/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -589,7 +589,7 @@ async fn readiness_survives_a_notification_listener_panic() {
client.inner.notifier_start_panics.load(Ordering::Acquire),
0
);
run.shutdown().await.unwrap();
run.stop().await.unwrap();
pool.close().await;
for suffix in ["", "-shm", "-wal"] {
let mut file = path.as_os_str().to_owned();
Expand Down
4 changes: 2 additions & 2 deletions rust/riverqueue/src/maintenance/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1008,7 +1008,7 @@ async fn periodic_start_hooks_and_run_on_start_follow_each_leadership_gain() {
wait_for(1).await;
client.request_resign().await.unwrap();
wait_for(2).await;
handle.shutdown().await.unwrap();
handle.stop().await.unwrap();
assert_eq!(starts.load(Ordering::SeqCst), 2);
assert_eq!(periodic_count().await, 2);
database.cleanup().await;
Expand Down Expand Up @@ -1089,7 +1089,7 @@ async fn maintenance_start_retries_then_resigns(poll_only: bool) {
.expect("maintenance start should be retried in a new term");
let second_term = elected_at().await.unwrap();
assert_ne!(second_term, first_term);
handle.shutdown().await.unwrap();
handle.stop().await.unwrap();
assert_eq!(attempts.load(Ordering::SeqCst), 4);
database.cleanup().await;
}
10 changes: 5 additions & 5 deletions rust/riverqueue/tests/client_handles.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ macro_rules! scenarios {
.unwrap();
let mut run = client.start().unwrap();
wait_for_completion(&client, id).await;
run.shutdown().await.unwrap();
run.stop().await.unwrap();

let attempted_by = client.jobs().get(id).await.unwrap().attempted_by;
let mut expected = previous[1..].to_vec();
Expand Down Expand Up @@ -134,7 +134,7 @@ macro_rules! scenarios {
.unwrap();
client.jobs().cancel(id).await.unwrap();
wait_for_completion(&client, id).await;
run.shutdown().await.unwrap();
run.stop().await.unwrap();
fixture.cleanup().await;
}

Expand Down Expand Up @@ -408,7 +408,7 @@ macro_rules! scenarios {
.unwrap();
let mut run = client.start().unwrap();
wait_for_completion(&client, ids[4]).await;
run.shutdown().await.unwrap();
run.stop().await.unwrap();

assert_eq!(
*worked.lock().unwrap(),
Expand Down Expand Up @@ -465,7 +465,7 @@ macro_rules! scenarios {
let finalized_at = job.finalized_at.unwrap();
assert!((chrono::Utc::now() - finalized_at).num_seconds().abs() < 2);

run.shutdown().await.unwrap();
run.stop().await.unwrap();
fixture.cleanup().await;
}

Expand Down Expand Up @@ -861,7 +861,7 @@ macro_rules! scenarios {
);
assert!(!client.local_queues().configs().contains_key("dynamic"));

run.shutdown().await.unwrap();
run.stop().await.unwrap();
fixture.cleanup().await;
}
};
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/tests/extension_services.rs
Original file line number Diff line number Diff line change
Expand Up @@ -312,7 +312,7 @@ async fn assert_maintenance_services_restart_within_their_term(builder: riverque
>= 3
})
.await;
run.shutdown().await.unwrap();
run.stop().await.unwrap();

let terms = pilot
.calls
Expand Down
2 changes: 1 addition & 1 deletion rust/riverqueue/tests/fetch_only_known_kinds.rs
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@ async fn claims_only_registered_kinds(backend: Backend) {
}
assert_eq!(completed, HashSet::from([known, alias]));

tokio::time::timeout(TIMEOUT, run.shutdown())
tokio::time::timeout(TIMEOUT, run.stop())
.await
.expect("the client should stop")
.unwrap();
Expand Down
6 changes: 3 additions & 3 deletions rust/riverqueue/tests/leader_election_disabled.rs
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,7 @@ async fn works_jobs_without_electing(backend: Backend, poll_only: bool) {
assert_eq!(next_completed(&mut events).await.id, inserted.job.row.id);
assert_eq!(backend.leader_id().await, None);

tokio::time::timeout(Duration::from_secs(5), run.shutdown())
tokio::time::timeout(Duration::from_secs(5), run.stop())
.await
.expect("the client should stop")
.unwrap();
Expand Down Expand Up @@ -258,7 +258,7 @@ async fn stays_ineligible_after_leader_stops(backend: Backend) {
Some("eligible_leader")
);

tokio::time::timeout(Duration::from_secs(5), leader_run.shutdown())
tokio::time::timeout(Duration::from_secs(5), leader_run.stop())
.await
.expect("the leader should stop")
.unwrap();
Expand All @@ -273,7 +273,7 @@ async fn stays_ineligible_after_leader_stops(backend: Backend) {
assert_eq!(pilot.maintenance_starts.load(Ordering::SeqCst), 0);
assert_eq!(pilot.runtime_starts.load(Ordering::SeqCst), 1);

run.shutdown().await.unwrap();
run.stop().await.unwrap();
backend.cleanup().await;
}

Expand Down
Loading
Loading