mod support;
use cano::prelude::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::{Duration, Instant};
use support::Flow;
#[derive(Clone, Copy)]
enum Probe {
Fail,
Pending,
Ready,
}
struct ScriptedProbe {
script: Vec<Probe>,
calls: Arc<AtomicU32>,
}
#[task::poll]
impl PollTask<Flow> for ScriptedProbe {
fn on_poll_error(&self) -> PollErrorPolicy {
PollErrorPolicy::RetryOnError { max_errors: 2 }
}
async fn poll(&self, _res: &Resources) -> Result<PollOutcome<Flow>, CanoError> {
let n = self.calls.fetch_add(1, Ordering::SeqCst) as usize;
match self.script.get(n).copied().unwrap_or(Probe::Ready) {
Probe::Fail => Err(CanoError::task_execution(format!("probe {n} refused"))),
Probe::Pending => Ok(PollOutcome::Pending { delay_ms: 0 }),
Probe::Ready => Ok(PollOutcome::Ready(TaskResult::Single(Flow::Done))),
}
}
}
#[tokio::test]
async fn poll_error_budget_resets_on_pending_and_kills_on_burst() {
let calls = Arc::new(AtomicU32::new(0));
let workflow = Workflow::bare()
.register(
Flow::Work,
ScriptedProbe {
script: vec![
Probe::Fail,
Probe::Fail,
Probe::Pending, Probe::Fail,
Probe::Fail,
Probe::Ready,
],
calls: calls.clone(),
},
)
.add_exit_state(Flow::Done);
let result = workflow
.orchestrate(Flow::Work, CancellationToken::disabled())
.await
.unwrap();
assert_eq!(result, Flow::Done);
assert_eq!(calls.load(Ordering::SeqCst), 6);
let calls = Arc::new(AtomicU32::new(0));
let workflow = Workflow::bare()
.register(
Flow::Work,
ScriptedProbe {
script: vec![Probe::Fail, Probe::Fail, Probe::Fail],
calls: calls.clone(),
},
)
.add_exit_state(Flow::Done);
let err = workflow
.orchestrate(Flow::Work, CancellationToken::disabled())
.await
.unwrap_err();
assert_eq!(err.category(), "task_execution", "got: {err}");
assert!(err.to_string().contains("probe 2 refused"), "got: {err}");
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"no fourth poll after the budget"
);
}
struct StuckProbe;
#[task::poll]
impl PollTask<Flow> for StuckProbe {
fn config(&self) -> TaskConfig {
TaskConfig::minimal().with_attempt_timeout(Duration::from_millis(200))
}
async fn poll(&self, _res: &Resources) -> Result<PollOutcome<Flow>, CanoError> {
Ok(PollOutcome::Pending { delay_ms: 5 })
}
}
#[tokio::test]
async fn poll_eternal_pending_is_bounded_by_attempt_timeout() {
let workflow = Workflow::bare()
.register(Flow::Work, StuckProbe)
.add_exit_state(Flow::Done);
let started = Instant::now();
let err = workflow
.orchestrate(Flow::Work, CancellationToken::disabled())
.await
.unwrap_err();
let elapsed = started.elapsed();
assert_eq!(err.category(), "timeout", "got: {err}");
assert!(
elapsed >= Duration::from_millis(200),
"fired early: {elapsed:?}"
);
assert!(
elapsed < Duration::from_secs(30),
"liveness: took {elapsed:?}"
);
}