mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
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:
parent
ef583c35ac
commit
5597040354
4 changed files with 39 additions and 48 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue