use std::num::NonZeroU32;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use taskvisor::prelude::*;
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let sup = Supervisor::new(SupervisorConfig::default(), vec![]);
let handle = sup.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!("=== Canceled ===");
let worker: TaskRef = TaskFn::arc("worker", |ctx| async move {
ctx.cancelled().await;
Ok(())
});
let (id, waiter) = handle.add_and_watch(TaskSpec::restartable(worker)).await?;
tokio::time::sleep(Duration::from_millis(100)).await;
println!(" cancelling worker...");
handle.cancel(id).await?;
println!(" worker -> {:?}\n", waiter.wait().await?);
handle.shutdown().await?;
println!("Done.");
Ok(())
}