Async Runtime Patterns
Use this recipe for Tokio services, background tasks, cancellation, graceful shutdown, channels, timeouts, and retries.
Runtime Selection
#[tokio::main]
async fn main() -> Result<()> {
run().await
}Use the current-thread runtime only for intentionally lightweight, single-threaded workloads:
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<()> {
run().await
}Bounded Channels
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel::<Event>(128);
tokio::spawn(async move {
while let Some(event) = rx.recv().await {
if let Err(error) = process(event).await {
tracing::warn!(?error, "failed to process event");
}
}
});Rules:
- Prefer bounded
mpsc::channeloverunbounded_channel. - Use
oneshotfor single request/response exchanges. - Use
broadcastwhen each subscriber must receive every event. - Use
watchfor latest-value state such as configuration. - Consider backpressure explicitly when choosing channel size.
Graceful Shutdown
use tokio::task::JoinSet;
pub async fn run() -> Result<()> {
let cancellation_token = tokio_util::sync::CancellationToken::new();
let mut tasks = JoinSet::new();
tasks.spawn(worker(cancellation_token.clone()));
tasks.spawn(server(cancellation_token.clone()));
tokio::select! {
// biased! matters here: without it, an always-ready work arm can starve the
// shutdown arm indefinitely. Keep the shutdown branch first under biased mode.
biased;
result = tokio::signal::ctrl_c() => {
result?;
tracing::info!("shutdown signal received");
}
Some(result) = tasks.join_next() => {
result??;
}
}
cancellation_token.cancel();
while let Some(result) = tasks.join_next().await {
result??;
}
Ok(())
}Add tokio-util only when cancellation tokens are needed:
tokio-util = "0.7"If avoiding tokio-util, use watch or broadcast for shutdown signals.
Compile note: the result? / result?? pattern above requires your domain error to implement From<std::io::Error> (for ctrl_c) and From<tokio::task::JoinError> (for task joins). Add those conversions to the error type from recipes/errors.md or map explicitly.
Timeouts
Never make network or external I/O calls without a timeout:
use std::time::Duration;
let response = tokio::time::timeout(Duration::from_secs(10), client.call(request))
.await
.map_err(|_| ServiceError::invalid_input("request timed out"))??;Retries
Retry only operations that are idempotent or explicitly safe to repeat:
use std::{future::Future, time::Duration};
pub async fn retry_transient<F, Fut, T>(mut operation: F) -> Result<T>
where
F: FnMut() -> Fut,
Fut: Future<Output = Result<T>>,
{
let max_attempts = 3;
let mut delay = Duration::from_millis(100);
for attempt in 1..=max_attempts {
match operation().await {
Ok(value) => return Ok(value),
Err(error) if attempt == max_attempts => return Err(error),
Err(error) => {
tracing::warn!(attempt, ?error, "transient operation failed; retrying");
tokio::time::sleep(delay).await;
delay *= 2;
}
}
}
Err(ServiceError::invalid_input("retry policy must allow at least one attempt"))
}Use jitter for high-scale services to avoid thundering herds.
Blocking Work
let parsed = tokio::task::spawn_blocking(move || parse_large_file(path))
.await
.map_err(|error| ServiceError::invalid_input(format!("parser task failed: {error}")))??;Do not perform CPU-heavy work or blocking filesystem/network calls directly on async executor threads.