use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use taskvisor::prelude::*;
use tokio::sync::{Mutex, mpsc};
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let (tx, rx) = mpsc::unbounded_channel::<String>();
for i in 1..=8 {
tx.send(format!("message-{i}"))?;
}
drop(tx);
let rx = Arc::new(Mutex::new(rx));
let attempts = Arc::new(AtomicU32::new(0));
let consumer: TaskRef = TaskFn::arc("queue-consumer", {
let rx = Arc::clone(&rx);
let attempts = Arc::clone(&attempts);
move |ctx| {
let rx = Arc::clone(&rx);
let attempts = Arc::clone(&attempts);
async move {
let attempt = attempts.fetch_add(1, Ordering::Relaxed) + 1;
if attempt == 1 {
println!(
"[consumer] connect failed (simulated), supervisor retries with backoff"
);
return Err(TaskError::fail("connection refused"));
}
println!("[consumer] connected on attempt #{attempt}");
let mut rx = rx.lock().await;
loop {
match ctx.run_until_cancelled(rx.recv()).await? {
Some(msg) => println!("[consumer] processed {msg}"),
None => {
println!("[consumer] backlog drained, done");
return Ok(());
}
}
}
}
}
});
let spec = TaskSpec::restartable(consumer).with_backoff(
BackoffPolicy::exponential(Duration::from_millis(100))
.with_max(Duration::from_secs(5))
.with_jitter(JitterPolicy::Equal),
);
let sup = Supervisor::new(SupervisorConfig::default(), vec![]);
sup.run(vec![spec]).await?;
println!("Done.");
Ok(())
}