use std::collections::VecDeque;
use zeph_common::fidelity::PlannedToolHint;
use super::command::TaskRef;
use super::error::OrchestrationError;
use super::graph::PredicateOutcome;
use super::graph::{
FailureStrategy, GraphStatus, TaskGraph, TaskId, TaskNode, TaskResult, TaskStatus,
};
#[must_use = "validation result must be checked"]
pub fn validate(
tasks: &[TaskNode],
max_tasks: usize,
default_failure_strategy: FailureStrategy,
) -> Result<(), OrchestrationError> {
if tasks.len() > max_tasks {
return Err(OrchestrationError::InvalidGraph(format!(
"graph has {} tasks, exceeding the limit of {max_tasks}",
tasks.len()
)));
}
if tasks.is_empty() {
return Err(OrchestrationError::InvalidGraph(
"graph has no tasks".to_string(),
));
}
let mut route_to_target_counts: std::collections::HashMap<TaskId, usize> =
std::collections::HashMap::new();
for (i, task) in tasks.iter().enumerate() {
let expected = u32::try_from(i).map_err(|_| {
OrchestrationError::InvalidGraph(format!("task index {i} overflows u32"))
})?;
if task.id != TaskId(expected) {
return Err(OrchestrationError::InvalidGraph(format!(
"task at index {i} has id {task_id} (expected {i})",
task_id = task.id
)));
}
for dep in &task.depends_on {
if *dep == task.id {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} has a self-reference"
)));
}
if dep.index() >= tasks.len() {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} references non-existent task {dep}"
)));
}
}
if task.recovery.is_some() && task.verify_predicate.is_some() {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} sets both recovery and verify_predicate — a predicate-gated \
task must not be recovery-eligible"
)));
}
if let Some(recovery) = &task.recovery {
if recovery.state_injection.is_some() && recovery.route_to.is_some() {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} sets both recovery.state_injection and recovery.route_to — \
Mode 1 and Mode 2 recovery are mutually exclusive"
)));
}
validate_route_to(
i,
task,
recovery,
tasks,
default_failure_strategy,
&mut route_to_target_counts,
)?;
let effective_strategy = task.failure_strategy.unwrap_or(default_failure_strategy);
if recovery.state_injection.is_some()
&& matches!(
effective_strategy,
FailureStrategy::Skip | FailureStrategy::Ask
)
{
tracing::warn!(
task_index = i,
strategy = ?effective_strategy,
"recovery configured but effective failure strategy is Skip/Ask — \
recovery is inert"
);
}
}
}
for (target, count) in route_to_target_counts {
if count > 1 {
return Err(OrchestrationError::InvalidGraph(format!(
"task {target} is the recovery.route_to target of {count} sources — \
exactly one source per target is required (N:1 shared fallback is deferred)"
)));
}
}
let sorted = toposort(tasks)?;
let has_root = tasks.iter().any(|t| t.depends_on.is_empty());
if !has_root {
return Err(OrchestrationError::CycleDetected);
}
let _ = sorted;
Ok(())
}
fn validate_route_to(
i: usize,
task: &TaskNode,
recovery: &crate::graph::RecoveryAction,
tasks: &[TaskNode],
default_failure_strategy: FailureStrategy,
route_to_target_counts: &mut std::collections::HashMap<TaskId, usize>,
) -> Result<(), OrchestrationError> {
let Some(target) = recovery.route_to else {
return Ok(());
};
if target == task.id {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} sets recovery.route_to to itself — self-reroute is not allowed"
)));
}
if target.index() >= tasks.len() {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} sets recovery.route_to to non-existent task {target}"
)));
}
let target_task = &tasks[target.index()];
if !target_task.depends_on.is_empty() {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} routes to task {target}, but {target} has a non-empty depends_on \
— a route_to target must only become ready via on-failure activation"
)));
}
if target_task
.recovery
.as_ref()
.is_some_and(|r| r.route_to.is_some())
{
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} routes to task {target}, but {target} itself sets recovery.route_to \
— chained `route_to` is not supported in v1"
)));
}
*route_to_target_counts.entry(target).or_insert(0) += 1;
let effective_strategy = task.failure_strategy.unwrap_or(default_failure_strategy);
if matches!(
effective_strategy,
FailureStrategy::Skip | FailureStrategy::Ask
) {
return Err(OrchestrationError::InvalidGraph(format!(
"task {i} sets recovery.route_to but its effective failure strategy is \
{effective_strategy:?} — route_to must never coincide with the Skip-BFS \
or Ask-pause arms"
)));
}
Ok(())
}
pub fn toposort(tasks: &[TaskNode]) -> Result<Vec<TaskId>, OrchestrationError> {
let n = tasks.len();
let mut in_degree = vec![0u32; n];
for task in tasks {
in_degree[task.id.index()] = u32::try_from(task.depends_on.len()).map_err(|_| {
OrchestrationError::InvalidGraph("dependency count overflows u32".to_string())
})?;
}
let mut queue: VecDeque<TaskId> = in_degree
.iter()
.enumerate()
.filter(|(_, d)| **d == 0)
.map(|(i, _)| u32::try_from(i).map(TaskId))
.collect::<Result<_, _>>()
.map_err(|_| OrchestrationError::InvalidGraph("task index overflows u32".to_string()))?;
let mut dependents: Vec<Vec<TaskId>> = vec![Vec::new(); n];
for task in tasks {
for dep in &task.depends_on {
dependents[dep.index()].push(task.id);
}
}
let mut order = Vec::with_capacity(n);
while let Some(id) = queue.pop_front() {
order.push(id);
for &dep_id in &dependents[id.index()] {
in_degree[dep_id.index()] -= 1;
if in_degree[dep_id.index()] == 0 {
queue.push_back(dep_id);
}
}
}
if order.len() != n {
return Err(OrchestrationError::CycleDetected);
}
Ok(order)
}
fn all_parents_predicate_clear(task: &TaskNode, graph: &TaskGraph) -> bool {
task.depends_on.iter().all(|parent_id| {
let parent = &graph.tasks[parent_id.index()];
matches!(
(&parent.verify_predicate, &parent.predicate_outcome),
(None, _)
| (Some(_), Some(PredicateOutcome { passed: true, .. }))
)
})
}
#[must_use]
pub fn ready_tasks(graph: &TaskGraph) -> Vec<TaskId> {
graph
.tasks
.iter()
.filter_map(|task| {
match task.status {
TaskStatus::Ready => {
if all_parents_predicate_clear(task, graph) {
Some(task.id)
} else {
None
}
}
TaskStatus::Pending => {
let all_deps_done = task
.depends_on
.iter()
.all(|dep_id| graph.tasks[dep_id.index()].status == TaskStatus::Completed);
if all_deps_done && all_parents_predicate_clear(task, graph) {
Some(task.id)
} else {
None
}
}
_ => None,
}
})
.collect()
}
fn try_recover(graph: &mut TaskGraph, failed_id: TaskId) -> bool {
let Some(injection) = graph.tasks[failed_id.index()]
.recovery
.as_ref()
.and_then(|r| r.state_injection.clone())
else {
return false;
};
let node = &mut graph.tasks[failed_id.index()];
node.status = TaskStatus::Completed;
node.result = Some(TaskResult {
output: injection,
artifacts: Vec::new(),
duration_ms: 0,
agent_id: None,
agent_def: Some("__recovery__".to_string()),
});
tracing::info!(
task_id = %failed_id,
"orchestration.dag.recover_task: Mode-1 recovery applied"
);
true
}
pub(crate) fn mark_dormant_route_to_targets(graph: &mut TaskGraph) {
let targets: Vec<TaskId> = graph
.tasks
.iter()
.filter_map(|t| t.recovery.as_ref().and_then(|r| r.route_to))
.collect();
for target in targets {
let node = &mut graph.tasks[target.index()];
if node.status == TaskStatus::Pending {
node.status = TaskStatus::Dormant;
}
}
}
fn try_reroute(graph: &mut TaskGraph, failed_id: TaskId) -> bool {
let Some(target) = graph.tasks[failed_id.index()]
.recovery
.as_ref()
.and_then(|r| r.route_to)
else {
return false;
};
if graph.tasks[target.index()].status != TaskStatus::Dormant {
tracing::warn!(
task_id = %failed_id,
target = %target,
target_status = %graph.tasks[target.index()].status,
"orchestration.dag.try_reroute: route_to target is not Dormant, skipping activation"
);
return false;
}
let node = &mut graph.tasks[target.index()];
node.status = TaskStatus::Ready;
node.routed_from = Some(failed_id);
tracing::info!(
task_id = %failed_id,
target = %target,
"orchestration.dag.try_reroute: Mode-2 reroute activated"
);
true
}
fn resolve_task_ref(graph: &TaskGraph, goto: &TaskRef) -> Result<TaskId, OrchestrationError> {
match goto {
TaskRef::ById(id) => {
if id.index() >= graph.tasks.len() {
return Err(OrchestrationError::InvalidHandoffTarget(format!(
"goto references non-existent task {id}"
)));
}
Ok(*id)
}
TaskRef::ByTitle(title) => {
let mut matches = graph.tasks.iter().filter(|t| &t.title == title);
let Some(first) = matches.next() else {
return Err(OrchestrationError::InvalidHandoffTarget(format!(
"goto title {title:?} does not match any task"
)));
};
if matches.next().is_some() {
return Err(OrchestrationError::InvalidHandoffTarget(format!(
"goto title {title:?} matches more than one task — ambiguous"
)));
}
Ok(first.id)
}
}
}
pub fn try_handoff(
graph: &mut TaskGraph,
source: TaskId,
goto: &TaskRef,
max_handoffs: u32,
) -> Result<TaskId, OrchestrationError> {
if graph.handoff_count >= max_handoffs {
return Err(OrchestrationError::HandoffBudgetExhausted {
handoff_count: graph.handoff_count,
max_handoffs,
});
}
let target = validate_handoff_target(graph, source, goto)?;
let node = &mut graph.tasks[target.index()];
if matches!(node.status, TaskStatus::Dormant | TaskStatus::Pending) {
node.status = TaskStatus::Ready;
}
node.commanded_from = Some(source);
graph.handoff_count += 1;
tracing::info!(
task_id = %source,
target = %target,
handoff_count = graph.handoff_count,
"orchestration.dag.try_handoff: Command handoff activated"
);
Ok(target)
}
pub fn validate_handoff_target(
graph: &TaskGraph,
source: TaskId,
goto: &TaskRef,
) -> Result<TaskId, OrchestrationError> {
let target = resolve_task_ref(graph, goto)?;
let target_task = &graph.tasks[target.index()];
if target_task.status == TaskStatus::Completed {
return Err(OrchestrationError::InvalidHandoffTarget(format!(
"task {target} is already completed — Command.goto is forward-only"
)));
}
if !target_task
.depends_on
.iter()
.all(|dep| graph.tasks[dep.index()].status == TaskStatus::Completed)
{
return Err(OrchestrationError::InvalidHandoffTarget(format!(
"task {target} has unsatisfied depends_on — a Command.goto target must have \
empty or fully-completed dependencies"
)));
}
if target_task.status == TaskStatus::Dormant {
let reservation_source = graph
.tasks
.iter()
.find(|t| t.recovery.as_ref().and_then(|r| r.route_to) == Some(target));
if let Some(reservation_source) = reservation_source
&& !reservation_source.status.is_terminal()
{
return Err(OrchestrationError::InvalidHandoffTarget(format!(
"task {target} holds a live route_to reservation from non-terminal source {src}",
src = reservation_source.id
)));
}
}
tracing::debug!(
task_id = %source,
target = %target,
"orchestration.dag.validate_handoff_target: target passed read-only validation"
);
Ok(target)
}
fn skip_subtree(graph: &mut TaskGraph, seed: TaskId, rev_adj: &[Vec<TaskId>]) -> Vec<TaskId> {
let mut to_cancel = Vec::new();
let mut queue: VecDeque<TaskId> = VecDeque::new();
queue.push_back(seed);
while let Some(current) = queue.pop_front() {
let dependents = rev_adj.get(current.index()).map_or(&[] as &[TaskId], |v| v);
for &dep_id in dependents {
if !graph.tasks[dep_id.index()].status.is_terminal() {
if graph.tasks[dep_id.index()].status == TaskStatus::Running {
to_cancel.push(dep_id);
}
graph.tasks[dep_id.index()].status = TaskStatus::Skipped;
queue.push_back(dep_id);
}
}
}
to_cancel
}
pub(crate) fn cancel_dangling_commanded_targets(
graph: &mut TaskGraph,
source: TaskId,
rev_adj: &[Vec<TaskId>],
) -> Vec<TaskId> {
let targets: Vec<TaskId> = graph
.tasks
.iter()
.filter(|t| t.commanded_from == Some(source))
.filter(|t| matches!(t.status, TaskStatus::Pending | TaskStatus::Ready))
.map(|t| t.id)
.collect();
let mut to_cancel = Vec::new();
for target in targets {
tracing::warn!(
source = %source,
target = %target,
"orchestration.dag.cancel_dangling_commanded_targets: skipping Command-handoff \
target whose source was corrected Completed -> Failed post-hoc (#6394)"
);
graph.tasks[target.index()].status = TaskStatus::Skipped;
to_cancel.extend(skip_subtree(graph, target, rev_adj));
}
to_cancel
}
pub(crate) fn resolve_dormant_after_terminal(
graph: &mut TaskGraph,
rev_adj: &[Vec<TaskId>],
) -> Vec<TaskId> {
let dormant_sources: Vec<(TaskId, TaskId)> = graph
.tasks
.iter()
.filter(|t| t.status == TaskStatus::Dormant)
.filter_map(|target| {
graph
.tasks
.iter()
.find(|t| t.recovery.as_ref().and_then(|r| r.route_to) == Some(target.id))
.map(|source| (target.id, source.id))
})
.collect();
let mut resolved = Vec::new();
for (target, source) in dormant_sources {
if graph.tasks[source.index()].status.is_terminal() {
graph.tasks[target.index()].status = TaskStatus::Skipped;
resolved.push(target);
skip_subtree(graph, target, rev_adj);
}
}
resolved
}
pub fn propagate_failure(
graph: &mut TaskGraph,
failed_id: TaskId,
rev_adj: &[Vec<TaskId>],
) -> Vec<TaskId> {
if graph.tasks[failed_id.index()].status != TaskStatus::Failed {
return Vec::new();
}
let strategy = graph.tasks[failed_id.index()]
.failure_strategy
.unwrap_or(graph.default_failure_strategy);
let max_retries = graph.tasks[failed_id.index()]
.max_retries
.unwrap_or(graph.default_max_retries);
match strategy {
FailureStrategy::Abort => {
if try_recover(graph, failed_id) {
return Vec::new();
}
if try_reroute(graph, failed_id) {
return Vec::new();
}
graph.status = GraphStatus::Failed;
graph
.tasks
.iter()
.filter(|t| t.status == TaskStatus::Running)
.map(|t| t.id)
.collect()
}
FailureStrategy::Skip => {
graph.tasks[failed_id.index()].status = TaskStatus::Skipped;
skip_subtree(graph, failed_id, rev_adj)
}
FailureStrategy::Retry => {
let retry_count = graph.tasks[failed_id.index()].retry_count;
if retry_count < max_retries {
graph.tasks[failed_id.index()].retry_count += 1;
graph.tasks[failed_id.index()].status = TaskStatus::Ready;
Vec::new()
} else {
if try_recover(graph, failed_id) {
return Vec::new();
}
if try_reroute(graph, failed_id) {
return Vec::new();
}
graph.status = GraphStatus::Failed;
graph
.tasks
.iter()
.filter(|t| t.status == TaskStatus::Running)
.map(|t| t.id)
.collect()
}
}
FailureStrategy::Ask => {
graph.status = GraphStatus::Paused;
Vec::new()
}
strategy => {
tracing::error!(
?strategy,
"unhandled failure strategy variant, defaulting to Abort"
);
graph.status = GraphStatus::Failed;
graph
.tasks
.iter()
.filter(|t| t.status == TaskStatus::Running)
.map(|t| t.id)
.collect()
}
}
}
pub(crate) fn propagate_failure_forced_terminal(
graph: &mut TaskGraph,
failed_id: TaskId,
) -> Vec<TaskId> {
if graph.tasks[failed_id.index()].status != TaskStatus::Failed {
return Vec::new();
}
graph.status = GraphStatus::Failed;
graph
.tasks
.iter()
.filter(|t| t.status == TaskStatus::Running)
.map(|t| t.id)
.collect()
}
pub fn reset_for_retry(
graph: &mut TaskGraph,
rev_adj: &[Vec<TaskId>],
) -> Result<(), OrchestrationError> {
use super::graph::GraphStatus;
if graph.status != GraphStatus::Failed && graph.status != GraphStatus::Paused {
return Err(OrchestrationError::InvalidGraph(format!(
"cannot retry graph in status {}; only Failed or Paused graphs can be retried",
graph.status
)));
}
let mut seeds: Vec<TaskId> = Vec::new();
for task in &mut graph.tasks {
if task.status == TaskStatus::Failed {
task.status = TaskStatus::Ready;
task.retry_count = 0;
seeds.push(task.id);
}
}
for task in &mut graph.tasks {
if task.status == TaskStatus::Canceled {
task.status = TaskStatus::Pending;
}
}
if seeds.is_empty() {
graph.status = GraphStatus::Running;
return Ok(());
}
let seeds_for_reroute = seeds.clone();
let mut queue: std::collections::VecDeque<TaskId> = seeds.into_iter().collect();
while let Some(current) = queue.pop_front() {
let dependents = rev_adj.get(current.index()).map_or(&[] as &[TaskId], |v| v);
for &dep_id in dependents {
if graph.tasks[dep_id.index()].status == TaskStatus::Skipped {
graph.tasks[dep_id.index()].status = TaskStatus::Pending;
queue.push_back(dep_id);
}
}
}
for s_id in seeds_for_reroute {
let Some(target) = graph.tasks[s_id.index()]
.recovery
.as_ref()
.and_then(|r| r.route_to)
else {
continue;
};
let target_node = &mut graph.tasks[target.index()];
target_node.status = TaskStatus::Dormant;
target_node.routed_from = None;
target_node.retry_count = 0;
target_node.result = None;
let mut re_arm_queue: VecDeque<TaskId> = VecDeque::new();
re_arm_queue.push_back(target);
while let Some(current) = re_arm_queue.pop_front() {
let dependents = rev_adj.get(current.index()).map_or(&[] as &[TaskId], |v| v);
for &dep_id in dependents {
let dep = &mut graph.tasks[dep_id.index()];
if dep.status != TaskStatus::Pending {
dep.status = TaskStatus::Pending;
dep.retry_count = 0;
dep.result = None;
re_arm_queue.push_back(dep_id);
}
}
}
}
graph.status = GraphStatus::Running;
Ok(())
}
const KEYWORD_STOPWORDS: &[&str] = &["the", "a", "an", "in", "of", "for", "to", "from", "with"];
#[must_use]
pub fn lookahead_tools(graph: &TaskGraph, depth: u8) -> Vec<PlannedToolHint> {
let _span = tracing::debug_span!("orch.dag.lookahead", depth = depth).entered();
if depth == 0 {
return vec![];
}
let tasks = &graph.tasks;
let n = tasks.len();
let mut forward_adj: Vec<Vec<usize>> = vec![Vec::new(); n];
for task in tasks {
for dep in &task.depends_on {
forward_adj[dep.index()].push(task.id.index());
}
}
let mut visited = vec![false; n];
let mut queue: VecDeque<(usize, u8)> = VecDeque::new();
for task in tasks {
if matches!(task.status, TaskStatus::Running | TaskStatus::Ready) {
visited[task.id.index()] = true;
queue.push_back((task.id.index(), 0));
}
}
if queue.is_empty() {
return vec![];
}
let mut hints: Vec<PlannedToolHint> = Vec::new();
while let Some((idx, dist)) = queue.pop_front() {
for &child_idx in &forward_adj[idx] {
if visited[child_idx] {
continue;
}
visited[child_idx] = true;
let child_dist = dist + 1;
if child_dist <= depth {
let child = &tasks[child_idx];
let tool_name = child.agent_hint.as_deref().unwrap_or(&child.title);
hints.push(PlannedToolHint::new(
tool_name,
extract_keywords(tool_name, &child.description),
child_dist,
));
queue.push_back((child_idx, child_dist));
}
}
}
hints.sort_by_key(|h| h.distance_from_current);
hints
}
fn extract_keywords(tool_name: &str, description: &str) -> Vec<String> {
let end = description.floor_char_boundary(200);
let desc_prefix = &description[..end];
let combined = format!("{tool_name} {desc_prefix}");
let mut seen = std::collections::HashSet::new();
let mut keywords: Vec<String> = Vec::new();
let full = tool_name.to_lowercase();
seen.insert(full.clone());
keywords.push(full);
for token in combined.split(|c: char| !c.is_alphanumeric()) {
if keywords.len() == 10 {
break;
}
if token.len() < 3 {
continue;
}
let lower = token.to_lowercase();
if KEYWORD_STOPWORDS.contains(&lower.as_str()) {
continue;
}
if seen.insert(lower.clone()) {
keywords.push(lower);
}
}
keywords
}
pub fn inject_tasks(
graph: &mut TaskGraph,
new_tasks: Vec<TaskNode>,
max_tasks: usize,
) -> Result<(), OrchestrationError> {
if new_tasks.is_empty() {
return Ok(());
}
let existing_len = graph.tasks.len();
let total = existing_len + new_tasks.len();
if total > max_tasks {
return Err(OrchestrationError::VerificationFailed(format!(
"inject_tasks would create {total} tasks, exceeding limit of {max_tasks}"
)));
}
for (i, task) in new_tasks.iter().enumerate() {
let expected = TaskId(u32::try_from(existing_len + i).map_err(|_| {
OrchestrationError::VerificationFailed("task index overflows u32".to_string())
})?);
if task.id != expected {
return Err(OrchestrationError::VerificationFailed(format!(
"injected task at position {} has id {} (expected {})",
i, task.id, expected
)));
}
}
graph.tasks.extend(new_tasks);
validate(&graph.tasks, max_tasks, graph.default_failure_strategy).map_err(|e| match e {
OrchestrationError::CycleDetected => {
OrchestrationError::VerificationFailed("inject_tasks introduced a cycle".to_string())
}
other => OrchestrationError::VerificationFailed(other.to_string()),
})?;
let n = graph.tasks.len();
for i in existing_len..n {
let all_deps_done = graph.tasks[i]
.depends_on
.iter()
.all(|dep| graph.tasks[dep.index()].status == TaskStatus::Completed);
if all_deps_done {
graph.tasks[i].status = TaskStatus::Ready;
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::graph::{FailureStrategy, GraphStatus, TaskGraph, TaskNode, TaskStatus};
use crate::topology::build_rev_adj;
use std::assert_matches;
fn make_node(id: u32, deps: &[u32]) -> TaskNode {
let mut n = TaskNode::new(id, format!("task-{id}"), "desc");
n.depends_on = deps.iter().map(|&d| TaskId(d)).collect();
n
}
fn graph_from_nodes(nodes: Vec<TaskNode>) -> TaskGraph {
let mut g = TaskGraph::new("test");
g.tasks = nodes;
g
}
fn make_rev_adj(graph: &TaskGraph) -> Vec<Vec<TaskId>> {
build_rev_adj(&graph.tasks)
}
#[test]
fn test_validate_empty_graph() {
let err = validate(&[], 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_exceeds_max_tasks() {
let tasks: Vec<TaskNode> = (0..5).map(|i| make_node(i, &[])).collect();
let err = validate(&tasks, 3, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_single_task_no_deps() {
let tasks = vec![make_node(0, &[])];
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_self_reference() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].depends_on = vec![TaskId(0)];
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_invalid_taskid_reference() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].depends_on = vec![TaskId(99)];
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_linear_chain() {
let tasks = vec![make_node(0, &[]), make_node(1, &[0]), make_node(2, &[1])];
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_diamond() {
let tasks = vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[0]),
make_node(3, &[1, 2]),
];
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_cycle_two_nodes() {
let tasks = vec![make_node(0, &[1]), make_node(1, &[0])];
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::CycleDetected);
}
#[test]
fn test_validate_cycle_three_nodes() {
let tasks = vec![make_node(0, &[2]), make_node(1, &[0]), make_node(2, &[1])];
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::CycleDetected);
}
#[test]
fn test_validate_taskid_invariant() {
let mut tasks = vec![make_node(0, &[]), make_node(1, &[0])];
tasks[1].id = TaskId(5);
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_rejects_recovery_with_verify_predicate() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback".to_string()),
route_to: None,
});
tasks[0].verify_predicate = Some(crate::graph::VerifyPredicate::Natural(
"criterion".to_string(),
));
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_recovery_alone_is_ok() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback".to_string()),
route_to: None,
});
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_recovery_with_skip_strategy_warns_but_ok() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback".to_string()),
route_to: None,
});
tasks[0].failure_strategy = Some(FailureStrategy::Skip);
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_recovery_with_ask_strategy_warns_but_ok() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback".to_string()),
route_to: None,
});
tasks[0].failure_strategy = Some(FailureStrategy::Ask);
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_recovery_with_abort_or_retry_no_warning_ok() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback".to_string()),
route_to: None,
});
tasks[0].failure_strategy = Some(FailureStrategy::Retry);
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
let mut tasks2 = vec![make_node(0, &[])];
tasks2[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback".to_string()),
route_to: None,
});
assert!(validate(&tasks2, 20, FailureStrategy::Abort).is_ok());
}
fn make_route_to_pair() -> Vec<TaskNode> {
let mut tasks = vec![make_node(0, &[]), make_node(1, &[])];
tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
tasks
}
#[test]
fn test_validate_route_to_valid_pair_ok() {
let tasks = make_route_to_pair();
assert!(validate(&tasks, 20, FailureStrategy::Abort).is_ok());
}
#[test]
fn test_validate_route_to_self_reroute_rejected() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_out_of_range_rejected() {
let mut tasks = vec![make_node(0, &[])];
tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(99)),
});
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_and_state_injection_mutually_exclusive() {
let mut tasks = make_route_to_pair();
tasks[1].recovery.as_mut().unwrap().state_injection = Some("fallback".to_string());
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_target_with_deps_rejected() {
let mut tasks = vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[1]), ];
tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(2)),
});
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_chained_rejected() {
let mut tasks = vec![
make_node(0, &[]), make_node(1, &[]), make_node(2, &[]), ];
tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
tasks[2].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(1)),
});
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_n_to_one_rejected() {
let mut tasks = vec![
make_node(0, &[]), make_node(1, &[]), make_node(2, &[]), ];
tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
tasks[2].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_under_skip_strategy_rejected() {
let mut tasks = make_route_to_pair();
tasks[1].failure_strategy = Some(FailureStrategy::Skip);
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_under_ask_strategy_rejected() {
let mut tasks = make_route_to_pair();
tasks[1].failure_strategy = Some(FailureStrategy::Ask);
let err = validate(&tasks, 20, FailureStrategy::Abort).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_validate_route_to_under_default_skip_strategy_rejected() {
let tasks = make_route_to_pair();
let err = validate(&tasks, 20, FailureStrategy::Skip).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_try_handoff_activates_pending_target() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Pending;
let target = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(1)), 16).unwrap();
assert_eq!(target, TaskId(1));
assert_eq!(graph.tasks[1].status, TaskStatus::Ready);
assert_eq!(graph.tasks[1].commanded_from, Some(TaskId(0)));
assert_eq!(graph.handoff_count, 1);
}
#[test]
fn test_try_handoff_activates_dormant_target() {
let mut tasks = vec![make_node(0, &[]), make_node(1, &[]), make_node(2, &[])];
tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(2)),
});
let mut graph = graph_from_nodes(tasks);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Completed; graph.tasks[2].status = TaskStatus::Dormant;
let target = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(2)), 16).unwrap();
assert_eq!(target, TaskId(2));
assert_eq!(graph.tasks[2].status, TaskStatus::Ready);
assert_eq!(graph.tasks[2].commanded_from, Some(TaskId(0)));
}
#[test]
fn test_try_handoff_by_title_resolves() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
let target = try_handoff(
&mut graph,
TaskId(0),
&TaskRef::ByTitle("task-1".to_string()),
16,
)
.unwrap();
assert_eq!(target, TaskId(1));
assert_eq!(graph.tasks[1].status, TaskStatus::Ready);
}
#[test]
fn test_try_handoff_by_title_ambiguous_rejected() {
let mut tasks = vec![make_node(0, &[]), make_node(1, &[]), make_node(2, &[])];
tasks[2].title = "task-1".to_string(); let mut graph = graph_from_nodes(tasks);
graph.tasks[0].status = TaskStatus::Completed;
let err = try_handoff(
&mut graph,
TaskId(0),
&TaskRef::ByTitle("task-1".to_string()),
16,
)
.unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
assert_eq!(
graph.handoff_count, 0,
"a rejected handoff must not consume budget"
);
}
#[test]
fn test_try_handoff_by_title_no_match_rejected() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
let err = try_handoff(
&mut graph,
TaskId(0),
&TaskRef::ByTitle("does-not-exist".to_string()),
16,
)
.unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
}
#[test]
fn test_try_handoff_by_id_out_of_range_rejected() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
let err = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(99)), 16).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
}
#[test]
fn test_try_handoff_rejects_completed_target_forward_only() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Completed;
let err = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(1)), 16).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
assert_eq!(graph.handoff_count, 0);
}
#[test]
fn test_try_handoff_rejects_unsatisfied_depends_on() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[1]),
]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[2].status = TaskStatus::Pending;
let err = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(2)), 16).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
assert_eq!(
graph.tasks[2].status,
TaskStatus::Pending,
"target must not be force-activated with unsatisfied deps"
);
}
#[test]
fn test_try_handoff_allows_fully_satisfied_depends_on() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[1]),
]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Completed;
graph.tasks[2].status = TaskStatus::Pending;
let target = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(2)), 16).unwrap();
assert_eq!(target, TaskId(2));
assert_eq!(graph.tasks[2].status, TaskStatus::Ready);
}
#[test]
fn test_try_handoff_rejects_live_route_to_reservation_from_non_terminal_source() {
let mut tasks = vec![make_node(0, &[]), make_node(1, &[]), make_node(2, &[])];
tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(2)),
});
let mut graph = graph_from_nodes(tasks);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Running; graph.tasks[2].status = TaskStatus::Dormant;
let err = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(2)), 16).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
assert_eq!(
graph.tasks[2].status,
TaskStatus::Dormant,
"target must remain parked while its route_to reservation is live"
);
}
#[test]
fn test_try_handoff_budget_exhausted() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.handoff_count = 16;
let err = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(1)), 16).unwrap_err();
assert_matches!(
err,
OrchestrationError::HandoffBudgetExhausted {
handoff_count: 16,
max_handoffs: 16,
}
);
assert_eq!(
graph.tasks[1].status,
TaskStatus::Pending,
"exhausted-budget rejection must not activate the target"
);
}
#[test]
fn test_validate_handoff_target_does_not_mutate_graph() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Pending;
let target = validate_handoff_target(&graph, TaskId(0), &TaskRef::ById(TaskId(1))).unwrap();
assert_eq!(target, TaskId(1));
assert_eq!(
graph.tasks[1].status,
TaskStatus::Pending,
"read-only pre-check must not activate the target"
);
assert_eq!(
graph.tasks[1].commanded_from, None,
"read-only pre-check must not set commanded_from"
);
assert_eq!(
graph.handoff_count, 0,
"read-only pre-check must not consume budget"
);
}
#[test]
fn test_validate_handoff_target_ignores_exhausted_budget() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.handoff_count = 16;
let target = validate_handoff_target(&graph, TaskId(0), &TaskRef::ById(TaskId(1))).unwrap();
assert_eq!(target, TaskId(1));
}
#[test]
fn test_validate_handoff_target_rejects_same_reasons_as_try_handoff() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Completed;
let err =
validate_handoff_target(&graph, TaskId(0), &TaskRef::ById(TaskId(1))).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidHandoffTarget(_));
}
#[test]
fn test_try_handoff_redundant_activation_on_already_ready_target_is_accepted() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Ready;
graph.tasks[1].routed_from = Some(TaskId(99));
let target = try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(1)), 16).unwrap();
assert_eq!(target, TaskId(1));
assert_eq!(graph.tasks[1].status, TaskStatus::Ready);
assert_eq!(graph.tasks[1].commanded_from, Some(TaskId(0)));
assert_eq!(
graph.tasks[1].routed_from,
Some(TaskId(99)),
"an unrelated earlier routed_from marker must be preserved, not clobbered"
);
assert_eq!(graph.handoff_count, 1);
}
#[test]
fn test_try_handoff_multiple_hops_increment_budget() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[]),
]);
graph.tasks[0].status = TaskStatus::Completed;
try_handoff(&mut graph, TaskId(0), &TaskRef::ById(TaskId(1)), 16).unwrap();
assert_eq!(graph.handoff_count, 1);
graph.tasks[1].status = TaskStatus::Completed;
try_handoff(&mut graph, TaskId(1), &TaskRef::ById(TaskId(2)), 16).unwrap();
assert_eq!(graph.handoff_count, 2);
}
#[test]
fn test_mark_dormant_route_to_targets_marks_pending_target() {
let mut graph = graph_from_nodes(make_route_to_pair());
mark_dormant_route_to_targets(&mut graph);
assert_eq!(graph.tasks[0].status, TaskStatus::Dormant);
}
#[test]
fn test_mark_dormant_route_to_targets_guard_skips_non_pending() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Completed;
mark_dormant_route_to_targets(&mut graph);
assert_eq!(
graph.tasks[0].status,
TaskStatus::Completed,
"guard must not re-dormant an already-terminal target"
);
}
#[test]
fn test_mark_dormant_route_to_targets_no_route_to_is_noop() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
mark_dormant_route_to_targets(&mut graph);
assert_eq!(graph.tasks[0].status, TaskStatus::Pending);
}
#[test]
fn test_propagate_failure_abort_reroutes_to_dormant_target() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Failed;
graph.tasks[1].failure_strategy = Some(FailureStrategy::Abort);
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(1), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(
graph.tasks[1].status,
TaskStatus::Failed,
"source stays terminal Failed"
);
assert_eq!(graph.tasks[0].status, TaskStatus::Ready, "target activated");
assert_eq!(graph.tasks[0].routed_from, Some(TaskId(1)));
assert_eq!(
graph.status,
GraphStatus::Running,
"graph.status must be left untouched by reroute"
);
}
#[test]
fn test_propagate_failure_retry_exhausted_reroutes_to_dormant_target() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Failed;
graph.tasks[1].failure_strategy = Some(FailureStrategy::Retry);
graph.tasks[1].max_retries = Some(3);
graph.tasks[1].retry_count = 3;
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(1), &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(graph.tasks[0].routed_from, Some(TaskId(1)));
assert_eq!(graph.status, GraphStatus::Running);
}
#[test]
fn test_propagate_failure_reroute_runtime_guard_refuses_non_dormant_target() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Ready; graph.tasks[1].status = TaskStatus::Failed;
graph.tasks[1].failure_strategy = Some(FailureStrategy::Abort);
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(1), &__ra);
assert_eq!(
graph.tasks[0].status,
TaskStatus::Ready,
"runtime guard must not mutate a non-Dormant target"
);
assert_eq!(graph.tasks[0].routed_from, None);
assert_eq!(
graph.status,
GraphStatus::Failed,
"must fall through to Abort when reroute is refused"
);
}
#[test]
fn test_route_to_target_dependent_becomes_ready_after_reroute() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0]),
]);
graph.tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Failed;
graph.tasks[1].failure_strategy = Some(FailureStrategy::Abort);
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(1), &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
graph.tasks[0].status = TaskStatus::Completed;
let ready = ready_tasks(&graph);
assert!(ready.contains(&TaskId(2)));
}
#[test]
fn test_resolve_dormant_after_terminal_skips_untriggered_fallback_on_source_success() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Completed;
let __ra = make_rev_adj(&graph);
let resolved = resolve_dormant_after_terminal(&mut graph, &__ra);
assert_eq!(resolved, vec![TaskId(0)]);
assert_eq!(graph.tasks[0].status, TaskStatus::Skipped);
}
#[test]
fn test_resolve_dormant_after_terminal_skips_subtree() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0]),
]);
graph.tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Completed;
graph.tasks[2].status = TaskStatus::Pending;
let __ra = make_rev_adj(&graph);
resolve_dormant_after_terminal(&mut graph, &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Skipped);
assert_eq!(
graph.tasks[2].status,
TaskStatus::Skipped,
"F's downstream subtree must be skipped when the fallback is never triggered"
);
}
#[test]
fn test_resolve_dormant_after_terminal_noop_while_source_running() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Running;
let __ra = make_rev_adj(&graph);
let resolved = resolve_dormant_after_terminal(&mut graph, &__ra);
assert!(resolved.is_empty());
assert_eq!(graph.tasks[0].status, TaskStatus::Dormant);
}
#[test]
fn test_resolve_dormant_after_terminal_ignores_activated_target() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Ready;
graph.tasks[0].routed_from = Some(TaskId(1));
graph.tasks[1].status = TaskStatus::Failed;
let __ra = make_rev_adj(&graph);
let resolved = resolve_dormant_after_terminal(&mut graph, &__ra);
assert!(resolved.is_empty());
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
}
#[test]
fn test_reset_for_retry_rearms_dormant_fallback_that_never_ran() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Failed;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(
graph.tasks[1].status,
TaskStatus::Ready,
"source reset for retry"
);
assert_eq!(graph.tasks[0].status, TaskStatus::Dormant);
assert_eq!(graph.tasks[0].routed_from, None);
}
#[test]
fn test_reset_for_retry_rearms_activated_fallback_case_a_source_now_succeeds() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Completed; graph.tasks[0].routed_from = Some(TaskId(1));
graph.tasks[0].result = Some(TaskResult {
output: "stale fallback output".to_string(),
artifacts: Vec::new(),
duration_ms: 5,
agent_id: None,
agent_def: None,
});
graph.tasks[1].status = TaskStatus::Failed; graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[1].status, TaskStatus::Ready);
assert_eq!(
graph.tasks[0].status,
TaskStatus::Dormant,
"activated fallback must be re-armed to Dormant on retry"
);
assert_eq!(
graph.tasks[0].routed_from, None,
"stale routed_from must be cleared"
);
assert!(
graph.tasks[0].result.is_none(),
"stale result must be cleared"
);
graph.tasks[1].status = TaskStatus::Completed;
let ready = ready_tasks(&graph);
assert!(
!ready.contains(&TaskId(0)),
"re-armed Dormant fallback must not dispatch when the source now succeeds"
);
}
#[test]
fn test_reset_for_retry_rearms_activated_fallback_case_b_source_fails_again() {
let mut graph = graph_from_nodes(make_route_to_pair());
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[0].routed_from = Some(TaskId(1));
graph.tasks[1].status = TaskStatus::Failed;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Dormant);
graph.tasks[1].status = TaskStatus::Failed;
graph.tasks[1].failure_strategy = Some(FailureStrategy::Abort);
let to_cancel = propagate_failure(&mut graph, TaskId(1), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(
graph.tasks[0].status,
TaskStatus::Ready,
"reroute must fire again after the re-arm"
);
assert_eq!(graph.tasks[0].routed_from, Some(TaskId(1)));
}
#[test]
fn test_reset_for_retry_rearm_resets_fallback_subtree() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0]),
]);
graph.tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[0].routed_from = Some(TaskId(1));
graph.tasks[1].status = TaskStatus::Failed;
graph.tasks[2].status = TaskStatus::Completed;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Dormant);
assert_eq!(
graph.tasks[2].status,
TaskStatus::Pending,
"F's downstream subtree must reset to Pending alongside the re-arm"
);
}
#[test]
fn test_reset_for_retry_does_not_rearm_untouched_route_to_source() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]), make_node(1, &[]), make_node(2, &[]), ]);
graph.tasks[1].recovery = Some(crate::graph::RecoveryAction {
state_injection: None,
route_to: Some(TaskId(0)),
});
graph.tasks[0].status = TaskStatus::Dormant;
graph.tasks[1].status = TaskStatus::Completed; graph.tasks[2].status = TaskStatus::Failed;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(
graph.tasks[1].status,
TaskStatus::Completed,
"S was not reset (not Failed)"
);
assert_eq!(
graph.tasks[0].status,
TaskStatus::Dormant,
"F must be untouched since its source was never reset"
);
}
#[test]
fn test_toposort_linear() {
let tasks = vec![make_node(0, &[]), make_node(1, &[0]), make_node(2, &[1])];
let order = toposort(&tasks).expect("should succeed");
assert_eq!(order, vec![TaskId(0), TaskId(1), TaskId(2)]);
}
#[test]
fn test_toposort_diamond() {
let tasks = vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[0]),
make_node(3, &[1, 2]),
];
let order = toposort(&tasks).expect("should succeed");
assert_eq!(order[0], TaskId(0));
assert_eq!(order[3], TaskId(3));
}
#[test]
fn test_toposort_wide_parallel() {
let tasks = vec![make_node(0, &[]), make_node(1, &[]), make_node(2, &[])];
let order = toposort(&tasks).expect("should succeed");
assert_eq!(order.len(), 3);
}
#[test]
fn test_toposort_single_node() {
let tasks = vec![make_node(0, &[])];
let order = toposort(&tasks).expect("should succeed");
assert_eq!(order, vec![TaskId(0)]);
}
#[test]
fn test_ready_tasks_initial_roots() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0, 1]),
]);
graph.tasks[0].status = TaskStatus::Pending;
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[2].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(ready.contains(&TaskId(0)));
assert!(ready.contains(&TaskId(1)));
assert!(!ready.contains(&TaskId(2)));
}
#[test]
fn test_ready_tasks_after_completion() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(ready.contains(&TaskId(1)));
}
#[test]
fn test_ready_tasks_skipped_does_not_unblock() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Skipped;
graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(!ready.contains(&TaskId(1)));
}
#[test]
fn test_ready_tasks_partial_deps_completed() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0, 1]),
]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Running;
graph.tasks[2].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(!ready.contains(&TaskId(2)));
}
#[test]
fn test_ready_tasks_all_terminal() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Completed;
let ready = ready_tasks(&graph);
assert!(ready.is_empty());
}
#[test]
fn test_ready_tasks_already_ready_included() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Ready; graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(ready.contains(&TaskId(0)));
}
#[test]
fn test_ready_tasks_predicate_gate_blocks_downstream() {
use crate::graph::VerifyPredicate;
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[0].verify_predicate = Some(VerifyPredicate::Natural(
"output must be non-empty".to_string(),
));
graph.tasks[0].predicate_outcome = None;
graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(
!ready.contains(&TaskId(1)),
"task 1 must be blocked by uncleared predicate on task 0"
);
}
#[test]
fn test_ready_tasks_predicate_gate_unblocks_on_pass() {
use crate::graph::{PredicateOutcome, VerifyPredicate};
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[0].verify_predicate = Some(VerifyPredicate::Natural("criterion".to_string()));
graph.tasks[0].predicate_outcome = Some(PredicateOutcome {
passed: true,
confidence: 0.9,
reason: "ok".to_string(),
});
graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(
ready.contains(&TaskId(1)),
"task 1 must be unblocked when predicate passed"
);
}
#[test]
fn test_ready_tasks_predicate_gate_remains_closed_on_fail() {
use crate::graph::{PredicateOutcome, VerifyPredicate};
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[0].verify_predicate = Some(VerifyPredicate::Natural("criterion".to_string()));
graph.tasks[0].predicate_outcome = Some(PredicateOutcome {
passed: false,
confidence: 0.1,
reason: "criterion not met".to_string(),
});
graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(
!ready.contains(&TaskId(1)),
"task 1 must remain blocked when predicate failed"
);
}
#[test]
fn test_ready_tasks_no_predicate_unblocks_normally() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Pending;
let ready = ready_tasks(&graph);
assert!(
ready.contains(&TaskId(1)),
"no predicate = gate always clear"
);
}
#[test]
fn test_propagate_failure_abort() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[0]),
]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[1].status = TaskStatus::Running;
graph.tasks[2].status = TaskStatus::Pending;
graph.default_failure_strategy = FailureStrategy::Abort;
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.status, GraphStatus::Failed);
assert!(to_cancel.contains(&TaskId(1)));
assert!(!to_cancel.contains(&TaskId(2)));
}
#[test]
fn test_propagate_failure_skip_single() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Skip);
graph.tasks[1].status = TaskStatus::Pending;
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(graph.tasks[0].status, TaskStatus::Skipped);
assert_eq!(graph.tasks[1].status, TaskStatus::Skipped);
}
#[test]
fn test_propagate_failure_skip_transitive() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[1]),
]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Skip);
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[2].status = TaskStatus::Pending;
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Skipped);
assert_eq!(graph.tasks[1].status, TaskStatus::Skipped);
assert_eq!(graph.tasks[2].status, TaskStatus::Skipped);
}
#[test]
fn test_propagate_failure_skip_running_dependent_returned() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Skip);
graph.tasks[1].status = TaskStatus::Running;
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(
to_cancel.contains(&TaskId(1)),
"Running dependent must be returned for cancellation"
);
assert_eq!(graph.tasks[1].status, TaskStatus::Skipped);
}
#[test]
fn test_propagate_failure_retry_under_max() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Retry);
graph.tasks[0].max_retries = Some(3);
graph.tasks[0].retry_count = 1;
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(graph.tasks[0].retry_count, 2);
}
#[test]
fn test_propagate_failure_retry_exhausted() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Retry);
graph.tasks[0].max_retries = Some(3);
graph.tasks[0].retry_count = 3;
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.status, GraphStatus::Failed);
}
#[test]
fn test_propagate_failure_ask() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Ask);
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(graph.status, GraphStatus::Paused);
}
#[test]
fn test_propagate_failure_per_task_override() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.default_failure_strategy = FailureStrategy::Abort;
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Skip);
graph.tasks[1].status = TaskStatus::Pending;
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Skipped);
assert_ne!(graph.status, GraphStatus::Failed);
}
#[test]
fn test_propagate_failure_already_terminal() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(graph.status, GraphStatus::Created);
}
#[test]
fn test_propagate_failure_abort_recovers_with_state_injection() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Abort);
graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(graph.tasks[0].status, TaskStatus::Completed);
assert_eq!(
graph.tasks[0].result.as_ref().unwrap().output,
"fallback output"
);
assert_eq!(
graph.tasks[0].result.as_ref().unwrap().agent_def.as_deref(),
Some("__recovery__")
);
assert_eq!(
graph.status,
GraphStatus::Running,
"graph.status must be left untouched by recovery"
);
}
#[test]
fn test_propagate_failure_retry_exhausted_recovers_with_state_injection() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Retry);
graph.tasks[0].max_retries = Some(3);
graph.tasks[0].retry_count = 3; graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
let __ra = make_rev_adj(&graph);
let to_cancel = propagate_failure(&mut graph, TaskId(0), &__ra);
assert!(to_cancel.is_empty());
assert_eq!(graph.tasks[0].status, TaskStatus::Completed);
assert_eq!(
graph.tasks[0].result.as_ref().unwrap().output,
"fallback output"
);
assert_eq!(graph.status, GraphStatus::Running);
}
#[test]
fn test_propagate_failure_abort_no_recovery_configured_is_unchanged() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Abort);
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.status, GraphStatus::Failed);
assert_eq!(graph.tasks[0].status, TaskStatus::Failed);
}
#[test]
fn test_propagate_failure_retry_exhausted_no_recovery_configured_is_unchanged() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Retry);
graph.tasks[0].max_retries = Some(3);
graph.tasks[0].retry_count = 3;
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.status, GraphStatus::Failed);
}
#[test]
fn test_recovered_task_dependent_becomes_ready() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.status = GraphStatus::Running;
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Abort);
graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
graph.tasks[1].status = TaskStatus::Pending;
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Completed);
let ready = ready_tasks(&graph);
assert!(
ready.contains(&TaskId(1)),
"dependent must unblock via the Pending arm after recovery"
);
}
#[test]
fn test_skip_strategy_with_recovery_configured_still_skips() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Skip);
graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.tasks[0].status, TaskStatus::Skipped);
}
#[test]
fn test_ask_strategy_with_recovery_configured_still_pauses() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].failure_strategy = Some(FailureStrategy::Ask);
graph.tasks[0].recovery = Some(crate::graph::RecoveryAction {
state_injection: Some("fallback output".to_string()),
route_to: None,
});
let __ra = make_rev_adj(&graph);
propagate_failure(&mut graph, TaskId(0), &__ra);
assert_eq!(graph.status, GraphStatus::Paused);
assert_ne!(graph.tasks[0].status, TaskStatus::Completed);
}
#[test]
fn test_reset_for_retry_resets_failed_to_ready() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(graph.status, GraphStatus::Running);
}
#[test]
fn test_reset_for_retry_resets_skipped_dependents_to_pending() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[1].status = TaskStatus::Skipped;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(graph.tasks[1].status, TaskStatus::Pending);
}
#[test]
fn test_reset_for_retry_transitive_skipped_reset() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[0]),
make_node(2, &[1]),
]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[1].status = TaskStatus::Skipped;
graph.tasks[2].status = TaskStatus::Skipped;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(graph.tasks[1].status, TaskStatus::Pending);
assert_eq!(graph.tasks[2].status, TaskStatus::Pending);
}
#[test]
fn test_reset_for_retry_completed_tasks_unchanged() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.tasks[1].status = TaskStatus::Failed;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Completed);
assert_eq!(graph.tasks[1].status, TaskStatus::Ready);
}
#[test]
fn test_reset_for_retry_rejects_running_graph() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Running;
graph.status = GraphStatus::Running;
let __ra = make_rev_adj(&graph);
let err = reset_for_retry(&mut graph, &__ra).unwrap_err();
assert_matches!(err, OrchestrationError::InvalidGraph(_));
}
#[test]
fn test_reset_for_retry_paused_graph_ok() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[1].status = TaskStatus::Skipped;
graph.status = GraphStatus::Paused;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.status, GraphStatus::Running);
}
#[test]
fn test_reset_for_retry_clears_retry_count() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[0].retry_count = 5;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].retry_count, 0);
}
#[test]
fn test_reset_for_retry_paused_no_failed_tasks() {
let mut graph = graph_from_nodes(vec![make_node(0, &[])]);
graph.tasks[0].status = TaskStatus::Completed;
graph.status = GraphStatus::Paused;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.status, GraphStatus::Running);
assert_eq!(graph.tasks[0].status, TaskStatus::Completed);
}
#[test]
fn test_reset_for_retry_canceled_tasks_reset_to_pending() {
let mut graph = graph_from_nodes(vec![
make_node(0, &[]),
make_node(1, &[]),
make_node(2, &[0, 1]),
]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[1].status = TaskStatus::Canceled; graph.tasks[2].status = TaskStatus::Pending;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(
graph.tasks[1].status,
TaskStatus::Pending,
"Canceled task must be reset to Pending (IC2)"
);
assert_eq!(graph.tasks[2].status, TaskStatus::Pending);
}
#[test]
fn test_reset_for_retry_canceled_unblocks_dependents() {
let mut graph = graph_from_nodes(vec![make_node(0, &[]), make_node(1, &[0])]);
graph.tasks[0].status = TaskStatus::Failed;
graph.tasks[1].status = TaskStatus::Canceled;
graph.status = GraphStatus::Failed;
let __ra = make_rev_adj(&graph);
reset_for_retry(&mut graph, &__ra).unwrap();
assert_eq!(graph.tasks[0].status, TaskStatus::Ready);
assert_eq!(graph.tasks[1].status, TaskStatus::Pending);
}
fn make_node_titled(id: u32, deps: &[u32], title: &str, desc: &str) -> TaskNode {
let mut n = TaskNode::new(id, title, desc);
n.depends_on = deps.iter().map(|&d| TaskId(d)).collect();
n
}
#[test]
fn lookahead_depth_zero_returns_empty() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "web_search", "Search the web for results"),
make_node_titled(1, &[0], "summarize", "Summarize findings"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 0);
assert!(hints.is_empty(), "depth=0 must return empty vec");
}
#[test]
fn lookahead_depth_one_emits_only_direct_child() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "task-a", "Root task"),
make_node_titled(1, &[0], "web_search", "Search the web"),
make_node_titled(2, &[1], "summarize", "Summarize search results"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[2].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 1);
assert_eq!(hints.len(), 1, "depth=1 should emit only B");
assert_eq!(hints[0].tool_name, "web_search");
assert_eq!(hints[0].distance_from_current, 1);
}
#[test]
fn lookahead_depth_two_emits_both_children() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "task-a", "Root task"),
make_node_titled(1, &[0], "web_search", "Search the web"),
make_node_titled(2, &[1], "summarize", "Summarize search results"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[2].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 2);
assert_eq!(hints.len(), 2, "depth=2 should emit B and C");
assert_eq!(hints[0].tool_name, "web_search");
assert_eq!(hints[0].distance_from_current, 1);
assert_eq!(hints[1].tool_name, "summarize");
assert_eq!(hints[1].distance_from_current, 2);
}
#[test]
fn lookahead_no_frontier_returns_empty() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "task-a", "Root"),
make_node_titled(1, &[0], "task-b", "Child"),
]);
graph.tasks[0].status = TaskStatus::Pending;
graph.tasks[1].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 2);
assert!(hints.is_empty(), "no frontier → empty");
}
#[test]
fn lookahead_frontier_not_emitted() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "running-tool", "Currently executing"),
make_node_titled(1, &[0], "next-tool", "Next step"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 3);
assert!(
hints.iter().all(|h| h.tool_name != "running-tool"),
"frontier task must not be emitted"
);
assert_eq!(hints.len(), 1);
}
#[test]
fn lookahead_uses_agent_hint_as_tool_name() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "dispatch", "Root"),
make_node_titled(1, &[0], "raw-title", "Execute shell command"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[1].agent_hint = Some("shell_executor".to_string());
let hints = lookahead_tools(&graph, 1);
assert_eq!(hints.len(), 1);
assert_eq!(
hints[0].tool_name, "shell_executor",
"agent_hint should take precedence over title"
);
}
#[test]
fn lookahead_results_sorted_by_distance() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "root", "Root"),
make_node_titled(1, &[0], "step-one", "Step one"),
make_node_titled(2, &[1], "step-two", "Step two"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
graph.tasks[2].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 2);
for w in hints.windows(2) {
assert!(
w[0].distance_from_current <= w[1].distance_from_current,
"hints must be sorted by distance"
);
}
}
#[test]
fn lookahead_keywords_extracted_and_deduped() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "root", "Root task"),
make_node_titled(1, &[0], "search", "search search search results web"),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 1);
assert_eq!(hints.len(), 1);
let count = hints[0]
.keywords
.iter()
.filter(|k| k.as_str() == "search")
.count();
assert_eq!(count, 1, "duplicate keywords must be deduplicated");
}
#[test]
fn lookahead_stopwords_filtered() {
let mut graph = graph_from_nodes(vec![
make_node_titled(0, &[], "root", "Root"),
make_node_titled(
1,
&[0],
"task",
"the result of the operation from the source",
),
]);
graph.tasks[0].status = TaskStatus::Running;
graph.tasks[1].status = TaskStatus::Pending;
let hints = lookahead_tools(&graph, 1);
assert_eq!(hints.len(), 1);
for kw in &hints[0].keywords {
assert!(
!KEYWORD_STOPWORDS.contains(&kw.as_str()),
"stopword '{kw}' must not appear in keywords"
);
}
}
}