diff --git a/lib/crates/fabro-workflows/src/core_adapter/handler.rs b/lib/crates/fabro-workflows/src/core_adapter/handler.rs index 81e9b37e8..e0fd7940a 100644 --- a/lib/crates/fabro-workflows/src/core_adapter/handler.rs +++ b/lib/crates/fabro-workflows/src/core_adapter/handler.rs @@ -1,6 +1,6 @@ use std::panic::AssertUnwindSafe; use std::path::PathBuf; -use std::sync::Arc; +use std::sync::{Arc, LazyLock}; use async_trait::async_trait; use futures::FutureExt; @@ -14,9 +14,13 @@ use fabro_core::retry::RetryPolicy as CoreRetryPolicy; use super::graph::WorkflowGraph; use super::WorkflowNode; use crate::engine; -use crate::handler::EngineServices; +use crate::handler::{format_panic_message, EngineServices}; use crate::outcome::{Outcome, StageStatus}; +/// Cached stub graph for handler dispatch (avoids allocating on every call). +static STUB_GRAPH: LazyLock = + LazyLock::new(|| fabro_graphviz::graph::types::Graph::new("stub")); + /// Production node handler that bridges fabro-core's NodeHandler to the /// existing fabro-workflows Handler trait via EngineServices. pub struct WorkflowNodeHandler { @@ -36,7 +40,6 @@ impl NodeHandler for WorkflowNodeHandler { let handler = self.services.registry.resolve(gv_node); let wf_context = crate::context::Context::new(); - let wf_graph = fabro_graphviz::graph::types::Graph::new("stub"); // Timeout from the node let node_timeout = gv_node.timeout(); @@ -47,7 +50,7 @@ impl NodeHandler for WorkflowNodeHandler { handler, gv_node, &wf_context, - &wf_graph, + &STUB_GRAPH, &run_dir, &self.services, ); @@ -81,13 +84,7 @@ impl NodeHandler for WorkflowNodeHandler { })) } Err(panic_payload) => { - let msg = if let Some(s) = panic_payload.downcast_ref::<&str>() { - format!("handler panicked: {s}") - } else if let Some(s) = panic_payload.downcast_ref::() { - format!("handler panicked: {s}") - } else { - "handler panicked".to_string() - }; + let msg = format_panic_message(panic_payload); Err(CoreError::handler(HandlerErrorDetail { message: msg, retryable: false, @@ -100,8 +97,7 @@ impl NodeHandler for WorkflowNodeHandler { fn retry_policy(&self, node: &WorkflowNode, _graph: &WorkflowGraph) -> CoreRetryPolicy { let gv_node = node.inner(); - let gv_graph = fabro_graphviz::graph::types::Graph::new("stub"); - let wf_policy = engine::build_retry_policy(gv_node, &gv_graph); + let wf_policy = engine::build_retry_policy(gv_node, &STUB_GRAPH); CoreRetryPolicy { max_attempts: wf_policy.max_attempts, backoff: wf_policy.backoff, diff --git a/lib/crates/fabro-workflows/src/core_adapter/lifecycle.rs b/lib/crates/fabro-workflows/src/core_adapter/lifecycle.rs index bc794d3a6..727312434 100644 --- a/lib/crates/fabro-workflows/src/core_adapter/lifecycle.rs +++ b/lib/crates/fabro-workflows/src/core_adapter/lifecycle.rs @@ -162,7 +162,7 @@ impl RunLifecycle for WorkflowLifecycle { name: gv.label().to_string(), index: stage_index, duration_ms: 0, - status: "success".to_string(), + status: StageStatus::Success.to_string(), preferred_label: None, suggested_next_ids: Vec::new(), usage: None, diff --git a/lib/crates/fabro-workflows/src/engine.rs b/lib/crates/fabro-workflows/src/engine.rs index 9d673e2e6..fb00c5ae5 100644 --- a/lib/crates/fabro-workflows/src/engine.rs +++ b/lib/crates/fabro-workflows/src/engine.rs @@ -76,17 +76,19 @@ pub struct RetryPolicy { } impl RetryPolicy { + const DEFAULT_BACKOFF: BackoffPolicy = BackoffPolicy { + initial_delay: Duration::from_millis(5_000), + factor: 2.0, + max_delay: Duration::from_millis(60_000), + jitter: true, + }; + /// No retries -- fail immediately. #[must_use] pub fn none() -> Self { Self { max_attempts: 1, - backoff: BackoffPolicy { - initial_delay: Duration::from_millis(5_000), - factor: 2.0, - max_delay: Duration::from_millis(60_000), - jitter: true, - }, + backoff: Self::DEFAULT_BACKOFF, } } @@ -95,12 +97,7 @@ impl RetryPolicy { pub fn standard() -> Self { Self { max_attempts: 5, - backoff: BackoffPolicy { - initial_delay: Duration::from_millis(5_000), - factor: 2.0, - max_delay: Duration::from_millis(60_000), - jitter: true, - }, + backoff: Self::DEFAULT_BACKOFF, } } @@ -111,9 +108,7 @@ impl RetryPolicy { max_attempts: 5, backoff: BackoffPolicy { initial_delay: Duration::from_millis(500), - factor: 2.0, - max_delay: Duration::from_millis(60_000), - jitter: true, + ..Self::DEFAULT_BACKOFF }, } } @@ -126,8 +121,7 @@ impl RetryPolicy { backoff: BackoffPolicy { initial_delay: Duration::from_millis(500), factor: 1.0, - max_delay: Duration::from_millis(60_000), - jitter: true, + ..Self::DEFAULT_BACKOFF }, } } @@ -140,8 +134,7 @@ impl RetryPolicy { backoff: BackoffPolicy { initial_delay: Duration::from_millis(2000), factor: 3.0, - max_delay: Duration::from_millis(60_000), - jitter: true, + ..Self::DEFAULT_BACKOFF }, } } @@ -168,12 +161,7 @@ pub(crate) fn build_retry_policy(node: &Node, graph: &Graph) -> RetryPolicy { let max_attempts = u32::try_from(max_retries + 1).unwrap_or(1).max(1); RetryPolicy { max_attempts, - backoff: BackoffPolicy { - initial_delay: Duration::from_millis(5_000), - factor: 2.0, - max_delay: Duration::from_millis(60_000), - jitter: true, - }, + backoff: RetryPolicy::DEFAULT_BACKOFF, } } @@ -1048,13 +1036,7 @@ impl WorkflowRunEngine { match timed_result { Ok(r) => r, Err(panic_payload) => { - let msg = if let Some(s) = panic_payload.downcast_ref::<&str>() { - format!("handler panicked: {s}") - } else if let Some(s) = panic_payload.downcast_ref::() { - format!("handler panicked: {s}") - } else { - "handler panicked".to_string() - }; + let msg = crate::handler::format_panic_message(panic_payload); let panic_dir = node_dir(run_dir, &node.id, visit); let _ = std::fs::create_dir_all(&panic_dir); let _ = std::fs::write(panic_dir.join("panic.txt"), &msg); @@ -1529,9 +1511,10 @@ impl WorkflowRunEngine { }; // Set up stall watchdog - let stall_token = graph.stall_timeout().map(|_| CancellationToken::new()); + let stall_timeout_opt = graph.stall_timeout(); + let stall_token = stall_timeout_opt.map(|_| CancellationToken::new()); let stall_shutdown = - if let (Some(stall_timeout), Some(ref token)) = (graph.stall_timeout(), &stall_token) { + if let (Some(stall_timeout), Some(ref token)) = (stall_timeout_opt, &stall_token) { let shutdown = CancellationToken::new(); let emitter = self.services.emitter.clone(); let token_clone = token.clone(); diff --git a/lib/crates/fabro-workflows/src/handler/mod.rs b/lib/crates/fabro-workflows/src/handler/mod.rs index e704951c7..449d7b4b9 100644 --- a/lib/crates/fabro-workflows/src/handler/mod.rs +++ b/lib/crates/fabro-workflows/src/handler/mod.rs @@ -10,6 +10,7 @@ pub mod prompt; pub mod start; pub mod wait; +use std::any::Any; use std::collections::HashMap; use std::path::Path; use std::sync::Arc; @@ -111,6 +112,17 @@ pub trait Handler: Send + Sync { } } +/// Extract a human-readable message from a panic payload. +pub(crate) fn format_panic_message(payload: Box) -> String { + if let Some(s) = payload.downcast_ref::<&str>() { + format!("handler panicked: {s}") + } else if let Some(s) = payload.downcast_ref::() { + format!("handler panicked: {s}") + } else { + "handler panicked".to_string() + } +} + /// Route to [`Handler::simulate`] when `services.dry_run` is true, otherwise /// [`Handler::execute`]. pub async fn dispatch_handler(