use std::num::NonZeroU32;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use taskvisor::prelude::*;
use tokio::sync::Notify;
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let supervisor = Supervisor::new(SupervisorConfig::default(), vec![]);
let handle = supervisor.serve();
println!("=== Completed ===");
let job: TaskRef = TaskFn::arc("import", |_ctx| async {
tokio::time::sleep(Duration::from_millis(100)).await;
Ok(())
});
let (_id, waiter) = handle.add_and_watch(TaskSpec::once(job)).await?;
println!(" import -> {:?}\n", waiter.wait().await?);
println!("=== Failed (retries exhausted) ===");
let attempts = Arc::new(AtomicU32::new(0));
let flaky: TaskRef = TaskFn::arc("sync", move |_ctx| {
let attempts = Arc::clone(&attempts);
async move {
let n = attempts.fetch_add(1, Ordering::Relaxed) + 1;
println!(" sync attempt #{n} failing...");
Err(TaskError::fail("upstream 503").with_exit_code(75))
}
});
let spec = TaskSpec::restartable(flaky)
.with_backoff(BackoffPolicy::constant(Duration::from_millis(20)))
.with_max_retries(NonZeroU32::new(2).unwrap());
match handle.add_and_watch(spec).await?.1.wait().await? {
TaskOutcome::Failed {
reason, exit_code, ..
} => {
println!(" sync -> Failed: {reason} (exit_code={exit_code:?})\n");
}
other => println!(" sync -> {other:?}\n"),
}
println!("=== Failed (attempt timed out) ===");
let slow: TaskRef = TaskFn::arc("slow-report", |_ctx| async {
tokio::time::sleep(Duration::from_secs(1)).await;
Ok(())
});
let timed = TaskSpec::once(slow).with_timeout(Duration::from_millis(20));
match handle.add_and_watch(timed).await?.1.wait().await? {
TaskOutcome::Failed { reason, .. } => {
println!(" slow-report -> Failed: {reason}\n");
}
other => println!(" slow-report -> {other:?}\n"),
}
println!("=== Canceled ===");
let started = Arc::new(Notify::new());
let worker: TaskRef = TaskFn::arc("worker", {
let started = Arc::clone(&started);
move |ctx| {
let started = Arc::clone(&started);
async move {
started.notify_one();
ctx.cancelled().await;
Err(TaskError::Canceled)
}
}
});
let (id, waiter) = handle.add_and_watch(TaskSpec::restartable(worker)).await?;
started.notified().await;
println!(" cancelling worker...");
handle.cancel(id).await?;
println!(" worker -> {:?}\n", waiter.wait().await?);
handle.shutdown().await?;
Ok(())
}