use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use taskvisor::prelude::*;
use tokio::sync::oneshot;
fn sum_of_primes(limit: u64) -> u64 {
(2..limit)
.filter(|n| (2..).take_while(|d| d * d <= *n).all(|d| n % d != 0))
.sum()
}
#[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();
let attempts = Arc::new(AtomicU32::new(0));
let job: TaskRef = TaskFn::arc("prime-sum", {
let attempts = Arc::clone(&attempts);
move |ctx| {
let attempts = Arc::clone(&attempts);
async move {
let attempt = attempts.fetch_add(1, Ordering::Relaxed) + 1;
let (tx, rx) = oneshot::channel();
rayon::spawn(move || {
let result = if attempt == 1 {
Err("transient compute failure (simulated)".to_string())
} else {
Ok(sum_of_primes(50_000))
};
let _ = tx.send(result);
});
println!("[prime-sum] attempt #{attempt}: computing on rayon...");
match ctx.run_until_cancelled(rx).await? {
Ok(Ok(sum)) => {
println!("[prime-sum] done: {sum}");
Ok(())
}
Ok(Err(reason)) => {
println!("[prime-sum] failed: {reason}");
Err(TaskError::fail(reason))
}
Err(_dropped) => Err(TaskError::fail("compute thread dropped the channel")),
}
}
}
});
let spec = TaskSpec::restartable(job)
.with_backoff(BackoffPolicy::constant(Duration::from_millis(200)));
let (_id, waiter) = handle.add_and_watch(spec, Duration::from_secs(1)).await?;
println!("outcome: {:?}", waiter.wait().await?);
handle.shutdown().await?;
Ok(())
}