use std::time::Duration;
use super::*;
use crate::scheduler::tests::*;
use crate::verifier::ToolCallSummary;
#[test]
fn test_tick_produces_spawn_for_ready() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
let actions = scheduler.tick();
let spawns: Vec<_> = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.collect();
assert_eq!(spawns.len(), 2);
}
#[test]
fn test_tick_dispatches_all_regardless_of_max_parallel() {
let graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[]),
make_node(3, &[]),
make_node(4, &[]),
]);
let mut config = make_config();
config.max_parallel = 2;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
let actions = scheduler.tick();
let spawn_count = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.count();
assert_eq!(
spawn_count, 2,
"max_parallel=2 caps dispatched tasks per tick"
);
}
#[test]
fn test_tick_detects_completion() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
let config = make_config();
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
let actions = scheduler.tick();
let has_done = actions.iter().any(|a| {
matches!(
a,
SchedulerAction::Done {
status: GraphStatus::Completed
}
)
});
assert!(
has_done,
"should emit Done(Completed) when all tasks are terminal"
);
}
#[test]
fn test_completion_event_marks_deps_ready() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "handle-0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-0".to_string(),
outcome: TaskOutcome::Completed {
output: "done".to_string(),
artifacts: vec![],
tool_trace: None,
},
};
scheduler.buffered_events.push_back(event);
let actions = scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Completed);
let has_spawn_1 = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(1)));
assert!(
has_spawn_1 || scheduler.graph.tasks[1].status == TaskStatus::Ready,
"task 1 should be spawned or marked Ready"
);
}
#[test]
fn test_handoff_event_marks_source_completed_and_activates_target() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "handle-0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-0".to_string(),
outcome: TaskOutcome::Handoff {
output: "handing off".to_string(),
goto: TaskRef::ById(TaskId(1)),
tool_trace: None,
},
});
let actions = scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Completed);
assert_eq!(
scheduler.graph.tasks[0]
.result
.as_ref()
.map(|r| r.output.as_str()),
Some("handing off"),
"the source node's own output must be preserved on the Handoff outcome"
);
assert_eq!(scheduler.graph.tasks[1].commanded_from, Some(TaskId(0)));
assert_eq!(scheduler.graph.handoff_count, 1);
assert!(
scheduler.graph.tasks[0].handoff_rejected.is_none(),
"a successful handoff must not set the rejection signal"
);
let has_spawn_1 = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(1)));
assert!(
has_spawn_1 || scheduler.graph.tasks[1].status == TaskStatus::Ready,
"handoff target should be spawned or marked Ready in the same tick"
);
}
#[test]
fn test_handoff_event_rejection_leaves_source_completed_not_failed() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[1].status = TaskStatus::Completed; let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "handle-0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-0".to_string(),
outcome: TaskOutcome::Handoff {
output: "handing off".to_string(),
goto: TaskRef::ById(TaskId(1)),
tool_trace: None,
},
});
let actions = scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Completed,
"a rejected handoff must not escalate the source node to Failed"
);
assert_eq!(
scheduler.graph.handoff_count, 0,
"a rejected handoff must not consume the budget"
);
assert!(
scheduler.graph.tasks[0].handoff_rejected.is_some(),
"critic finding C1: a rejected handoff must be recorded as a graph-visible, \
persisted signal, not just a log line"
);
assert!(
actions.iter().any(|a| matches!(
a,
SchedulerAction::CheckToolOutcome { task_id, .. } if *task_id == TaskId(0)
)),
"CheckToolOutcome must still be emitted for a rejected handoff's own task_id: \
{actions:?}"
);
}
#[test]
fn test_handoff_event_emits_check_tool_outcome_action() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "handle-0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-0".to_string(),
outcome: TaskOutcome::Handoff {
output: "handing off".to_string(),
goto: TaskRef::ById(TaskId(1)),
tool_trace: None,
},
});
let actions = scheduler.tick();
let has_check = actions.iter().any(|a| {
matches!(
a,
SchedulerAction::CheckToolOutcome { task_id, .. } if *task_id == TaskId(0)
)
});
assert!(
has_check,
"a Handoff outcome must emit CheckToolOutcome for its own task_id (#6394)"
);
assert!(
!actions
.iter()
.any(|a| matches!(a, SchedulerAction::Verify { .. })),
"Verify must NOT be emitted when verify_completeness is disabled (#6394): {actions:?}"
);
}
#[test]
fn test_handoff_event_emits_verify_action_when_verify_completeness_enabled() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut config = make_config();
config.verify_completeness = true;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "handle-0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-0".to_string(),
outcome: TaskOutcome::Handoff {
output: "handing off".to_string(),
goto: TaskRef::ById(TaskId(1)),
tool_trace: None,
},
});
let actions = scheduler.tick();
let verify_output = actions.iter().find_map(|a| match a {
SchedulerAction::Verify {
task_id, output, ..
} if *task_id == TaskId(0) => Some(output.clone()),
_ => None,
});
assert_eq!(
verify_output.as_deref(),
Some("handing off"),
"Verify must be emitted for the handoff node and carry its own output (#6394)"
);
}
#[test]
fn test_handoff_event_all_tools_failed_marks_task_failed_not_handoff() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "handle-0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-0".to_string(),
outcome: TaskOutcome::Handoff {
output: "claiming completion".to_string(),
goto: TaskRef::ById(TaskId(1)),
tool_trace: Some(vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}]),
},
});
scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Failed,
"a Handoff outcome with an all-failed tool trace must be Failed, not Completed"
);
assert!(
scheduler.graph.tasks[1].commanded_from.is_none(),
"try_handoff must never run when the handoff node itself is corrected to Failed"
);
assert_eq!(
scheduler.graph.handoff_count, 0,
"no handoff budget must be consumed when the outcome is corrected to Failed"
);
assert!(
scheduler.graph.tasks[0].handoff_rejected.is_none(),
"handoff_rejected is a try_handoff-rejection signal, not used for the \
all-tool-calls-failed short-circuit"
);
}
#[cfg(feature = "llm-planning")]
#[test]
fn test_plan_with_verify_criteria_and_predicate_disabled_reaches_completed() {
use crate::graph::PlanSlug;
use crate::planner::{PlannedTask, PlannerResponse, convert_response_pub};
let response = PlannerResponse {
tasks: vec![
PlannedTask {
task_id: PlanSlug::from("parent"),
title: "Parent".to_string(),
description: "do parent work".to_string(),
agent_hint: None,
depends_on: vec![],
failure_strategy: None,
execution_mode: None,
verify_criteria: Some("output must be valid JSON".to_string()),
tool_allowlist: None,
},
PlannedTask {
task_id: PlanSlug::from("child"),
title: "Child".to_string(),
description: "do child work".to_string(),
agent_hint: None,
depends_on: vec![PlanSlug::from("parent")],
failure_strategy: None,
execution_mode: None,
verify_criteria: None,
tool_allowlist: None,
},
],
};
let graph = convert_response_pub(response, "goal", &[make_def("worker")], 20, false).unwrap();
assert!(
graph.tasks[0].verify_predicate.is_none(),
"verify_predicate must be dropped when verify_predicate_enabled is false"
);
let mut scheduler = make_scheduler(graph);
let actions = scheduler.tick();
assert!(
actions
.iter()
.any(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(0)))
);
scheduler.record_spawn(
TaskId(0),
"handle-parent".to_string(),
"worker".to_string(),
None,
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(0),
agent_handle_id: "handle-parent".to_string(),
outcome: TaskOutcome::Completed {
output: "parent done".to_string(),
artifacts: vec![],
tool_trace: None,
},
});
let actions = scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Completed);
assert!(
actions
.iter()
.any(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(1))),
"child must be dispatched once its only dependency completes"
);
scheduler.record_spawn(
TaskId(1),
"handle-child".to_string(),
"worker".to_string(),
None,
);
scheduler.buffered_events.push_back(TaskEvent {
task_id: TaskId(1),
agent_handle_id: "handle-child".to_string(),
outcome: TaskOutcome::Completed {
output: "child done".to_string(),
artifacts: vec![],
tool_trace: None,
},
});
let actions = scheduler.tick();
assert!(
actions.iter().any(|a| matches!(
a,
SchedulerAction::Done {
status: GraphStatus::Completed
}
)),
"graph should complete successfully, not deadlock"
);
assert_eq!(scheduler.graph.status, GraphStatus::Completed);
}
#[test]
fn test_failure_abort_cancels_running() {
let graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0, 1]),
]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.graph.tasks[1].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(1),
RunningTask {
agent_handle_id: "h1".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Failed {
error: "boom".to_string(),
},
};
scheduler.buffered_events.push_back(event);
let actions = scheduler.tick();
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
let cancel_ids: Vec<_> = actions
.iter()
.filter_map(|a| {
if let SchedulerAction::Cancel { agent_handle_id } = a {
Some(agent_handle_id.as_str())
} else {
None
}
})
.collect();
assert!(cancel_ids.contains(&"h1"), "task 1 should be canceled");
assert!(
actions
.iter()
.any(|a| matches!(a, SchedulerAction::Done { .. }))
);
}
#[test]
fn test_failure_skip_propagates() {
use crate::graph::FailureStrategy;
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(FailureStrategy::Skip);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Failed {
error: "skip me".to_string(),
},
};
scheduler.buffered_events.push_back(event);
scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Skipped);
assert_eq!(scheduler.graph.tasks[1].status, TaskStatus::Skipped);
}
#[test]
fn test_failure_retry_reschedules() {
use crate::graph::FailureStrategy;
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(FailureStrategy::Retry);
scheduler.graph.tasks[0].max_retries = Some(3);
scheduler.graph.tasks[0].retry_count = 0;
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Failed {
error: "transient".to_string(),
},
};
scheduler.buffered_events.push_back(event);
let actions = scheduler.tick();
let has_spawn = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(0)));
assert!(
has_spawn || scheduler.graph.tasks[0].status == TaskStatus::Ready,
"retry should produce spawn or Ready status"
);
assert_eq!(scheduler.graph.tasks[0].retry_count, 1);
}
#[test]
fn test_process_event_failed_retry() {
use crate::graph::FailureStrategy;
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(FailureStrategy::Retry);
scheduler.graph.tasks[0].max_retries = Some(2);
scheduler.graph.tasks[0].retry_count = 0;
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Failed {
error: "first failure".to_string(),
},
};
scheduler.buffered_events.push_back(event);
let actions = scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].retry_count, 1);
let spawned = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(0)));
assert!(
spawned || scheduler.graph.tasks[0].status == TaskStatus::Ready,
"retry should emit Spawn or set Ready"
);
assert_eq!(scheduler.graph.status, GraphStatus::Running);
}
#[test]
fn test_cascade_chain_threshold_preempts_recovery() {
let graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[1]),
]);
let mut config = make_config();
config.cascade_chain_threshold = 3;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].failure_strategy = Some(crate::graph::FailureStrategy::Retry);
scheduler.graph.tasks[0].max_retries = Some(5);
scheduler.graph.tasks[1].failure_strategy = Some(crate::graph::FailureStrategy::Retry);
scheduler.graph.tasks[1].max_retries = Some(5);
scheduler.graph.tasks[2].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("should never be applied".to_string()),
route_to: None,
});
for (id, handle) in [(TaskId(0), "h0"), (TaskId(1), "h1"), (TaskId(2), "h2")] {
scheduler.graph.tasks[id.index()].status = TaskStatus::Running;
scheduler.running.insert(
id,
RunningTask {
agent_handle_id: handle.to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
}
for (id, handle) in [(TaskId(0), "h0"), (TaskId(1), "h1"), (TaskId(2), "h2")] {
scheduler.buffered_events.push_back(TaskEvent {
task_id: id,
agent_handle_id: handle.to_string(),
outcome: TaskOutcome::Failed {
error: "boom".to_string(),
},
});
}
scheduler.tick();
assert_eq!(
scheduler.graph.status,
GraphStatus::Failed,
"cascade chain threshold must abort the graph"
);
assert_eq!(
scheduler.graph.tasks[2].status,
TaskStatus::Failed,
"the recovery-configured node must NOT be recovered — cascade-abort preempts \
propagate_failure() (and thus try_recover()) entirely"
);
assert_eq!(
scheduler.graph.tasks[2]
.result
.as_ref()
.map(|r| r.output.as_str()),
Some("boom"),
"result must hold the plain failure error, not a Mode-1 recovery substitution \
(agent_id/agent_def would also be set to the recovery marker if try_recover had run)"
);
assert_eq!(
scheduler.graph.tasks[2]
.result
.as_ref()
.and_then(|r| r.agent_def.as_deref()),
None,
"no synthetic recovery TaskResult (with its recovery marker agent_def) should ever \
have been set"
);
}
#[test]
fn test_timeout_cancels_stalled() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 1; let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_secs(2))
.unwrap(), admission_permit: None,
last_progress_at: None,
},
);
let actions = scheduler.tick();
let has_cancel = actions.iter().any(
|a| matches!(a, SchedulerAction::Cancel { agent_handle_id } if agent_handle_id == "h0"),
);
assert!(has_cancel, "timed-out task should emit Cancel action");
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
}
#[test]
fn test_per_task_timeout_override_fires_before_global_default() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300; let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Skip);
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: Some(1),
idle_timeout_secs: None,
});
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.graph.tasks[1].status = TaskStatus::Running;
let started_2s_ago = std::time::Instant::now()
.checked_sub(Duration::from_secs(2))
.unwrap();
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: started_2s_ago,
admission_permit: None,
last_progress_at: None,
},
);
scheduler.running.insert(
TaskId(1),
RunningTask {
agent_handle_id: "h1".to_string(),
agent_def_name: "worker".to_string(),
started_at: started_2s_ago,
admission_permit: None,
last_progress_at: None,
},
);
let actions = scheduler.tick();
let canceled: Vec<&str> = actions
.iter()
.filter_map(|a| match a {
SchedulerAction::Cancel { agent_handle_id } => Some(agent_handle_id.as_str()),
_ => None,
})
.collect();
assert_eq!(
canceled,
vec!["h0"],
"only the overridden task should time out; the other respects the longer global default"
);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Skipped);
assert_eq!(scheduler.graph.tasks[1].status, TaskStatus::Running);
}
#[test]
fn test_no_overrides_timing_matches_pre_feature_behavior() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 1;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_secs(2))
.unwrap(),
admission_permit: None,
last_progress_at: None,
},
);
let actions = scheduler.tick();
let has_cancel = actions.iter().any(
|a| matches!(a, SchedulerAction::Cancel { agent_handle_id } if agent_handle_id == "h0"),
);
assert!(
has_cancel,
"global timeout must still fire with no override"
);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
}
fn progress_handle_now() -> std::sync::Arc<std::sync::atomic::AtomicU64> {
std::sync::Arc::new(std::sync::atomic::AtomicU64::new(
zeph_common::monotonic_millis(),
))
}
#[test]
fn test_idle_timeout_fires_on_stale_progress() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300; let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: None,
idle_timeout_secs: Some(1), });
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(), admission_permit: None,
last_progress_at: Some(progress_handle_now()),
},
);
std::thread::sleep(Duration::from_millis(1_100));
let actions = scheduler.tick();
let has_cancel = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Cancel { .. }));
assert!(
has_cancel,
"stale progress heartbeat must fire idle timeout"
);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
let output = scheduler.graph.tasks[0]
.result
.as_ref()
.expect("timeout must populate TaskResult")
.output
.clone();
assert!(
output.contains("idle timeout"),
"TaskResult.output must name the idle cause, got: {output}"
);
}
#[test]
fn test_idle_timeout_does_not_fire_while_progress_continues() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: None,
idle_timeout_secs: Some(60), });
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: Some(progress_handle_now()),
},
);
let actions = scheduler.tick();
let has_cancel = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Cancel { .. }));
assert!(
!has_cancel,
"a task with a fresh heartbeat must not be idle-killed"
);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Running);
}
#[test]
fn test_idle_timeout_exempt_without_progress_handle() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300; let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: None,
idle_timeout_secs: Some(1), });
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_secs(5))
.unwrap(),
admission_permit: None,
last_progress_at: None, },
);
let actions = scheduler.tick();
let has_cancel = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Cancel { .. }));
assert!(
!has_cancel,
"a task with no progress handle must never be idle-killed, regardless of age"
);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Running);
}
#[test]
fn test_idle_timeout_multi_task_heartbeats_are_independent() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300; let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
let short_idle = crate::graph::TimeoutPolicy {
run_timeout_secs: None,
idle_timeout_secs: Some(1),
};
scheduler.graph.tasks[0].timeout = Some(short_idle.clone());
scheduler.graph.tasks[1].timeout = Some(short_idle);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Skip);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.graph.tasks[1].status = TaskStatus::Running;
let handle0 = progress_handle_now();
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: Some(handle0),
},
);
std::thread::sleep(Duration::from_millis(1_100));
let handle1 = progress_handle_now();
scheduler.running.insert(
TaskId(1),
RunningTask {
agent_handle_id: "h1".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: Some(handle1),
},
);
let actions = scheduler.tick();
let canceled: Vec<&str> = actions
.iter()
.filter_map(|a| match a {
SchedulerAction::Cancel { agent_handle_id } => Some(agent_handle_id.as_str()),
_ => None,
})
.collect();
assert_eq!(
canceled,
vec!["h0"],
"only the task with the stale heartbeat should be idle-killed"
);
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Skipped,
"task 0 (stale heartbeat) must be killed — Skip strategy turns the timeout-Failed \
status into Skipped via propagate_failure, same as the existing per-task-timeout test"
);
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Running,
"task 1's own fresh heartbeat must keep it alive, unaffected by task 0's staleness"
);
}
#[test]
fn test_run_timeout_wins_precedence_over_idle_on_same_tick() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 1; let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: None,
idle_timeout_secs: Some(1),
});
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_secs(10))
.unwrap(),
admission_permit: None,
last_progress_at: Some(progress_handle_now()),
},
);
scheduler.tick();
let output = scheduler.graph.tasks[0]
.result
.as_ref()
.expect("timeout must populate TaskResult")
.output
.clone();
assert!(
output.contains("run timeout"),
"run timeout must win precedence on a same-tick tie, got: {output}"
);
assert!(
!output.contains("idle timeout"),
"only one cause must be reported per NFR-OB-01, got: {output}"
);
}
#[test]
fn test_run_timeout_populates_task_result_cause() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 1;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_secs(2))
.unwrap(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.tick();
let result = scheduler.graph.tasks[0]
.result
.as_ref()
.expect("run-timeout kill must populate TaskResult (F5 fix)");
assert!(result.output.contains("run timeout"));
assert_eq!(result.agent_id.as_deref(), Some("h0"));
assert_eq!(result.agent_def.as_deref(), Some("worker"));
}
#[test]
fn test_effective_run_timeout_falls_back_to_global_default() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300;
let defs = vec![make_def("worker")];
let scheduler = DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
assert_eq!(
scheduler.effective_run_timeout(TaskId(0)),
Duration::from_mins(5)
);
}
#[test]
fn test_effective_run_timeout_uses_per_task_override() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut config = make_config();
config.task_timeout_secs = 300;
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: Some(45),
idle_timeout_secs: None,
});
assert_eq!(
scheduler.effective_run_timeout(TaskId(0)),
Duration::from_secs(45)
);
}
#[test]
fn test_effective_idle_timeout_none_when_unset_anywhere() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let defs = vec![make_def("worker")];
let scheduler =
DagScheduler::new(graph, &make_config(), Box::new(FirstRouter), defs, None).unwrap();
assert_eq!(scheduler.effective_idle_timeout(TaskId(0)), None);
}
#[test]
fn test_effective_idle_timeout_falls_back_to_global_default() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
default_idle_timeout_secs: Some(30),
..make_config()
};
let defs = vec![make_def("worker")];
let scheduler = DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
assert_eq!(
scheduler.effective_idle_timeout(TaskId(0)),
Some(Duration::from_secs(30))
);
}
#[test]
fn test_effective_idle_timeout_uses_per_task_override_over_global_default() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
default_idle_timeout_secs: Some(30),
..make_config()
};
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: None,
idle_timeout_secs: Some(5),
});
assert_eq!(
scheduler.effective_idle_timeout(TaskId(0)),
Some(Duration::from_secs(5)),
"per-task override must win over the global default (30s), not merge with it"
);
}
#[test]
fn test_cancel_all() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.graph.tasks[1].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(1),
RunningTask {
agent_handle_id: "h1".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let actions = scheduler.cancel_all();
assert_eq!(scheduler.graph.status, GraphStatus::Canceled);
assert!(scheduler.running.is_empty());
let cancel_count = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Cancel { .. }))
.count();
assert_eq!(cancel_count, 2);
assert!(actions.iter().any(|a| matches!(
a,
SchedulerAction::Done {
status: GraphStatus::Canceled
}
)));
}
#[test]
fn test_record_spawn_failure() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let error = SubAgentError::Spawn("spawn error".to_string());
let actions = scheduler.record_spawn_failure(TaskId(0), &error);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
assert!(
actions
.iter()
.any(|a| matches!(a, SchedulerAction::Done { .. }))
);
}
#[test]
fn test_record_spawn_failure_concurrency_limit_reverts_to_ready() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let error = SubAgentError::ConcurrencyLimit { active: 4, max: 4 };
let actions = scheduler.record_spawn_failure(TaskId(0), &error);
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Ready,
"task must revert to Ready so the next tick can retry"
);
assert_eq!(
scheduler.graph.status,
GraphStatus::Running,
"graph must stay Running, not transition to Failed"
);
assert!(
actions.is_empty(),
"no cancel or done actions expected for a transient deferral"
);
}
#[test]
fn test_record_spawn_failure_concurrency_limit_variant_spawn_for_task() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let error = SubAgentError::ConcurrencyLimit { active: 1, max: 1 };
let actions = scheduler.record_spawn_failure(TaskId(0), &error);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Ready);
assert!(actions.is_empty());
}
#[test]
fn test_concurrency_deferral_does_not_affect_running_task() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.graph.tasks[1].status = TaskStatus::Running;
let error = SubAgentError::ConcurrencyLimit { active: 1, max: 1 };
let actions = scheduler.record_spawn_failure(TaskId(1), &error);
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Running,
"task 0 must remain Running"
);
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Ready,
"task 1 must revert to Ready"
);
assert_eq!(
scheduler.graph.status,
GraphStatus::Running,
"graph must stay Running"
);
assert!(actions.is_empty(), "no cancel or done actions expected");
}
#[test]
fn test_max_concurrent_zero_no_infinite_loop() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
max_parallel: 0,
..make_config()
};
let mut scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
let actions1 = scheduler.tick();
assert!(
actions1
.iter()
.all(|a| !matches!(a, SchedulerAction::Spawn { .. })),
"no Spawn expected when max_parallel=0"
);
assert!(
actions1
.iter()
.all(|a| !matches!(a, SchedulerAction::Done { .. })),
"no Done(Failed) expected — ready tasks exist, so no deadlock"
);
assert_eq!(scheduler.graph.status, GraphStatus::Running);
let actions2 = scheduler.tick();
assert!(
actions2
.iter()
.all(|a| !matches!(a, SchedulerAction::Done { .. })),
"second tick must not emit Done(Failed) — ready tasks still exist"
);
assert_eq!(
scheduler.graph.status,
GraphStatus::Running,
"graph must remain Running"
);
}
#[test]
fn test_all_tasks_deferred_graph_stays_running() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
let actions = scheduler.tick();
assert_eq!(
actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.count(),
2,
"expected 2 Spawn actions on first tick"
);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Running);
assert_eq!(scheduler.graph.tasks[1].status, TaskStatus::Running);
let error = SubAgentError::ConcurrencyLimit { active: 2, max: 2 };
let r0 = scheduler.record_spawn_failure(TaskId(0), &error);
let r1 = scheduler.record_spawn_failure(TaskId(1), &error);
assert!(r0.is_empty() && r1.is_empty(), "no cancel/done on deferral");
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(scheduler.graph.tasks[1].status, TaskStatus::Ready);
assert_eq!(scheduler.graph.status, GraphStatus::Running);
let retry_actions = scheduler.tick();
let spawn_count = retry_actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.count();
assert!(
spawn_count > 0,
"second tick must re-emit Spawn for deferred tasks"
);
assert!(
retry_actions.iter().all(|a| !matches!(
a,
SchedulerAction::Done {
status: GraphStatus::Failed,
..
}
)),
"no Done(Failed) expected"
);
}
#[test]
fn test_no_agent_routes_inline() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler_with_router(graph, Box::new(NoneRouter));
let actions = scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Running);
assert!(
actions
.iter()
.any(|a| matches!(a, SchedulerAction::RunInline { .. }))
);
}
#[test]
fn test_stale_event_rejected() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "current-handle".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let stale_event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "old-handle".to_string(),
outcome: TaskOutcome::Completed {
output: "stale output".to_string(),
artifacts: vec![],
tool_trace: None,
},
};
scheduler.buffered_events.push_back(stale_event);
let actions = scheduler.tick();
assert_ne!(
scheduler.graph.tasks[0].status,
TaskStatus::Completed,
"stale event must not complete the task"
);
let has_done = actions
.iter()
.any(|a| matches!(a, SchedulerAction::Done { .. }));
assert!(
!has_done,
"no Done action should be emitted for a stale event"
);
assert!(
scheduler.running.contains_key(&TaskId(0)),
"running task must remain after stale event"
);
}
#[test]
fn test_duration_ms_computed_correctly() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_millis(50))
.unwrap(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Completed {
output: "result".to_string(),
artifacts: vec![],
tool_trace: None,
},
};
scheduler.buffered_events.push_back(event);
scheduler.tick();
let result = scheduler.graph.tasks[0].result.as_ref().unwrap();
assert!(
result.duration_ms > 0,
"duration_ms should be > 0, got {}",
result.duration_ms
);
}
#[test]
fn test_consecutive_spawn_failures_increments_on_concurrency_limit() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
assert_eq!(scheduler.consecutive_spawn_failures, 0, "starts at zero");
let error = SubAgentError::ConcurrencyLimit { active: 4, max: 4 };
scheduler.record_spawn_failure(TaskId(0), &error);
scheduler.record_batch_backoff(false, true);
assert_eq!(
scheduler.consecutive_spawn_failures, 1,
"first deferral tick: consecutive_spawn_failures must be 1"
);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.record_spawn_failure(TaskId(0), &error);
scheduler.record_batch_backoff(false, true);
assert_eq!(
scheduler.consecutive_spawn_failures, 2,
"second deferral tick: consecutive_spawn_failures must be 2"
);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.record_spawn_failure(TaskId(0), &error);
scheduler.record_batch_backoff(false, true);
assert_eq!(
scheduler.consecutive_spawn_failures, 3,
"third deferral tick: consecutive_spawn_failures must be 3"
);
}
#[test]
fn test_consecutive_spawn_failures_resets_on_success() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let error = SubAgentError::ConcurrencyLimit { active: 1, max: 1 };
scheduler.record_spawn_failure(TaskId(0), &error);
scheduler.record_batch_backoff(false, true);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.record_spawn_failure(TaskId(0), &error);
scheduler.record_batch_backoff(false, true);
assert_eq!(scheduler.consecutive_spawn_failures, 2);
scheduler.record_spawn(
TaskId(0),
"handle-0".to_string(),
"worker".to_string(),
None,
);
assert_eq!(
scheduler.consecutive_spawn_failures, 0,
"record_spawn must reset consecutive_spawn_failures to 0"
);
}
#[tokio::test]
async fn test_exponential_backoff_duration() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
deferral_backoff_ms: 50,
..make_config()
};
let mut scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
assert_eq!(scheduler.consecutive_spawn_failures, 0);
let start = tokio::time::Instant::now();
scheduler.wait_event().await;
let elapsed0 = start.elapsed();
assert!(
elapsed0.as_millis() >= 50,
"backoff with 0 deferrals must be >= base (50ms), got {}ms",
elapsed0.as_millis()
);
scheduler.consecutive_spawn_failures = 3;
let start = tokio::time::Instant::now();
scheduler.wait_event().await;
let elapsed3 = start.elapsed();
assert!(
elapsed3.as_millis() >= 400,
"backoff with 3 deferrals must be >= 400ms (50 * 8), got {}ms",
elapsed3.as_millis()
);
scheduler.consecutive_spawn_failures = 20;
let start = tokio::time::Instant::now();
scheduler.wait_event().await;
let elapsed_capped = start.elapsed();
assert!(
elapsed_capped.as_millis() >= 5000,
"backoff must be capped at 5000ms with high deferrals, got {}ms",
elapsed_capped.as_millis()
);
}
#[tokio::test]
async fn test_wait_event_nearest_deadline_reflects_per_task_override() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let config = zeph_config::OrchestrationConfig {
task_timeout_secs: 300,
..make_config()
};
let mut scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
scheduler.graph.tasks[0].timeout = Some(crate::graph::TimeoutPolicy {
run_timeout_secs: Some(1),
idle_timeout_secs: None,
});
let now = std::time::Instant::now();
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: now.checked_sub(Duration::from_millis(950)).unwrap(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.running.insert(
TaskId(1),
RunningTask {
agent_handle_id: "h1".to_string(),
agent_def_name: "worker".to_string(),
started_at: now,
admission_permit: None,
last_progress_at: None,
},
);
let start = tokio::time::Instant::now();
scheduler.wait_event().await;
let elapsed = start.elapsed();
assert!(
elapsed.as_millis() < 5000,
"wait_event must return promptly based on task 0's near-elapsed override, \
not the 300s global default; got {}ms",
elapsed.as_millis()
);
}
#[tokio::test]
async fn test_wait_event_sleeps_deferral_backoff_when_running_empty() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
deferral_backoff_ms: 50,
..make_config()
};
let mut scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
assert!(scheduler.running.is_empty());
let start = tokio::time::Instant::now();
scheduler.wait_event().await;
let elapsed = start.elapsed();
assert!(
elapsed.as_millis() >= 50,
"wait_event must sleep at least deferral_backoff (50ms) when running is empty, but only slept {}ms",
elapsed.as_millis()
);
}
#[test]
fn test_current_deferral_backoff_exponential_growth() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
deferral_backoff_ms: 250,
..make_config()
};
let mut scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
assert_eq!(
scheduler.current_deferral_backoff(),
Duration::from_millis(250)
);
scheduler.consecutive_spawn_failures = 1;
assert_eq!(
scheduler.current_deferral_backoff(),
Duration::from_millis(500)
);
scheduler.consecutive_spawn_failures = 2;
assert_eq!(scheduler.current_deferral_backoff(), Duration::from_secs(1));
scheduler.consecutive_spawn_failures = 3;
assert_eq!(scheduler.current_deferral_backoff(), Duration::from_secs(2));
scheduler.consecutive_spawn_failures = 4;
assert_eq!(scheduler.current_deferral_backoff(), Duration::from_secs(4));
scheduler.consecutive_spawn_failures = 5;
assert_eq!(scheduler.current_deferral_backoff(), Duration::from_secs(5));
scheduler.consecutive_spawn_failures = 100;
assert_eq!(scheduler.current_deferral_backoff(), Duration::from_secs(5));
}
#[test]
fn test_record_spawn_resets_consecutive_failures() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = DagScheduler::new(
graph,
&make_config(),
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
scheduler.consecutive_spawn_failures = 3;
let task_id = TaskId(0);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.record_spawn(task_id, "handle-1".into(), "worker".into(), None);
assert_eq!(scheduler.consecutive_spawn_failures, 0);
}
#[test]
fn test_record_spawn_failure_reverts_to_ready_no_counter_change() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = DagScheduler::new(
graph,
&make_config(),
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
assert_eq!(scheduler.consecutive_spawn_failures, 0);
let task_id = TaskId(0);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let error = SubAgentError::ConcurrencyLimit { active: 1, max: 1 };
scheduler.record_spawn_failure(task_id, &error);
assert_eq!(scheduler.consecutive_spawn_failures, 0);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Ready);
}
#[test]
fn test_parallel_dispatch_all_ready() {
let nodes: Vec<_> = (0..6).map(|i| make_node(i, &[])).collect();
let graph = graph_from_nodes(nodes);
let config = zeph_config::OrchestrationConfig {
max_parallel: 2,
..make_config()
};
let mut scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
let actions = scheduler.tick();
let spawn_count = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.count();
assert_eq!(
spawn_count, 2,
"only max_parallel=2 tasks dispatched per tick"
);
let running_count = scheduler
.graph
.tasks
.iter()
.filter(|t| t.status == TaskStatus::Running)
.count();
assert_eq!(running_count, 2, "only 2 tasks marked Running");
}
#[test]
fn test_batch_backoff_partial_success() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.consecutive_spawn_failures = 3;
scheduler.record_batch_backoff(true, true);
assert_eq!(
scheduler.consecutive_spawn_failures, 0,
"any success in batch must reset counter"
);
}
#[test]
fn test_batch_backoff_all_failed() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.consecutive_spawn_failures = 2;
scheduler.record_batch_backoff(false, true);
assert_eq!(
scheduler.consecutive_spawn_failures, 3,
"all-failure tick must increment counter"
);
}
#[test]
fn test_batch_backoff_no_spawns() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.consecutive_spawn_failures = 5;
scheduler.record_batch_backoff(false, false);
assert_eq!(
scheduler.consecutive_spawn_failures, 5,
"no spawns must not change counter"
);
}
#[test]
fn test_buffer_guard_uses_task_count() {
let nodes: Vec<_> = (0..10).map(|i| make_node(i, &[])).collect();
let graph = graph_from_nodes(nodes);
let config = zeph_config::OrchestrationConfig {
max_parallel: 2,
..make_config()
};
let scheduler = DagScheduler::new(
graph,
&config,
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
assert_eq!(scheduler.graph.tasks.len() * 2, 20);
assert_eq!(scheduler.max_parallel * 2, 4);
}
#[test]
fn test_batch_mixed_concurrency_and_fatal_failure() {
use crate::graph::FailureStrategy;
let mut nodes = vec![make_node(0, &[]), make_node(1, &[])];
nodes[1].failure_strategy = Some(FailureStrategy::Skip);
let graph = graph_from_nodes(nodes);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.graph.tasks[1].status = TaskStatus::Running;
let concurrency_err = SubAgentError::ConcurrencyLimit { active: 1, max: 1 };
let actions0 = scheduler.record_spawn_failure(TaskId(0), &concurrency_err);
assert!(
actions0.is_empty(),
"ConcurrencyLimit must produce no extra actions"
);
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Ready,
"task 0 must revert to Ready"
);
let fatal_err = SubAgentError::Spawn("provider unavailable".to_string());
let actions1 = scheduler.record_spawn_failure(TaskId(1), &fatal_err);
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Skipped,
"task 1: Skip strategy turns Failed into Skipped via propagate_failure"
);
assert!(
actions1
.iter()
.all(|a| !matches!(a, SchedulerAction::Done { .. })),
"no Done action expected: task 0 is still Ready"
);
scheduler.consecutive_spawn_failures = 0;
scheduler.record_batch_backoff(false, true);
assert_eq!(
scheduler.consecutive_spawn_failures, 1,
"batch with only ConcurrencyLimit must increment counter"
);
}
#[test]
fn test_deadlock_marks_non_terminal_tasks_canceled() {
let mut nodes = vec![make_node(0, &[]), make_node(1, &[0]), make_node(2, &[0])];
nodes[0].status = TaskStatus::Failed;
nodes[1].status = TaskStatus::Pending;
nodes[2].status = TaskStatus::Pending;
let mut graph = graph_from_nodes(nodes);
graph.status = GraphStatus::Failed;
let mut scheduler = DagScheduler::resume_from(
graph,
&make_config(),
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
let actions = scheduler.tick();
assert!(
actions.iter().any(|a| matches!(
a,
SchedulerAction::Done {
status: GraphStatus::Failed
}
)),
"deadlock must emit Done(Failed); got: {actions:?}"
);
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Canceled,
"Pending task must be Canceled on deadlock"
);
assert_eq!(
scheduler.graph.tasks[2].status,
TaskStatus::Canceled,
"Pending task must be Canceled on deadlock"
);
}
#[test]
fn test_deadlock_not_triggered_when_task_running() {
let mut nodes = vec![make_node(0, &[]), make_node(1, &[0])];
nodes[0].status = TaskStatus::Running;
nodes[0].assigned_agent = Some("handle-1".into());
nodes[1].status = TaskStatus::Pending;
let mut graph = graph_from_nodes(nodes);
graph.status = GraphStatus::Failed;
let mut scheduler = DagScheduler::resume_from(
graph,
&make_config(),
Box::new(FirstRouter),
vec![make_def("worker")],
None,
)
.unwrap();
let actions = scheduler.tick();
assert!(
actions
.iter()
.all(|a| !matches!(a, SchedulerAction::Done { .. })),
"no Done action expected when a task is running; got: {actions:?}"
);
assert_eq!(scheduler.graph.status, GraphStatus::Running);
}
fn make_def_with_provider(name: &str, provider: &str) -> zeph_subagent::SubAgentDef {
let mut d = zeph_subagent::SubAgentDef::for_test(name);
d.model = Some(zeph_subagent::ModelSpec::Named(provider.to_string()));
d
}
#[test]
fn admission_gate_saturated_defers_task() {
let gate = crate::admission::AdmissionGate::new(&[("quality".to_string(), 1usize)]);
let _held_permit = gate
.try_acquire("quality")
.expect("first permit must succeed");
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = make_config();
let defs = vec![make_def_with_provider("worker", "quality")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, Some(gate)).unwrap();
let actions = scheduler.tick();
let spawn_count = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.count();
assert_eq!(
spawn_count, 0,
"saturated gate must defer task — no Spawn emitted"
);
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Ready,
"deferred task must stay Ready"
);
}
#[test]
fn admission_gate_permit_transferred_to_running() {
let gate = crate::admission::AdmissionGate::new(&[("quality".to_string(), 2usize)]);
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = make_config();
let defs = vec![make_def_with_provider("worker", "quality")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, Some(gate)).unwrap();
let actions = scheduler.tick();
let spawned = actions.iter().find_map(|a| {
if let SchedulerAction::Spawn { task_id, .. } = a {
Some(*task_id)
} else {
None
}
});
let task_id = spawned.expect("task must be spawned");
assert!(
scheduler.pending_permits.contains_key(&task_id),
"permit must be in pending_permits after dispatch"
);
scheduler.graph.tasks[task_id.index()].status = TaskStatus::Running;
scheduler.record_spawn(task_id, "handle-1".into(), "worker".into(), None);
assert!(
!scheduler.pending_permits.contains_key(&task_id),
"pending_permits must be empty after record_spawn"
);
assert!(
scheduler.running[&task_id].admission_permit.is_some(),
"admission_permit must be set in RunningTask"
);
}
#[test]
fn admission_gate_bypass_for_ungated_provider() {
let gate = crate::admission::AdmissionGate::new(&[("quality".to_string(), 1usize)]);
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = make_config();
let defs = vec![make_def_with_provider("worker", "fast")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, Some(gate)).unwrap();
let actions = scheduler.tick();
let spawn_count = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { .. }))
.count();
assert_eq!(spawn_count, 1, "ungated provider must not be blocked");
}
#[test]
fn record_spawn_failure_releases_pending_permit() {
let gate = crate::admission::AdmissionGate::new(&[("quality".to_string(), 2usize)]);
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = make_config();
let defs = vec![make_def_with_provider("worker", "quality")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, Some(gate)).unwrap();
let permit = scheduler
.admission_gate
.as_ref()
.unwrap()
.try_acquire("quality")
.expect("permit must be available");
let task_id = TaskId(0);
scheduler.pending_permits.insert(task_id, permit);
scheduler.graph.tasks[0].status = TaskStatus::Running;
assert!(scheduler.pending_permits.contains_key(&task_id));
let fatal = zeph_subagent::SubAgentError::Spawn("provider unavailable".to_string());
scheduler.record_spawn_failure(task_id, &fatal);
assert!(
!scheduler.pending_permits.contains_key(&task_id),
"pending permit must be removed after fatal spawn failure"
);
}
#[test]
fn graph_dirty_clear_at_construction() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let scheduler = make_scheduler(graph);
assert!(
!scheduler.graph_dirty,
"graph_dirty must be false immediately after construction"
);
}
#[test]
fn take_graph_dirty_returns_false_when_no_mutations() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must return false when no mutations occurred"
);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must remain false after a second call"
);
}
#[test]
fn take_graph_dirty_true_after_task_completes() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Completed {
output: "done".to_string(),
artifacts: vec![],
tool_trace: None,
},
};
scheduler.buffered_events.push_back(event);
scheduler.tick();
assert!(
scheduler.take_graph_dirty(),
"take_graph_dirty must return true after a task completes"
);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must return false on the second call (reset invariant)"
);
}
#[test]
fn take_graph_dirty_true_after_task_fails() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
let event = TaskEvent {
task_id: TaskId(0),
agent_handle_id: "h0".to_string(),
outcome: TaskOutcome::Failed {
error: "boom".to_string(),
},
};
scheduler.buffered_events.push_back(event);
scheduler.tick();
assert!(
scheduler.take_graph_dirty(),
"take_graph_dirty must return true after a task fails"
);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must return false on the second call (reset invariant)"
);
}
#[test]
fn take_graph_dirty_true_after_fatal_spawn_failure() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let fatal = zeph_subagent::SubAgentError::Spawn("provider gone".to_string());
scheduler.record_spawn_failure(TaskId(0), &fatal);
assert!(
scheduler.take_graph_dirty(),
"take_graph_dirty must return true after a fatal spawn failure marks task Failed"
);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must reset to false on second call"
);
}
#[test]
fn take_graph_dirty_false_after_transient_concurrency_failure() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
let transient = zeph_subagent::SubAgentError::ConcurrencyLimit { active: 1, max: 1 };
scheduler.record_spawn_failure(TaskId(0), &transient);
assert!(
!scheduler.take_graph_dirty(),
"transient concurrency deferral must not set graph_dirty (no terminal mutation)"
);
}
#[test]
fn take_graph_dirty_true_after_cancel_all() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.cancel_all();
assert!(
scheduler.take_graph_dirty(),
"take_graph_dirty must return true after cancel_all"
);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must reset to false on second call"
);
}
#[test]
fn take_graph_dirty_true_after_timeout() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
task_timeout_secs: 1,
..make_config()
};
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
scheduler.graph.tasks[0].status = TaskStatus::Running;
scheduler.running.insert(
TaskId(0),
RunningTask {
agent_handle_id: "h0".to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now()
.checked_sub(Duration::from_secs(2))
.unwrap(),
admission_permit: None,
last_progress_at: None,
},
);
scheduler.tick();
assert!(
scheduler.take_graph_dirty(),
"take_graph_dirty must return true after a task times out"
);
assert!(
!scheduler.take_graph_dirty(),
"take_graph_dirty must reset to false on second call"
);
}
fn make_running_task(scheduler: &mut DagScheduler, task_id: TaskId, handle_id: &str) {
scheduler.graph.tasks[task_id.index()].status = TaskStatus::Running;
scheduler.running.insert(
task_id,
RunningTask {
agent_handle_id: handle_id.to_string(),
agent_def_name: "worker".to_string(),
started_at: std::time::Instant::now(),
admission_permit: None,
last_progress_at: None,
},
);
}
fn completed_event(
task_id: TaskId,
handle_id: &str,
tool_trace: Option<Vec<ToolCallSummary>>,
) -> TaskEvent {
TaskEvent {
task_id,
agent_handle_id: handle_id.to_string(),
outcome: TaskOutcome::Completed {
output: "narration".to_string(),
artifacts: vec![],
tool_trace,
},
}
}
#[test]
fn handle_completed_outcome_all_tools_failed_marks_task_failed() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler.buffered_events.push_back(completed_event(
TaskId(0),
"h0",
Some(vec![
ToolCallSummary {
tool: "create_directory".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
]),
));
scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Failed,
"a task whose every tool call failed must be marked Failed, not Completed"
);
}
#[test]
fn handle_completed_outcome_mixed_trace_write_failed_marks_task_failed() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler.buffered_events.push_back(completed_event(
TaskId(0),
"h0",
Some(vec![
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
]),
));
scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Failed,
"a mixed trace where every write-type call failed must be corrected to Failed, \
regardless of a successful read-type call"
);
}
#[test]
fn handle_completed_outcome_mixed_trace_write_succeeded_preserves_completed() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler.buffered_events.push_back(completed_event(
TaskId(0),
"h0",
Some(vec![
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: true,
is_read_only: false,
},
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
},
]),
));
scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Completed,
"a mixed trace where the write-type call succeeded must preserve Completed even if \
a read-type call failed"
);
}
#[test]
fn handle_completed_outcome_empty_tool_trace_preserves_completed() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler
.buffered_events
.push_back(completed_event(TaskId(0), "h0", Some(vec![])));
scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Completed,
"an empty tool trace (no tool calls made) must preserve Completed"
);
}
#[test]
fn handle_completed_outcome_none_tool_trace_preserves_completed() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler
.buffered_events
.push_back(completed_event(TaskId(0), "h0", None));
scheduler.tick();
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Completed,
"tool_trace: None must preserve Completed (no synchronous trace available)"
);
}
#[test]
fn handle_completed_outcome_all_tools_failed_does_not_double_count_cascade() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
cascade_routing: true,
topology_selection: true,
..make_config()
};
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
assert!(
scheduler.cascade_detector.is_some(),
"test precondition: cascade_detector must be enabled"
);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler.buffered_events.push_back(completed_event(
TaskId(0),
"h0",
Some(vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}]),
));
scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
let health = scheduler
.cascade_detector
.as_ref()
.unwrap()
.region_health()
.get(&TaskId(0))
.expect("region health must be recorded for this task's region");
assert_eq!(
health.total_tasks, 1,
"task must be recorded exactly once in RegionHealth, not double-counted \
(success then failure): {health:?}"
);
assert_eq!(
health.failed_tasks, 1,
"the single recorded outcome must be a failure"
);
}
fn scheduler_with_completed_task() -> (DagScheduler, TaskId) {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
let task_id = TaskId(0);
scheduler.graph.tasks[task_id.index()].status = TaskStatus::Completed;
scheduler.graph.tasks[task_id.index()].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-1".to_string()),
agent_def: Some("worker".to_string()),
});
let _ = scheduler.take_graph_dirty(); (scheduler, task_id)
}
#[test]
fn correct_completed_to_failed_noop_on_none_trace() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let corrected = scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, None);
assert!(!corrected);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Completed
);
assert!(!scheduler.take_graph_dirty());
}
#[test]
fn correct_completed_to_failed_noop_on_all_ok_trace() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(!corrected);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Completed
);
assert!(!scheduler.take_graph_dirty());
}
#[test]
fn correct_completed_to_failed_all_read_only_calls_failed_still_corrects() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
corrected,
"a trace with no write-type calls must fall back to the 'every call failed' rule"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
}
#[test]
fn correct_completed_to_failed_mixed_trace_write_failed_corrects() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
corrected,
"a mixed trace where the only write-type call failed must correct to Failed"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
}
#[test]
fn correct_completed_to_failed_mixed_trace_write_succeeded_preserves_completed() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
},
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: true,
is_read_only: false,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
!corrected,
"a successful write-type call must preserve Completed even with a failed read"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Completed
);
}
#[test]
fn correct_completed_to_failed_partial_write_success_preserves_completed() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: true,
is_read_only: false,
},
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
!corrected,
"partial write success (one write ok, one write failed) must preserve Completed"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Completed
);
}
#[test]
fn correct_completed_to_failed_blocked_fetch_alongside_successful_read_corrects() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
ToolCallSummary {
tool: "fetch".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
corrected,
"a blocked fetch (quarantine-denied despite is_read_only == true) alongside a \
successful plain read must still correct to Failed"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
}
#[test]
fn correct_completed_to_failed_blocked_invoke_skill_alongside_successful_read_corrects() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "list_directory".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
ToolCallSummary {
tool: "invoke_skill".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
corrected,
"a blocked invoke_skill alongside a successful list_directory must correct to Failed"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
}
#[test]
fn correct_completed_to_failed_successful_fetch_alongside_failed_read_preserves_completed() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
},
ToolCallSummary {
tool: "fetch".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
!corrected,
"a successful fetch (quarantine-denied-class call) must preserve Completed even with \
a failed plain read"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Completed
);
}
#[test]
fn correct_completed_to_failed_plain_read_only_trace_all_failed_still_corrects() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![ToolCallSummary {
tool: "list_directory".to_string(),
args_summary: None,
ok: false,
is_read_only: true,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
corrected,
"a trace with zero counting calls (plain, non-denied reads only) must fall back to \
the 'every call failed' rule"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
}
#[test]
fn correct_completed_to_failed_noop_on_empty_trace() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&[]));
assert!(!corrected);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Completed
);
}
#[test]
fn correct_completed_to_failed_noop_when_task_not_completed() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
scheduler.graph.tasks[task_id.index()].status = TaskStatus::Failed;
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(
!corrected,
"correction must no-op when the task's current status is not Completed"
);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
}
#[test]
fn correct_completed_to_failed_success_flips_status_and_marks_output() {
let (mut scheduler, task_id) = scheduler_with_completed_task();
let trace = vec![
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
ToolCallSummary {
tool: "bash".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(task_id, Some(&trace));
assert!(corrected);
assert_eq!(
scheduler.graph.tasks[task_id.index()].status,
TaskStatus::Failed
);
let output = &scheduler.graph.tasks[task_id.index()]
.result
.as_ref()
.unwrap()
.output;
assert!(
output.starts_with("original output"),
"correction must append to, not replace, the original output: {output}"
);
assert!(
output.contains("corrected") && output.contains('2'),
"correction marker must be present and mention the failed call count: {output}"
);
assert!(
scheduler.take_graph_dirty(),
"a genuine status correction must set graph_dirty"
);
}
#[test]
fn propagate_corrected_task_failure_cancels_running_dependent_and_fails_graph() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
make_running_task(&mut scheduler, TaskId(1), "h1");
let _ = scheduler.take_graph_dirty();
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected, "test precondition: task 0 must be corrected");
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Cancel { agent_handle_id } if agent_handle_id == "h1")
),
"the already-Running dependent must be canceled: {actions:?}"
);
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Done { status } if *status == GraphStatus::Failed)
),
"GraphStatus leaving Running must emit a Done action: {actions:?}"
);
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
assert!(scheduler.graph.finished_at.is_some());
assert!(
!scheduler.running.contains_key(&TaskId(1)),
"the canceled dependent must be removed from the running map"
);
}
#[test]
fn propagate_corrected_task_failure_cancels_pending_commanded_from_target() {
let graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[1]),
]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "claimed done, handed off".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
scheduler.graph.tasks[1].status = TaskStatus::Ready;
scheduler.graph.tasks[1].commanded_from = Some(TaskId(0));
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected, "test precondition: task 0 must be corrected");
let _ = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Skipped,
"a Pending/Ready commanded_from target must be cancelled when its source is \
corrected to Failed post-hoc"
);
assert_eq!(
scheduler.graph.tasks[2].status,
TaskStatus::Skipped,
"a depends_on dependent of the cancelled commanded_from target must also be \
skipped, or it would strand in Pending forever"
);
}
#[test]
fn propagate_corrected_task_failure_leaves_running_commanded_from_target_untouched() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Skip);
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "claimed done, handed off".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
scheduler.graph.tasks[1].commanded_from = Some(TaskId(0));
make_running_task(&mut scheduler, TaskId(1), "h1");
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected, "test precondition: task 0 must be corrected");
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Running,
"an already-Running commanded_from target is out of this fix's scope"
);
assert!(
!actions.iter().any(
|a| matches!(a, SchedulerAction::Cancel { agent_handle_id } if agent_handle_id == "h1")
),
"a Running commanded_from target must not be cancelled by this fix: {actions:?}"
);
}
#[test]
fn propagate_corrected_task_failure_does_not_double_count_cascade() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let config = zeph_config::OrchestrationConfig {
cascade_routing: true,
topology_selection: true,
..make_config()
};
let defs = vec![make_def("worker")];
let mut scheduler =
DagScheduler::new(graph, &config, Box::new(FirstRouter), defs, None).unwrap();
assert!(
scheduler.cascade_detector.is_some(),
"test precondition: cascade_detector must be enabled"
);
make_running_task(&mut scheduler, TaskId(0), "h0");
scheduler
.buffered_events
.push_back(completed_event(TaskId(0), "h0", None));
scheduler.tick();
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Completed);
let region_health_before = scheduler
.cascade_detector
.as_ref()
.unwrap()
.region_health()
.get(&TaskId(0))
.expect("region health must be recorded")
.clone();
assert_eq!(region_health_before.total_tasks, 1);
assert_eq!(region_health_before.failed_tasks, 0);
let trace = vec![
ToolCallSummary {
tool: "read".to_string(),
args_summary: None,
ok: true,
is_read_only: true,
},
ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
},
];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let _ = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
let region_health_after = scheduler
.cascade_detector
.as_ref()
.unwrap()
.region_health()
.get(&TaskId(0))
.expect("region health must still be present")
.clone();
assert_eq!(
region_health_after.total_tasks, region_health_before.total_tasks,
"propagate_corrected_task_failure must not change RegionHealth: {region_health_after:?}"
);
assert_eq!(
region_health_after.failed_tasks, region_health_before.failed_tasks,
"propagate_corrected_task_failure must not re-record this task as a failure: \
{region_health_after:?}"
);
}
#[test]
fn propagate_corrected_task_failure_noop_when_no_dependents_and_graph_stays_running() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Skip);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert!(
actions.is_empty(),
"Skip strategy with no running dependents and a still-Running graph must emit no \
actions: {actions:?}"
);
assert_eq!(
scheduler.graph.status,
GraphStatus::Running,
"GraphStatus must remain Running while task 1 is still non-terminal"
);
}
#[test]
fn propagate_corrected_task_failure_retry_strategy_forces_terminal_failure() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Retry);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
assert_eq!(scheduler.graph.tasks[0].retry_count, 0);
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected, "test precondition: task 0 must be corrected");
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Failed,
"Retry must NOT resurrect an already-completed corrected task to Ready -- that would \
trigger a redundant redispatch and double-count RegionHealth (critic finding S1)"
);
assert_eq!(
scheduler.graph.tasks[0].retry_count, 0,
"no retry attempt was made, so retry_count must not increment"
);
assert_eq!(
scheduler.graph.status,
GraphStatus::Failed,
"Retry is forced to terminal Abort-style behavior for a post-hoc correction"
);
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Done { status } if *status == GraphStatus::Failed)
),
"GraphStatus leaving Running must emit a Done action: {actions:?}"
);
}
#[test]
fn propagate_corrected_task_failure_retry_then_tick_emits_no_spawn() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Retry);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let _ = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(scheduler.graph.tasks[0].status, TaskStatus::Failed);
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
let actions = scheduler.tick();
let spawn_count = actions
.iter()
.filter(|a| matches!(a, SchedulerAction::Spawn { task_id, .. } if *task_id == TaskId(0)))
.count();
assert_eq!(
spawn_count, 0,
"a task forced to terminal Failed under Retry must never be redispatched: {actions:?}"
);
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Done { status } if *status == GraphStatus::Failed)
),
"tick() on an already-terminal graph must return Done: {actions:?}"
);
}
#[test]
fn propagate_corrected_task_failure_ask_strategy_forces_terminal_failure_instead_of_pause() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Ask);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.status,
GraphStatus::Failed,
"Ask must not pause the graph for a post-hoc correction -- it is forced to terminal \
Failed instead (critic finding S1)"
);
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Done { status } if *status == GraphStatus::Failed)
),
"GraphStatus leaving Running must emit a Done action: {actions:?}"
);
}
#[test]
fn propagate_corrected_task_failure_retry_with_state_injection_does_not_recover() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Retry);
scheduler.graph.tasks[0].max_retries = Some(3);
scheduler.graph.tasks[0].retry_count = 3; scheduler.graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let _ = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Failed,
"the forced-terminal path must NOT invoke Mode-1 recovery -- a configured \
state_injection fallback must not silently flip the just-corrected Failed status \
back to Completed, which no already-unblocked dependent would ever see"
);
assert_ne!(
scheduler.graph.tasks[0]
.result
.as_ref()
.unwrap()
.agent_def
.as_deref(),
Some("__recovery__"),
"the recovery marker must never appear -- confirms try_recover was not invoked"
);
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
}
#[test]
fn propagate_corrected_task_failure_abort_with_state_injection_does_not_recover() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
assert_eq!(
scheduler.graph.default_failure_strategy,
zeph_config::FailureStrategy::Abort,
"test precondition: default strategy must be Abort (unconfigured failure_strategy)"
);
scheduler.graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Failed,
"default Abort strategy + state_injection must NOT invoke Mode-1 recovery on a \
post-hoc correction -- the injected fallback would never be seen by dependents that \
already consumed the original (now-invalidated) output"
);
assert_ne!(
scheduler.graph.tasks[0]
.result
.as_ref()
.unwrap()
.agent_def
.as_deref(),
Some("__recovery__"),
"the recovery marker must never appear -- confirms try_recover was not invoked"
);
assert_eq!(scheduler.graph.status, GraphStatus::Failed);
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Done { status } if *status == GraphStatus::Failed)
),
"GraphStatus leaving Running must emit a Done action: {actions:?}"
);
}
#[test]
fn propagate_corrected_task_failure_skip_with_state_injection_still_uses_generic_propagate() {
let graph = graph_from_nodes(vec![make_node(0, &[])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].failure_strategy = Some(zeph_config::FailureStrategy::Skip);
scheduler.graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let _ = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[0].status,
TaskStatus::Skipped,
"Skip's arm must run unmodified (flips to Skipped, not left Failed) even with \
state_injection configured -- Skip never calls try_recover, so it needs no \
forced-terminal special-case"
);
}
#[test]
fn propagate_corrected_task_failure_leaves_already_completed_dependent_untouched() {
let graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
let mut scheduler = make_scheduler(graph);
scheduler.graph.tasks[0].status = TaskStatus::Completed;
scheduler.graph.tasks[0].result = Some(TaskResult {
output: "original output".to_string(),
artifacts: vec![],
duration_ms: 10,
agent_id: Some("agent-0".to_string()),
agent_def: Some("worker".to_string()),
});
scheduler.graph.tasks[1].status = TaskStatus::Completed;
scheduler.graph.tasks[1].result = Some(TaskResult {
output: "dependent already finished using task 0's now-invalidated output".to_string(),
artifacts: vec![],
duration_ms: 5,
agent_id: Some("agent-1".to_string()),
agent_def: Some("worker".to_string()),
});
let trace = vec![ToolCallSummary {
tool: "write".to_string(),
args_summary: None,
ok: false,
is_read_only: false,
}];
let corrected =
scheduler.correct_completed_to_failed_if_all_tool_calls_failed(TaskId(0), Some(&trace));
assert!(corrected);
let actions = scheduler.propagate_corrected_task_failure(TaskId(0));
assert_eq!(
scheduler.graph.tasks[1].status,
TaskStatus::Completed,
"a dependent that already completed before the correction landed is not \
retroactively unwound (documented limitation)"
);
assert!(
actions.iter().any(
|a| matches!(a, SchedulerAction::Done { status } if *status == GraphStatus::Failed)
),
"the graph must still terminalize as Failed even though the dependent stayed \
Completed -- task 0 itself is Failed and no non-terminal task remains: {actions:?}"
);
}