use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use super::*;
#[tokio::test]
async fn sequential_contains_one_failure() {
let bad = agent(Behavior::ErrHandle, 1);
let bad_id = bad.id();
let agents = vec![
agent(Behavior::Complete, 1),
bad,
agent(Behavior::Complete, 1),
];
let mut reactor: Reactor<_, _, TestAgent> =
Reactor::new(MockInference::end_turns(3), MemStore::default(), agents);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 2);
assert_eq!(report.failed, 1);
assert!(report.errors.contains_key(&bad_id));
}
#[tokio::test]
async fn batch_contains_one_failure() {
let bad = batch_agent(Behavior::ErrHandle, 1);
let bad_id = bad.id();
let agents = vec![
batch_agent(Behavior::Complete, 1),
bad,
batch_agent(Behavior::Complete, 1),
];
let mut reactor: Reactor<_, _, TestAgent> =
Reactor::new(MockInference::default(), MemStore::default(), agents);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 2);
assert_eq!(report.failed, 1);
assert!(report.errors.contains_key(&bad_id));
}
#[tokio::test]
async fn teardown_runs_on_drive_error_both_paths() {
let seq = agent(Behavior::ErrHandle, 1);
let bat = batch_agent(Behavior::ErrHandle, 1);
let (seq_td, bat_td) = (seq.teardowns.clone(), bat.teardowns.clone());
let (seq_id, bat_id) = (seq.id(), bat.id());
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
MockInference::end_turns(1),
MemStore::default(),
vec![seq, bat],
);
let report = reactor.run().await.unwrap();
assert_eq!(report.failed, 2);
assert_eq!(seq_td.load(Ordering::SeqCst), 1, "agent-major tore down");
assert_eq!(bat_td.load(Ordering::SeqCst), 1, "round-major tore down");
assert_eq!(
report.errors[&seq_id].kind,
ErrorKind::Agent,
"drive error wins"
);
assert_eq!(report.errors[&bat_id].kind, ErrorKind::Agent);
}
#[tokio::test]
async fn failed_models_probe_retains_cohort() {
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
FlakyModels::failing(1),
MemStore::default(),
vec![agent(Behavior::Complete, 1), agent(Behavior::Complete, 1)],
);
let err = reactor.run().await.expect_err("first probe fails the run");
assert_eq!(err.kind, ErrorKind::Inference);
assert_eq!(
err.retry_after(),
Some(Duration::from_millis(1)),
"the retry classification survives the erasure"
);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 2, "cohort survived the failed probe");
}
#[tokio::test]
async fn fatal_item_fails_without_burning_retries() {
let calls = Arc::new(AtomicUsize::new(0));
let transport = FailingBatch {
calls: calls.clone(),
fatal: true,
};
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
transport,
MemStore::default(),
vec![batch_agent(Behavior::Complete, 1)],
);
let report = reactor.run().await.unwrap();
assert_eq!(report.failed, 1);
assert_eq!(calls.load(Ordering::SeqCst), 1, "failed on the first round");
}
#[tokio::test]
async fn transient_item_retries_to_cap() {
let calls = Arc::new(AtomicUsize::new(0));
let transport = FailingBatch {
calls: calls.clone(),
fatal: false,
};
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
transport,
MemStore::default(),
vec![batch_agent(Behavior::Complete, 1)],
);
let report = reactor.run().await.unwrap();
assert_eq!(report.failed, 1);
assert_eq!(
calls.load(Ordering::SeqCst),
3,
"re-batched up to MAX_BATCH_ITEM_RETRIES"
);
}
#[tokio::test]
async fn dead_transport_attributes_error_to_all_live_agents() {
let agents = vec![
batch_agent(Behavior::Complete, 1),
batch_agent(Behavior::Complete, 2),
batch_agent(Behavior::Complete, 3),
];
let ids: Vec<AgentId> = agents.iter().map(|a| a.id()).collect();
let mut reactor: Reactor<_, _, TestAgent> =
Reactor::new(DeadBatch, MemStore::default(), agents);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 0);
assert_eq!(report.failed, 3);
for id in ids {
let err = report.errors.get(&id).expect("every live agent attributed");
assert_eq!(err.kind, ErrorKind::Inference);
assert_eq!(
err.retry_after,
Some(Duration::from_millis(1)),
"the retry hint survives per agent"
);
}
}
#[tokio::test]
async fn rejected_agent_serialize_failure_is_recorded() {
let mut greedy = agent(Behavior::Complete, 1);
greedy.model.max_input_tokens = 1;
greedy.state.poison = Some(());
let greedy_id = greedy.id();
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
MockInference::end_turns(1),
MemStore::default(),
vec![agent(Behavior::Complete, 1), greedy],
);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 1, "the satisfiable agent still ran");
assert!(report.rejected.is_empty(), "no snapshot could be taken");
let err = report.errors.get(&greedy_id).expect("never silent");
assert_eq!(err.kind, ErrorKind::Storage);
}
#[tokio::test]
async fn sequential_retries_transient_infer() {
let inference = FlakyInfer::failing(2);
let calls = inference.calls.clone();
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
inference,
MemStore::default(),
vec![agent(Behavior::Complete, 1)],
);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 1, "{:?}", report.errors);
assert_eq!(report.failed, 0);
assert_eq!(calls.load(Ordering::SeqCst), 3);
}
#[tokio::test]
async fn sequential_transient_infer_exhausts_budget() {
let inference = FlakyInfer::failing(usize::MAX);
let calls = inference.calls.clone();
let mut reactor: Reactor<_, _, TestAgent> = Reactor::new(
inference,
MemStore::default(),
vec![agent(Behavior::Complete, 1)],
);
let report = reactor.run().await.unwrap();
assert_eq!(report.done, 0);
assert_eq!(report.failed, 1);
let (_, e) = report.errors.iter().next().unwrap();
assert_eq!(e.kind, ErrorKind::Inference);
assert_eq!(calls.load(Ordering::SeqCst), MAX_INFER_RETRIES as usize + 1);
}