mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-09-07 08:27:12 +00:00
fix: resolve fabro-workflow nextest timeouts and regressions
This commit is contained in:
parent
d490dbe4fa
commit
efa4670d9e
11 changed files with 359 additions and 46 deletions
|
|
@ -89,7 +89,7 @@ impl SlateRunStore {
|
|||
self.inner.record.clone()
|
||||
}
|
||||
|
||||
pub(crate) fn created_at(&self) -> DateTime<Utc> {
|
||||
pub fn created_at(&self) -> DateTime<Utc> {
|
||||
self.inner.record.created_at
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -189,6 +189,8 @@ pub enum WorkflowRunEvent {
|
|||
delay_ms: u64,
|
||||
},
|
||||
ParallelStarted {
|
||||
node_id: String,
|
||||
visit: u32,
|
||||
branch_count: usize,
|
||||
join_policy: String,
|
||||
},
|
||||
|
|
@ -205,6 +207,8 @@ pub enum WorkflowRunEvent {
|
|||
head_sha: Option<String>,
|
||||
},
|
||||
ParallelCompleted {
|
||||
node_id: String,
|
||||
visit: u32,
|
||||
duration_ms: u64,
|
||||
success_count: usize,
|
||||
failure_count: usize,
|
||||
|
|
@ -680,6 +684,7 @@ impl WorkflowRunEvent {
|
|||
Self::ParallelStarted {
|
||||
branch_count,
|
||||
join_policy,
|
||||
..
|
||||
} => {
|
||||
debug!(branch_count, join_policy, "Parallel execution started");
|
||||
}
|
||||
|
|
@ -703,6 +708,7 @@ impl WorkflowRunEvent {
|
|||
success_count,
|
||||
failure_count,
|
||||
results,
|
||||
..
|
||||
} => {
|
||||
debug!(
|
||||
duration_ms,
|
||||
|
|
@ -1302,6 +1308,8 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields {
|
|||
| WorkflowRunEvent::SubgraphCompleted { .. }
|
||||
| WorkflowRunEvent::AssetCaptured { .. }
|
||||
| WorkflowRunEvent::PromptCompleted { .. }
|
||||
| WorkflowRunEvent::ParallelStarted { .. }
|
||||
| WorkflowRunEvent::ParallelCompleted { .. }
|
||||
| WorkflowRunEvent::CommandStarted { .. }
|
||||
| WorkflowRunEvent::CommandCompleted { .. }
|
||||
| WorkflowRunEvent::AgentCliStarted { .. }
|
||||
|
|
|
|||
|
|
@ -4,7 +4,10 @@ use std::sync::Arc;
|
|||
|
||||
use async_trait::async_trait;
|
||||
use fabro_agent::Sandbox;
|
||||
use fabro_model::Provider;
|
||||
use fabro_store::NodeVisitRef;
|
||||
use fabro_types::RunId;
|
||||
use tokio::fs;
|
||||
|
||||
use crate::context::keys;
|
||||
use crate::context::{Context, WorkflowContext};
|
||||
|
|
@ -13,6 +16,7 @@ use crate::event::EventEmitter;
|
|||
use crate::outcome::{
|
||||
FailureCategory, FailureDetail, Outcome, OutcomeExt, StageStatus, StageUsage,
|
||||
};
|
||||
use crate::run_dir::visit_from_context;
|
||||
use crate::transforms::variable_expansion::expand_vars;
|
||||
use fabro_graphviz::graph::{Graph, Node};
|
||||
|
||||
|
|
@ -195,6 +199,91 @@ pub(crate) fn truncate(s: &str, max_chars: usize) -> &str {
|
|||
}
|
||||
}
|
||||
|
||||
pub(crate) fn stage_dir(run_dir: &Path, node_id: &str, visit: u32) -> std::path::PathBuf {
|
||||
let node_dir = if visit <= 1 {
|
||||
node_id.to_string()
|
||||
} else {
|
||||
format!("{node_id}-visit_{visit}")
|
||||
};
|
||||
run_dir.join("nodes").join(node_dir)
|
||||
}
|
||||
|
||||
pub(crate) fn status_json_value(outcome: &Outcome) -> serde_json::Value {
|
||||
let mut status = serde_json::Map::new();
|
||||
status.insert(
|
||||
"status".to_string(),
|
||||
serde_json::Value::String(outcome.status.to_string()),
|
||||
);
|
||||
status.insert(
|
||||
"outcome".to_string(),
|
||||
serde_json::Value::String(outcome.status.to_string()),
|
||||
);
|
||||
if let Some(label) = outcome.preferred_label.as_ref() {
|
||||
status.insert(
|
||||
"preferred_next_label".to_string(),
|
||||
serde_json::Value::String(label.clone()),
|
||||
);
|
||||
}
|
||||
if !outcome.suggested_next_ids.is_empty() {
|
||||
status.insert(
|
||||
"suggested_next_ids".to_string(),
|
||||
serde_json::json!(outcome.suggested_next_ids),
|
||||
);
|
||||
}
|
||||
if let Some(failure) = outcome.failure.as_ref() {
|
||||
status.insert(
|
||||
"failure_reason".to_string(),
|
||||
serde_json::Value::String(failure.message.clone()),
|
||||
);
|
||||
}
|
||||
if !outcome.context_updates.is_empty() {
|
||||
status.insert(
|
||||
"context_updates".to_string(),
|
||||
serde_json::Value::Object(
|
||||
outcome
|
||||
.context_updates
|
||||
.iter()
|
||||
.map(|(key, value)| (key.clone(), value.clone()))
|
||||
.collect(),
|
||||
),
|
||||
);
|
||||
}
|
||||
serde_json::Value::Object(status)
|
||||
}
|
||||
|
||||
pub(crate) async fn write_provider_used_file(
|
||||
services: &EngineServices,
|
||||
stage_dir: &Path,
|
||||
node_id: &str,
|
||||
visit: u32,
|
||||
fallback: Option<serde_json::Value>,
|
||||
) -> Result<(), FabroError> {
|
||||
let provider_used = services
|
||||
.run_store
|
||||
.state()
|
||||
.await
|
||||
.ok()
|
||||
.and_then(|state| {
|
||||
let node = NodeVisitRef { node_id, visit };
|
||||
state
|
||||
.node(&node)
|
||||
.and_then(|node_state| node_state.provider_used.clone())
|
||||
})
|
||||
.or(fallback);
|
||||
let Some(provider_used) = provider_used else {
|
||||
return Ok(());
|
||||
};
|
||||
fs::write(
|
||||
stage_dir.join("provider_used.json"),
|
||||
serde_json::to_vec_pretty(&provider_used).map_err(|err| {
|
||||
FabroError::handler(format!("failed to serialize provider_used.json: {err}"))
|
||||
})?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write provider_used.json: {err}")))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Shared simulate implementation for LLM-backed handlers (agent & prompt).
|
||||
/// Produces a simulated outcome with standard context updates.
|
||||
pub(crate) fn simulate_llm_handler(node: &Node) -> Outcome {
|
||||
|
|
@ -248,6 +337,15 @@ impl Handler for AgentHandler {
|
|||
format!("{preamble}\n\n{expanded}")
|
||||
};
|
||||
|
||||
let visit = u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX);
|
||||
let stage_dir = stage_dir(_run_dir, &node.id, visit);
|
||||
fs::create_dir_all(&stage_dir)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?;
|
||||
fs::write(stage_dir.join("prompt.md"), &prompt)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write prompt file: {err}")))?;
|
||||
|
||||
// 3. Call LLM backend (agent loop)
|
||||
let thread_id = context.thread_id();
|
||||
let run_id = context
|
||||
|
|
@ -350,6 +448,31 @@ impl Handler for AgentHandler {
|
|||
}
|
||||
outcome.usage = stage_usage;
|
||||
outcome.files_touched = backend_files_touched;
|
||||
fs::write(stage_dir.join("response.md"), &response_text)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write response file: {err}")))?;
|
||||
fs::write(
|
||||
stage_dir.join("status.json"),
|
||||
serde_json::to_vec_pretty(&status_json_value(&outcome))
|
||||
.map_err(|err| FabroError::handler(format!("failed to serialize status: {err}")))?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write status file: {err}")))?;
|
||||
write_provider_used_file(
|
||||
services,
|
||||
&stage_dir,
|
||||
&node.id,
|
||||
visit,
|
||||
Some(serde_json::json!({
|
||||
"mode": if node.backend() == Some("cli") { "cli" } else { "agent" },
|
||||
"provider": node
|
||||
.provider()
|
||||
.map(String::from)
|
||||
.unwrap_or_else(|| Provider::default_from_env().as_str().to_string()),
|
||||
"model": node.model().map(String::from).unwrap_or_default(),
|
||||
})),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(outcome)
|
||||
}
|
||||
|
|
@ -390,6 +513,7 @@ mod tests {
|
|||
.await
|
||||
.unwrap();
|
||||
let services = EngineServices {
|
||||
emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)),
|
||||
run_store: run_store.clone(),
|
||||
..EngineServices::test_default()
|
||||
};
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ use crate::event::WorkflowRunEvent;
|
|||
use crate::outcome::{Outcome, OutcomeExt};
|
||||
use async_trait::async_trait;
|
||||
use fabro_graphviz::graph::{Graph, Node};
|
||||
use tokio::fs;
|
||||
|
||||
use super::{EngineServices, Handler};
|
||||
|
||||
|
|
@ -59,7 +60,7 @@ impl Handler for CommandHandler {
|
|||
node: &Node,
|
||||
_context: &Context,
|
||||
_graph: &Graph,
|
||||
_run_dir: &Path,
|
||||
run_dir: &Path,
|
||||
services: &EngineServices,
|
||||
) -> Result<Outcome, FabroError> {
|
||||
let script = node
|
||||
|
|
@ -97,6 +98,22 @@ impl Handler for CommandHandler {
|
|||
} else {
|
||||
script.to_string()
|
||||
};
|
||||
let stage_dir = run_dir.join("nodes").join(&node.id);
|
||||
fs::create_dir_all(&stage_dir)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?;
|
||||
fs::write(
|
||||
stage_dir.join("script_invocation.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::json!({
|
||||
"script": script,
|
||||
"command": command,
|
||||
"language": language,
|
||||
"timeout_ms": timeout_ms(node),
|
||||
}))
|
||||
.map_err(|err| FabroError::handler(format!("failed to serialize invocation: {err}")))?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write invocation file: {err}")))?;
|
||||
|
||||
let timeout_ms = node
|
||||
.timeout()
|
||||
|
|
@ -121,6 +138,23 @@ impl Handler for CommandHandler {
|
|||
duration_ms: result.duration_ms,
|
||||
timed_out: result.timed_out,
|
||||
});
|
||||
fs::write(stage_dir.join("stdout.log"), &result.stdout)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write stdout log: {err}")))?;
|
||||
fs::write(stage_dir.join("stderr.log"), &result.stderr)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write stderr log: {err}")))?;
|
||||
fs::write(
|
||||
stage_dir.join("script_timing.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::json!({
|
||||
"duration_ms": result.duration_ms,
|
||||
"exit_code": (!result.timed_out).then_some(result.exit_code),
|
||||
"timed_out": result.timed_out,
|
||||
}))
|
||||
.map_err(|err| FabroError::handler(format!("failed to serialize timing: {err}")))?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write timing file: {err}")))?;
|
||||
|
||||
if result.timed_out {
|
||||
return Err(FabroError::handler(format!(
|
||||
|
|
|
|||
|
|
@ -86,12 +86,22 @@ impl EngineServices {
|
|||
sandbox: Arc::new(fabro_agent::LocalSandbox::new(
|
||||
std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
|
||||
)),
|
||||
run_store: futures::executor::block_on(async {
|
||||
store
|
||||
.create_run(&fabro_types::RunId::new(), chrono::Utc::now(), None)
|
||||
.await
|
||||
.expect("slate-backed test run store should initialize")
|
||||
}),
|
||||
// Build the test run store on a dedicated runtime so this helper
|
||||
// remains safe to call from both sync tests and #[tokio::test].
|
||||
run_store: std::thread::spawn(move || {
|
||||
tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
.expect("test runtime should initialize")
|
||||
.block_on(async {
|
||||
store
|
||||
.create_run(&fabro_types::RunId::new(), chrono::Utc::now(), None)
|
||||
.await
|
||||
.expect("slate-backed test run store should initialize")
|
||||
})
|
||||
})
|
||||
.join()
|
||||
.expect("test run store thread should join"),
|
||||
git_state: std::sync::RwLock::new(None),
|
||||
hook_runner: None,
|
||||
env: HashMap::new(),
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ use std::time::Instant;
|
|||
use async_trait::async_trait;
|
||||
use fabro_agent::{Sandbox, WorktreeOptions, WorktreeSandbox};
|
||||
use fabro_types::RunId;
|
||||
use tokio::fs;
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
use crate::context::keys;
|
||||
|
|
@ -151,6 +152,8 @@ impl Handler for ParallelHandler {
|
|||
);
|
||||
|
||||
services.emitter.emit(&WorkflowRunEvent::ParallelStarted {
|
||||
node_id: node.id.clone(),
|
||||
visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX),
|
||||
branch_count: branches.len(),
|
||||
join_policy: join_policy.to_string(),
|
||||
});
|
||||
|
|
@ -473,10 +476,26 @@ impl Handler for ParallelHandler {
|
|||
entry
|
||||
})
|
||||
.collect();
|
||||
let stage_dir = run_dir.join("nodes").join(&node.id);
|
||||
fs::create_dir_all(&stage_dir)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?;
|
||||
fs::write(
|
||||
stage_dir.join("parallel_results.json"),
|
||||
serde_json::to_vec_pretty(&results_json).map_err(|err| {
|
||||
FabroError::handler(format!("failed to serialize parallel results: {err}"))
|
||||
})?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| {
|
||||
FabroError::handler(format!("failed to write parallel results file: {err}"))
|
||||
})?;
|
||||
context.set(keys::PARALLEL_RESULTS, serde_json::json!(results_json));
|
||||
context.set(keys::PARALLEL_BRANCH_COUNT, serde_json::json!(total));
|
||||
|
||||
services.emitter.emit(&WorkflowRunEvent::ParallelCompleted {
|
||||
node_id: node.id.clone(),
|
||||
visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX),
|
||||
duration_ms: millis_u64(parallel_start.elapsed()),
|
||||
success_count,
|
||||
failure_count: fail_count,
|
||||
|
|
@ -678,9 +697,12 @@ mod tests {
|
|||
.await
|
||||
.unwrap();
|
||||
let services = EngineServices {
|
||||
emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)),
|
||||
run_store: run_store.clone(),
|
||||
..EngineServices::test_default()
|
||||
};
|
||||
let logger = crate::event::StoreProgressLogger::new(run_store.clone());
|
||||
logger.register(services.emitter.as_ref());
|
||||
let mut node = Node::new("par");
|
||||
node.attrs.insert(
|
||||
"shape".to_string(),
|
||||
|
|
@ -703,6 +725,7 @@ mod tests {
|
|||
.execute(&node, &context, &graph, tmp.path(), &services)
|
||||
.await
|
||||
.unwrap();
|
||||
logger.flush().await;
|
||||
|
||||
let state = run_store.state().await.unwrap();
|
||||
let node_state = state
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
use std::path::Path;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use tokio::fs;
|
||||
|
||||
use crate::context::keys;
|
||||
use crate::context::{Context, WorkflowContext};
|
||||
|
|
@ -12,7 +13,8 @@ use fabro_graphviz::graph::{Graph, Node};
|
|||
use fabro_model::Provider;
|
||||
|
||||
use super::agent::{
|
||||
CodergenBackend, CodergenResult, expand_variables, extract_status_fields, truncate,
|
||||
CodergenBackend, CodergenResult, expand_variables, extract_status_fields, stage_dir,
|
||||
status_json_value, truncate, write_provider_used_file,
|
||||
};
|
||||
use super::{EngineServices, Handler};
|
||||
|
||||
|
|
@ -46,7 +48,7 @@ impl Handler for PromptHandler {
|
|||
node: &Node,
|
||||
context: &Context,
|
||||
graph: &Graph,
|
||||
_run_dir: &Path,
|
||||
run_dir: &Path,
|
||||
services: &EngineServices,
|
||||
) -> Result<Outcome, FabroError> {
|
||||
// 1. Build prompt (prepend fidelity preamble if present)
|
||||
|
|
@ -61,6 +63,14 @@ impl Handler for PromptHandler {
|
|||
} else {
|
||||
format!("{preamble}\n\n{expanded}")
|
||||
};
|
||||
let visit = u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX);
|
||||
let stage_dir = stage_dir(run_dir, &node.id, visit);
|
||||
fs::create_dir_all(&stage_dir)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?;
|
||||
fs::write(stage_dir.join("prompt.md"), &prompt)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write prompt file: {err}")))?;
|
||||
|
||||
// 1b. Discover project docs for system prompt when project_memory is enabled
|
||||
let system_prompt = if node.project_memory() {
|
||||
|
|
@ -166,6 +176,28 @@ impl Handler for PromptHandler {
|
|||
extract_status_fields(&response_text, &mut outcome);
|
||||
outcome.usage = stage_usage;
|
||||
outcome.files_touched = backend_files_touched;
|
||||
fs::write(stage_dir.join("response.md"), &response_text)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write response file: {err}")))?;
|
||||
fs::write(
|
||||
stage_dir.join("status.json"),
|
||||
serde_json::to_vec_pretty(&status_json_value(&outcome))
|
||||
.map_err(|err| FabroError::handler(format!("failed to serialize status: {err}")))?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| FabroError::handler(format!("failed to write status file: {err}")))?;
|
||||
write_provider_used_file(
|
||||
services,
|
||||
&stage_dir,
|
||||
&node.id,
|
||||
visit,
|
||||
Some(serde_json::json!({
|
||||
"mode": "prompt",
|
||||
"provider": prompt_provider.unwrap_or_default(),
|
||||
"model": prompt_model.unwrap_or_default(),
|
||||
})),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(outcome)
|
||||
}
|
||||
|
|
@ -205,6 +237,7 @@ mod tests {
|
|||
.await
|
||||
.unwrap();
|
||||
let services = EngineServices {
|
||||
emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)),
|
||||
run_store: run_store.clone(),
|
||||
..EngineServices::test_default()
|
||||
};
|
||||
|
|
|
|||
|
|
@ -1,10 +1,12 @@
|
|||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use async_trait::async_trait;
|
||||
use fabro_store::SlateRunStore;
|
||||
use fabro_types::RunId;
|
||||
use tokio::fs;
|
||||
|
||||
use fabro_core::error::CoreError;
|
||||
use fabro_core::error::Result as CoreResult;
|
||||
use fabro_core::graph::NodeSpec;
|
||||
use fabro_core::lifecycle::RunLifecycle;
|
||||
|
|
@ -12,10 +14,11 @@ use fabro_core::outcome::NodeResult;
|
|||
use fabro_core::state::ExecutionState;
|
||||
|
||||
use super::circuit_breaker::CircuitBreakerLifecycle;
|
||||
use super::git::GitCheckpointResult;
|
||||
use crate::event::{EventEmitter, RunNoticeLevel, WorkflowRunEvent, append_workflow_event};
|
||||
use crate::graph::WorkflowGraph;
|
||||
use crate::graph::WorkflowNode;
|
||||
use crate::outcome::StageUsage;
|
||||
use crate::outcome::{OutcomeExt, StageUsage};
|
||||
use crate::run_options::RunOptions;
|
||||
use fabro_graphviz::graph::types::Graph as GvGraph;
|
||||
|
||||
|
|
@ -30,10 +33,38 @@ pub(crate) struct DiskLifecycle {
|
|||
pub graph: Arc<GvGraph>,
|
||||
pub run_options: Arc<RunOptions>,
|
||||
pub emitter: Arc<EventEmitter>,
|
||||
pub checkpoint_git_result: Arc<Mutex<Option<GitCheckpointResult>>>,
|
||||
pub circuit_breaker: Arc<CircuitBreakerLifecycle>,
|
||||
pub checkpoint_enabled: bool,
|
||||
}
|
||||
|
||||
pub(super) fn build_checkpoint(
|
||||
node: &WorkflowNode,
|
||||
result: &WfNodeResult,
|
||||
next_node_id: Option<&str>,
|
||||
state: &WfRunState,
|
||||
loop_failure_signatures: std::collections::HashMap<fabro_types::FailureSignature, usize>,
|
||||
restart_failure_signatures: std::collections::HashMap<fabro_types::FailureSignature, usize>,
|
||||
git_commit_sha: Option<String>,
|
||||
) -> fabro_types::Checkpoint {
|
||||
let mut node_outcomes = state.node_outcomes.clone();
|
||||
node_outcomes.insert(node.id().to_string(), result.outcome.clone());
|
||||
|
||||
fabro_types::Checkpoint {
|
||||
timestamp: chrono::Utc::now(),
|
||||
current_node: node.id().to_string(),
|
||||
completed_nodes: state.completed_nodes.clone(),
|
||||
node_outcomes,
|
||||
node_retries: state.node_retries.clone(),
|
||||
context_values: state.context.snapshot(),
|
||||
next_node_id: next_node_id.map(String::from),
|
||||
git_commit_sha,
|
||||
node_visits: state.node_visits.clone(),
|
||||
loop_failure_signatures,
|
||||
restart_failure_signatures,
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl RunLifecycle<WorkflowGraph> for DiskLifecycle {
|
||||
async fn on_run_start(&self, _graph: &WorkflowGraph, _state: &WfRunState) -> CoreResult<()> {
|
||||
|
|
@ -55,10 +86,45 @@ impl RunLifecycle<WorkflowGraph> for DiskLifecycle {
|
|||
|
||||
async fn after_node(
|
||||
&self,
|
||||
_node: &WorkflowNode,
|
||||
_result: &mut WfNodeResult,
|
||||
_state: &WfRunState,
|
||||
node: &WorkflowNode,
|
||||
result: &mut WfNodeResult,
|
||||
state: &WfRunState,
|
||||
) -> CoreResult<()> {
|
||||
let visit =
|
||||
u32::try_from(*state.node_visits.get(node.id()).unwrap_or(&1)).unwrap_or(u32::MAX);
|
||||
let node_dir = if visit <= 1 {
|
||||
self.run_dir.join("nodes").join(node.id())
|
||||
} else {
|
||||
self.run_dir
|
||||
.join("nodes")
|
||||
.join(format!("{}-visit_{visit}", node.id()))
|
||||
};
|
||||
fs::create_dir_all(&node_dir)
|
||||
.await
|
||||
.map_err(|err| CoreError::Other(format!("failed to create node dir: {err}")))?;
|
||||
let mut status = serde_json::Map::new();
|
||||
status.insert(
|
||||
"status".to_string(),
|
||||
serde_json::Value::String(result.outcome.status.to_string()),
|
||||
);
|
||||
status.insert(
|
||||
"outcome".to_string(),
|
||||
serde_json::Value::String(result.outcome.status.to_string()),
|
||||
);
|
||||
if let Some(failure_reason) = result.outcome.failure_reason() {
|
||||
status.insert(
|
||||
"failure_reason".to_string(),
|
||||
serde_json::Value::String(failure_reason.to_string()),
|
||||
);
|
||||
}
|
||||
fs::write(
|
||||
node_dir.join("status.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::Value::Object(status)).map_err(|err| {
|
||||
CoreError::Other(format!("failed to serialize status.json: {err}"))
|
||||
})?,
|
||||
)
|
||||
.await
|
||||
.map_err(|err| CoreError::Other(format!("failed to write status.json: {err}")))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
|
@ -73,25 +139,27 @@ impl RunLifecycle<WorkflowGraph> for DiskLifecycle {
|
|||
return Ok(());
|
||||
}
|
||||
|
||||
let git_commit_sha = self
|
||||
.checkpoint_git_result
|
||||
.lock()
|
||||
.unwrap()
|
||||
.as_ref()
|
||||
.and_then(|result| result.commit_sha.clone());
|
||||
let (loop_sigs, restart_sigs) = self.circuit_breaker.snapshot();
|
||||
|
||||
// Build checkpoint from state
|
||||
let mut node_outcomes = state.node_outcomes.clone();
|
||||
node_outcomes.insert(node.id().to_string(), result.outcome.clone());
|
||||
|
||||
let _checkpoint = fabro_types::Checkpoint {
|
||||
timestamp: chrono::Utc::now(),
|
||||
current_node: node.id().to_string(),
|
||||
completed_nodes: state.completed_nodes.clone(),
|
||||
node_outcomes,
|
||||
node_retries: state.node_retries.clone(),
|
||||
context_values: state.context.snapshot(),
|
||||
next_node_id: next_node_id.map(String::from),
|
||||
git_commit_sha: None,
|
||||
node_visits: state.node_visits.clone(),
|
||||
loop_failure_signatures: loop_sigs,
|
||||
restart_failure_signatures: restart_sigs,
|
||||
};
|
||||
let checkpoint = build_checkpoint(
|
||||
node,
|
||||
result,
|
||||
next_node_id,
|
||||
state,
|
||||
loop_sigs,
|
||||
restart_sigs,
|
||||
git_commit_sha,
|
||||
);
|
||||
let checkpoint_bytes = serde_json::to_vec_pretty(&checkpoint)
|
||||
.map_err(|err| CoreError::Other(format!("failed to serialize checkpoint: {err}")))?;
|
||||
fs::write(self.run_dir.join("checkpoint.json"), checkpoint_bytes)
|
||||
.await
|
||||
.map_err(|err| CoreError::Other(format!("failed to write checkpoint.json: {err}")))?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,9 +1,11 @@
|
|||
use std::collections::HashMap;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use async_trait::async_trait;
|
||||
use fabro_store::SlateRunStore;
|
||||
use fabro_types::RunId;
|
||||
use tokio::fs;
|
||||
|
||||
use fabro_core::error::{CoreError, Result as CoreResult};
|
||||
use fabro_core::graph::NodeSpec;
|
||||
|
|
@ -16,6 +18,7 @@ use crate::event::{EventEmitter, RunNoticeLevel, WorkflowRunEvent};
|
|||
use crate::git::MetadataStore;
|
||||
use crate::graph::WorkflowGraph;
|
||||
use crate::graph::WorkflowNode;
|
||||
use crate::lifecycle::disk::build_checkpoint;
|
||||
use crate::outcome::{Outcome, StageStatus, StageUsage};
|
||||
use crate::run_dump::RunDump;
|
||||
use crate::run_options::RunOptions;
|
||||
|
|
@ -92,7 +95,7 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
|
|||
&self,
|
||||
node: &WorkflowNode,
|
||||
result: &WfNodeResult,
|
||||
_next_node_id: Option<&str>,
|
||||
next_node_id: Option<&str>,
|
||||
state: &WfRunState,
|
||||
) -> CoreResult<()> {
|
||||
let node_id = node.id();
|
||||
|
|
@ -113,15 +116,16 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
|
|||
) {
|
||||
let git_author = self.run_options.git_author();
|
||||
let store = MetadataStore::new(repo_path, &git_author);
|
||||
// Build checkpoint JSON for shadow branch
|
||||
if let Some(cp_json) = self
|
||||
.run_store
|
||||
.state()
|
||||
.await
|
||||
.ok()
|
||||
.and_then(|state| state.checkpoint)
|
||||
.and_then(|checkpoint| serde_json::to_vec_pretty(&checkpoint).ok())
|
||||
{
|
||||
let checkpoint = build_checkpoint(
|
||||
node,
|
||||
result,
|
||||
next_node_id,
|
||||
state,
|
||||
HashMap::new(),
|
||||
HashMap::new(),
|
||||
None,
|
||||
);
|
||||
if let Ok(cp_json) = serde_json::to_vec_pretty(&checkpoint) {
|
||||
let mut extra_entries: Vec<(String, Vec<u8>)> = {
|
||||
let artifact_store = self.artifact_store.lock().unwrap();
|
||||
artifact_store
|
||||
|
|
@ -296,6 +300,13 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
|
|||
match git_diff(&*self.sandbox, &base_sha).await {
|
||||
Ok(patch) if !patch.is_empty() => {
|
||||
*self.final_patch.lock().unwrap() = Some(patch.clone());
|
||||
if let Err(err) = fs::write(self.run_dir.join("final.patch"), patch).await {
|
||||
self.emitter.emit(&WorkflowRunEvent::RunNotice {
|
||||
level: RunNoticeLevel::Warn,
|
||||
code: "final_patch_write_failed".to_string(),
|
||||
message: format!("failed to write final.patch: {err}"),
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(_) => {
|
||||
*self.final_patch.lock().unwrap() = None;
|
||||
|
|
|
|||
|
|
@ -150,6 +150,7 @@ impl WorkflowLifecycle {
|
|||
graph: Arc::clone(&graph),
|
||||
run_options: Arc::clone(run_options),
|
||||
emitter: Arc::clone(emitter),
|
||||
checkpoint_git_result: Arc::clone(&checkpoint_git_result),
|
||||
circuit_breaker: Arc::clone(&circuit_breaker),
|
||||
checkpoint_enabled: true,
|
||||
};
|
||||
|
|
@ -407,10 +408,10 @@ impl RunLifecycle<WorkflowGraph> for WorkflowLifecycle {
|
|||
next_node_id: Option<&str>,
|
||||
state: &WfRunState,
|
||||
) -> CoreResult<()> {
|
||||
self.disk
|
||||
self.git
|
||||
.on_checkpoint(node, result, next_node_id, state)
|
||||
.await?;
|
||||
self.git
|
||||
self.disk
|
||||
.on_checkpoint(node, result, next_node_id, state)
|
||||
.await?;
|
||||
self.event
|
||||
|
|
|
|||
|
|
@ -33,9 +33,10 @@ pub(crate) async fn load_from_store(
|
|||
.state()
|
||||
.await
|
||||
.map_err(|err| FabroError::engine(err.to_string()))?;
|
||||
let run_record = state
|
||||
let mut run_record = state
|
||||
.run
|
||||
.ok_or_else(|| FabroError::Precondition("run record missing from store".to_string()))?;
|
||||
run_record.created_at = run_store.created_at();
|
||||
let graph = run_record.graph.clone();
|
||||
let source = state.graph_source.unwrap_or_default();
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue