Deduplicate panic formatting, backoff construction, and stub graph allocation

- Extract format_panic_message() helper, used by both engine.rs and core_adapter
- Add RetryPolicy::DEFAULT_BACKOFF const, replacing 6 identical BackoffPolicy literals
- Cache stub graph via LazyLock to avoid per-call allocation in core_adapter handler
- Replace magic "success" string with StageStatus::Success.to_string()
- Bind graph.stall_timeout() once in run_via_core instead of calling twice

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-03-24 11:27:31 -04:00
parent 0df56a6c3e
commit c977a42bab
No known key found for this signature in database
4 changed files with 39 additions and 48 deletions

View file

@ -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<fabro_graphviz::graph::types::Graph> =
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<WorkflowGraph> 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<WorkflowGraph> for WorkflowNodeHandler {
handler,
gv_node,
&wf_context,
&wf_graph,
&STUB_GRAPH,
&run_dir,
&self.services,
);
@ -81,13 +84,7 @@ impl NodeHandler<WorkflowGraph> 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::<String>() {
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<WorkflowGraph> 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,

View file

@ -162,7 +162,7 @@ impl RunLifecycle<WorkflowGraph> 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,

View file

@ -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::<String>() {
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();

View file

@ -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<dyn Any + Send>) -> String {
if let Some(s) = payload.downcast_ref::<&str>() {
format!("handler panicked: {s}")
} else if let Some(s) = payload.downcast_ref::<String>() {
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(