diff --git a/lib/components/fabro-validate/src/rules/on_failure_valid.rs b/lib/components/fabro-validate/src/rules/on_failure_valid.rs index 8c84a0b72..683d7bb75 100644 --- a/lib/components/fabro-validate/src/rules/on_failure_valid.rs +++ b/lib/components/fabro-validate/src/rules/on_failure_valid.rs @@ -16,39 +16,40 @@ impl LintRule for Rule { fn apply(&self, graph: &Graph) -> Vec { let mut diagnostics = Vec::new(); - let expected_values = OnFailure::expected_values(); - match graph.attrs.get("on_failure") { - None => {} - Some(AttrValue::String(value)) if value.parse::().is_ok() => {} - Some(AttrValue::String(value)) => diagnostics.push(Diagnostic { + let invalid_value_message = match graph.attrs.get("on_failure") { + None => None, + Some(AttrValue::String(value)) if value.parse::().is_ok() => None, + Some(AttrValue::String(value)) => { + Some(format!("Graph has invalid on_failure value '{value}'")) + } + Some(_) => Some("Graph attribute 'on_failure' must be a string".to_string()), + }; + if let Some(message) = invalid_value_message { + diagnostics.push(Diagnostic { rule: self.name().to_string(), severity: Severity::Error, - message: format!("Graph has invalid on_failure value '{value}'"), - fix: Some(format!("Use one of: {expected_values}")), + message, + fix: Some(format!("Use one of: {}", OnFailure::expected_values())), ..Diagnostic::default() - }), - Some(_) => diagnostics.push(Diagnostic { - rule: self.name().to_string(), - severity: Severity::Error, - message: "Graph attribute 'on_failure' must be a string".to_string(), - fix: Some(format!("Use one of: {expected_values}")), - ..Diagnostic::default() - }), + }); } + let misplaced = |subject: String| Diagnostic { + rule: self.name().to_string(), + severity: Severity::Warning, + message: format!( + "{subject} sets 'on_failure', which has no effect outside graph scope" + ), + fix: Some("Move 'on_failure' to the graph attributes".to_string()), + ..Diagnostic::default() + }; + for node in graph.nodes.values() { if node.attrs.contains_key("on_failure") { diagnostics.push(Diagnostic { - rule: self.name().to_string(), - severity: Severity::Warning, - message: format!( - "Node '{}' sets 'on_failure', which has no effect outside graph scope", - node.id - ), node_id: Some(node.id.clone()), - fix: Some("Move 'on_failure' to the graph attributes".to_string()), - ..Diagnostic::default() + ..misplaced(format!("Node '{}'", node.id)) }); } } @@ -56,15 +57,8 @@ impl LintRule for Rule { for edge in &graph.edges { if edge.attrs.contains_key("on_failure") { diagnostics.push(Diagnostic { - rule: self.name().to_string(), - severity: Severity::Warning, - message: format!( - "Edge '{} -> {}' sets 'on_failure', which has no effect outside graph scope", - edge.from, edge.to - ), edge: Some((edge.from.clone(), edge.to.clone())), - fix: Some("Move 'on_failure' to the graph attributes".to_string()), - ..Diagnostic::default() + ..misplaced(format!("Edge '{} -> {}'", edge.from, edge.to)) }); } } @@ -75,10 +69,10 @@ impl LintRule for Rule { #[cfg(test)] mod tests { - use fabro_graphviz::graph::{AttrValue, Edge, Node}; + use fabro_graphviz::graph::{AttrValue, Edge}; use super::Rule; - use crate::rules::test_support::minimal_graph; + use crate::rules::test_support::{minimal_graph, node_with_attrs}; use crate::{LintRule, Severity}; #[test] @@ -137,12 +131,10 @@ mod tests { #[test] fn warns_for_node_and_edge_placement() { let mut graph = minimal_graph(); - let mut node = Node::new("work"); - node.attrs.insert( - "on_failure".to_string(), - AttrValue::String("exit".to_string()), + graph.nodes.insert( + "work".to_string(), + node_with_attrs("work", &[("on_failure", "exit")]), ); - graph.nodes.insert("work".to_string(), node); let mut edge = Edge::new("start", "work"); edge.attrs.insert( diff --git a/lib/components/fabro-workflow/tests/it/integration.rs b/lib/components/fabro-workflow/tests/it/integration.rs index 1fdd577df..7786dd49d 100644 --- a/lib/components/fabro-workflow/tests/it/integration.rs +++ b/lib/components/fabro-workflow/tests/it/integration.rs @@ -1150,34 +1150,21 @@ impl Handler for OnFailureRecordingHandler { } } -fn on_failure_graph(policy: Option<&str>) -> Graph { - let mut graph = Graph::new("OnFailureTest"); - if let Some(policy) = policy { - graph.attrs.insert( - "on_failure".to_string(), - AttrValue::String(policy.to_string()), - ); - } - - let mut start = Node::new("start"); - start.attrs.insert( - "shape".to_string(), - AttrValue::String("Mdiamond".to_string()), +/// Builds the linear on_failure test graph, splicing `extra` statements +/// (policy attribute, recovery edges, node attributes) into the DOT source so +/// tests exercise the real parser path for the `on_failure` attribute. +fn on_failure_graph(extra: &str) -> Graph { + let input = format!( + r"digraph OnFailureTest {{ + {extra} + start [shape=Mdiamond] + exit [shape=Msquare] + work + downstream + start -> work -> downstream -> exit + }}" ); - graph.nodes.insert("start".to_string(), start); - - let mut exit = Node::new("exit"); - exit.attrs.insert( - "shape".to_string(), - AttrValue::String("Msquare".to_string()), - ); - graph.nodes.insert("exit".to_string(), exit); - - for node_id in ["work", "downstream", "recovery"] { - graph.nodes.insert(node_id.to_string(), Node::new(node_id)); - } - graph.edges.push(Edge::new("start", "work")); - graph + parse(&input).expect("on_failure test graph should parse") } fn on_failure_registry(visits: Arc>>) -> HandlerRegistry { @@ -1187,36 +1174,51 @@ fn on_failure_registry(visits: Arc>>) -> HandlerReg registry } -#[tokio::test] -async fn on_failure_exit_stops_linear_workflow_and_records_failed_lifecycle() { - let mut graph = on_failure_graph(Some("exit")); - graph.edges.push(Edge::new("work", "downstream")); - graph.edges.push(Edge::new("downstream", "exit")); +struct OnFailureRun { + outcome: Outcome, + state: fabro_store::RunProjection, + visits: Arc>>, + _run_dir: tempfile::TempDir, +} +async fn run_on_failure(graph: &Graph, emitter: Emitter) -> OnFailureRun { let visits = Arc::new(std::sync::Mutex::new(Vec::new())); - let emitter = Emitter::default(); - let events = collect_events(&emitter); let engine = WorkflowRunner::new( on_failure_registry(Arc::clone(&visits)), Arc::new(emitter), local_env(), ); - let dir = tempfile::tempdir().unwrap(); + let run_dir = tempfile::tempdir().expect("temporary run dir should be created"); let (outcome, state) = engine - .run_with_state(&graph, &make_run_options(dir.path())) + .run_with_state(graph, &make_run_options(run_dir.path())) .await - .expect("policy termination should return a failed workflow outcome"); + .expect("on_failure run should complete without engine errors"); + OnFailureRun { + outcome, + state, + visits, + _run_dir: run_dir, + } +} - assert_eq!(outcome.status, StageOutcome::Failed { +#[tokio::test] +async fn on_failure_exit_stops_linear_workflow_and_records_failed_lifecycle() { + let graph = on_failure_graph(r#"graph [on_failure="exit"]"#); + let emitter = Emitter::default(); + let events = collect_events(&emitter); + let run = run_on_failure(&graph, emitter).await; + + assert_eq!(run.outcome.status, StageOutcome::Failed { retry_requested: false, }); assert_eq!( - outcome.failure_reason(), + run.outcome.failure_reason(), Some("stage work failed and graph on_failure=exit stopped routing") ); - assert_eq!(*visits.lock().unwrap(), vec!["work"]); + assert_eq!(*run.visits.lock().unwrap(), vec!["work"]); - let checkpoint = state + let checkpoint = run + .state .current_checkpoint() .expect("failed work should be checkpointed"); assert_eq!(checkpoint.current_node, "work"); @@ -1246,27 +1248,14 @@ async fn on_failure_exit_stops_linear_workflow_and_records_failed_lifecycle() { #[tokio::test] async fn on_failure_route_and_absent_policy_preserve_unconditional_fallback() { - for policy in [None, Some("route")] { - let mut graph = on_failure_graph(policy); - graph.edges.push(Edge::new("work", "downstream")); - graph.edges.push(Edge::new("downstream", "exit")); + for policy_attr in ["", r#"graph [on_failure="route"]"#] { + let graph = on_failure_graph(policy_attr); + let run = run_on_failure(&graph, Emitter::default()).await; - let visits = Arc::new(std::sync::Mutex::new(Vec::new())); - let engine = WorkflowRunner::new( - on_failure_registry(Arc::clone(&visits)), - Arc::new(Emitter::default()), - local_env(), - ); - let dir = tempfile::tempdir().unwrap(); - let (outcome, state) = engine - .run_with_state(&graph, &make_run_options(dir.path())) - .await - .expect("route policy should preserve unconditional fallback"); - - assert_eq!(outcome.status, StageOutcome::Succeeded); - assert_eq!(*visits.lock().unwrap(), vec!["work", "downstream"]); + assert_eq!(run.outcome.status, StageOutcome::Succeeded); + assert_eq!(*run.visits.lock().unwrap(), vec!["work", "downstream"]); assert!( - state + run.state .current_checkpoint() .expect("downstream should be checkpointed") .node_outcomes @@ -1277,58 +1266,30 @@ async fn on_failure_route_and_absent_policy_preserve_unconditional_fallback() { #[tokio::test] async fn on_failure_exit_allows_explicit_failure_recovery_edge() { - let mut graph = on_failure_graph(Some("exit")); - graph.edges.push(Edge::new("work", "downstream")); - let mut recovery_edge = Edge::new("work", "recovery"); - recovery_edge.attrs.insert( - "condition".to_string(), - AttrValue::String("outcome=failed".to_string()), + let graph = on_failure_graph( + r#"graph [on_failure="exit"] + recovery + work -> recovery [condition="outcome=failed"] + recovery -> exit"#, ); - graph.edges.push(recovery_edge); - graph.edges.push(Edge::new("recovery", "exit")); - graph.edges.push(Edge::new("downstream", "exit")); + let run = run_on_failure(&graph, Emitter::default()).await; - let visits = Arc::new(std::sync::Mutex::new(Vec::new())); - let engine = WorkflowRunner::new( - on_failure_registry(Arc::clone(&visits)), - Arc::new(Emitter::default()), - local_env(), - ); - let dir = tempfile::tempdir().unwrap(); - let (outcome, _) = engine - .run_with_state(&graph, &make_run_options(dir.path())) - .await - .expect("explicit failure recovery should complete"); - - assert_eq!(outcome.status, StageOutcome::Succeeded); - assert_eq!(*visits.lock().unwrap(), vec!["work", "recovery"]); + assert_eq!(run.outcome.status, StageOutcome::Succeeded); + assert_eq!(*run.visits.lock().unwrap(), vec!["work", "recovery"]); } #[tokio::test] async fn on_failure_exit_uses_retry_target_instead_of_unconditional_edge() { - let mut graph = on_failure_graph(Some("exit")); - graph.nodes.get_mut("work").unwrap().attrs.insert( - "retry_target".to_string(), - AttrValue::String("recovery".to_string()), + let graph = on_failure_graph( + r#"graph [on_failure="exit"] + work [retry_target="recovery"] + recovery + recovery -> exit"#, ); - graph.edges.push(Edge::new("work", "downstream")); - graph.edges.push(Edge::new("recovery", "exit")); - graph.edges.push(Edge::new("downstream", "exit")); + let run = run_on_failure(&graph, Emitter::default()).await; - let visits = Arc::new(std::sync::Mutex::new(Vec::new())); - let engine = WorkflowRunner::new( - on_failure_registry(Arc::clone(&visits)), - Arc::new(Emitter::default()), - local_env(), - ); - let dir = tempfile::tempdir().unwrap(); - let (outcome, _) = engine - .run_with_state(&graph, &make_run_options(dir.path())) - .await - .expect("retry target should run before policy termination"); - - assert_eq!(outcome.status, StageOutcome::Succeeded); - assert_eq!(*visits.lock().unwrap(), vec!["work", "recovery"]); + assert_eq!(run.outcome.status, StageOutcome::Succeeded); + assert_eq!(*run.visits.lock().unwrap(), vec!["work", "recovery"]); } #[tokio::test] diff --git a/lib/foundation/fabro-core/src/executor.rs b/lib/foundation/fabro-core/src/executor.rs index 9aba44156..d1d8a784a 100644 --- a/lib/foundation/fabro-core/src/executor.rs +++ b/lib/foundation/fabro-core/src/executor.rs @@ -2260,40 +2260,6 @@ mod tests { assert_eq!(log.run_end.as_ref(), Some(&outcome)); } - #[tokio::test] - async fn executor_exit_policy_uses_retry_target_before_termination() { - let handler = Arc::new(CountingHandler::new(vec![ - Ok(Outcome::fail("boom")), - Ok(Outcome::success()), - ])); - let graph = TestGraph::new( - vec![ - TestNode::new("work"), - TestNode::new("downstream"), - TestNode::new("recovery"), - TestNode::terminal("end"), - ], - vec![ - TestEdge::new("work", "downstream"), - TestEdge::new("downstream", "end"), - TestEdge::new("recovery", "end"), - ], - "work", - ) - .with_retry_target("work", "recovery") - .with_on_failure(OnFailure::Exit); - let state = ExecutionState::new(&graph).unwrap(); - let executor = - ExecutorBuilder::new(Arc::clone(&handler) as Arc>).build(); - - let (outcome, state) = executor.run(&graph, state).await.unwrap(); - - assert_eq!(outcome.status, StageOutcome::Succeeded); - assert_eq!(handler.calls(), 2); - assert!(state.node_outcomes.contains_key("recovery")); - assert!(!state.node_outcomes.contains_key("downstream")); - } - #[tokio::test] async fn executor_goal_gate_retry_target_to_terminal_fails_without_looping() { let terminal_visits = Arc::new(AtomicU32::new(0));