diff --git a/lib/crates/fabro-sandbox/src/local.rs b/lib/crates/fabro-sandbox/src/local.rs index 8b69b7377..27c89da0c 100644 --- a/lib/crates/fabro-sandbox/src/local.rs +++ b/lib/crates/fabro-sandbox/src/local.rs @@ -282,14 +282,18 @@ impl Sandbox for LocalSandbox { let stdout_task = tokio::spawn(async move { let mut buf = String::new(); if let Some(ref mut r) = stdout_pipe { - let _ = r.read_to_string(&mut buf).await; + if let Err(err) = r.read_to_string(&mut buf).await { + tracing::warn!(error = %err, stream = "stdout", "Failed to drain child stdout"); + } } buf }); let stderr_task = tokio::spawn(async move { let mut buf = String::new(); if let Some(ref mut r) = stderr_pipe { - let _ = r.read_to_string(&mut buf).await; + if let Err(err) = r.read_to_string(&mut buf).await { + tracing::warn!(error = %err, stream = "stderr", "Failed to drain child stderr"); + } } buf }); diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index 676e2f142..7e3a4be0c 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -52,6 +52,8 @@ pub trait CodergenBackend: Send + Sync { _node: &Node, _prompt: &str, _system_prompt: Option<&str>, + _emitter: &Arc, + _stage_scope: &StageScope, ) -> Result { Err(Error::Validation( "one_shot mode not supported by this backend".into(), diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 5d21c7d32..9247c19ec 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -285,6 +285,8 @@ impl CodergenBackend for AgentApiBackend { node: &Node, prompt: &str, system_prompt: Option<&str>, + emitter: &Arc, + stage_scope: &StageScope, ) -> Result { let client = Client::from_source(self.source.as_ref()) .await @@ -358,14 +360,16 @@ impl CodergenBackend for AgentApiBackend { let mut found = None; for target in fallback_chain { - tracing::warn!( - stage = node.id.as_str(), - from_provider = from_provider.as_str(), - from_model = from_model.as_str(), - to_provider = target.provider.as_str(), - to_model = target.model.as_str(), - error = error_msg.as_str(), - "LLM provider failover (prompt)" + emitter.emit_scoped( + &Event::Failover { + stage: node.id.clone(), + from_provider: from_provider.clone(), + from_model: from_model.clone(), + to_provider: target.provider.clone(), + to_model: target.model.clone(), + error: error_msg.clone(), + }, + stage_scope, ); let max_tokens = node.max_tokens().or_else(|| { diff --git a/lib/crates/fabro-workflow/src/handler/llm/cli.rs b/lib/crates/fabro-workflow/src/handler/llm/cli.rs index 5036eec97..7fcf60ed8 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/cli.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/cli.rs @@ -810,9 +810,13 @@ impl CodergenBackend for BackendRouter { node: &Node, prompt: &str, system_prompt: Option<&str>, + emitter: &Arc, + stage_scope: &StageScope, ) -> Result { // CLI backend doesn't support one_shot, always route to API - self.api_backend.one_shot(node, prompt, system_prompt).await + self.api_backend + .one_shot(node, prompt, system_prompt, emitter, stage_scope) + .await } } diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 39862586c..42bddbbff 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -12,7 +12,7 @@ use tokio::sync::Semaphore; use super::{EngineServices, Handler}; use crate::context::{Context, WorkflowContext, keys}; use crate::error::Error; -use crate::event::{Event, StageScope}; +use crate::event::{Event, RunNoticeLevel, StageScope}; use crate::git::sanitize_ref_component; use crate::hook_context::set_hook_node; use crate::millis_u64; @@ -207,6 +207,11 @@ impl Handler for ParallelHandler { error = %fabro_sandbox::display_for_log(&e), "parallel base checkpoint failed" ); + services.run.emitter.notice( + RunNoticeLevel::Warn, + "parallel_base_checkpoint_failed", + format!("Could not checkpoint base state before parallel branches: {e}"), + ); None } } @@ -875,4 +880,107 @@ mod tests { let branch_count = context.get(keys::PARALLEL_BRANCH_COUNT); assert_eq!(branch_count, Some(serde_json::json!(2))); } + + #[expect( + clippy::disallowed_methods, + reason = "Test sets up a temporary git repo to exercise the parallel checkpoint failure path." + )] + #[tokio::test] + async fn parallel_base_checkpoint_failure_emits_notice() { + // Set up a git repo then lock the index to make `git add` fail. + let repo_dir = tempfile::tempdir().unwrap(); + let init = std::process::Command::new("git") + .args(["init", "-b", "main"]) + .current_dir(repo_dir.path()) + .output() + .unwrap(); + assert!(init.status.success()); + for (key, value) in [("user.name", "Test"), ("user.email", "test@test.com")] { + let config = std::process::Command::new("git") + .args(["config", key, value]) + .current_dir(repo_dir.path()) + .output() + .unwrap(); + assert!(config.status.success()); + } + let commit = std::process::Command::new("git") + .args(["commit", "--allow-empty", "-m", "initial"]) + .current_dir(repo_dir.path()) + .output() + .unwrap(); + assert!(commit.status.success()); + // Lock the index to make `git add` fail + let lock_path = repo_dir.path().join(".git/index.lock"); + std::fs::write(&lock_path, b"locked").unwrap(); + + let store = test_store(); + let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); + let mut services = EngineServices::test_default(); + let emitter = Arc::new(crate::event::Emitter::new(fixtures::RUN_1)); + let seen = Arc::new(std::sync::Mutex::new(Vec::new())); + emitter.on_event({ + let seen = Arc::clone(&seen); + move |event| seen.lock().unwrap().push(event.clone()) + }); + services.run = services + .run + .with_emitter(emitter) + .with_run_store(run_store.clone().into()) + .with_sandbox(Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + ))); + // Supply git_state so the parallel handler enters the git path + services.set_git_state(Some(Arc::new(crate::sandbox_git::GitState { + run_id: fixtures::RUN_1, + base_sha: "abc123".to_string(), + run_branch: Some("fabro/run/test".to_string()), + meta_branch: None, + checkpoint_exclude_globs: Vec::new(), + git_author: crate::git::GitAuthor::default(), + }))); + let logger = crate::event::StoreProgressLogger::new(run_store.clone()); + logger.register(services.run.emitter.as_ref()); + + let mut node = Node::new("par"); + node.attrs.insert( + "shape".to_string(), + AttrValue::String("component".to_string()), + ); + let context = test_context(); + let mut graph = Graph::new("test"); + graph.nodes.insert("par".to_string(), node.clone()); + graph + .nodes + .insert("branch_a".to_string(), Node::new("branch_a")); + graph.edges.push(Edge::new("par", "branch_a")); + + let outcome = ParallelHandler + .execute(&node, &context, &graph, repo_dir.path(), &services) + .await + .unwrap(); + logger.flush().await; + + // Clean up the lock so tempdir cleanup doesn't leave artifacts + let _ = std::fs::remove_file(&lock_path); + + // Should still succeed (the checkpoint failure is non-fatal) + assert_eq!(outcome.status, StageOutcome::Succeeded); + + let events = seen.lock().unwrap(); + let notice = events.iter().find(|e| { + matches!( + &e.body, + fabro_types::EventBody::RunNotice(props) + if props.code == "parallel_base_checkpoint_failed" + ) + }); + assert!( + notice.is_some(), + "expected parallel_base_checkpoint_failed notice, got events: {:?}", + events + .iter() + .map(fabro_types::RunEvent::event_name) + .collect::>() + ); + } } diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index a76099b80..502b9666e 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -105,7 +105,13 @@ impl Handler for PromptHandler { let (response_text, stage_usage, backend_files_touched) = if let Some(backend) = &self.backend { let result = backend - .one_shot(node, &prompt, system_prompt.as_deref()) + .one_shot( + node, + &prompt, + system_prompt.as_deref(), + &services.run.emitter, + &stage_scope, + ) .await; match result { Ok(CodergenResult::Full(outcome)) => return Ok(outcome), @@ -279,6 +285,8 @@ mod tests { _node: &Node, _prompt: &str, _system_prompt: Option<&str>, + _emitter: &Arc, + _stage_scope: &StageScope, ) -> Result { Ok(CodergenResult::Text { text: "one-shot response".to_string(), @@ -339,6 +347,8 @@ mod tests { _node: &Node, _prompt: &str, _system_prompt: Option<&str>, + _emitter: &Arc, + _stage_scope: &StageScope, ) -> Result { Ok(CodergenResult::Text { text: "one-shot response".to_string(), @@ -396,6 +406,8 @@ mod tests { _node: &Node, prompt: &str, system_prompt: Option<&str>, + _emitter: &Arc, + _stage_scope: &StageScope, ) -> Result { *self.captured_prompt.lock().unwrap() = Some(prompt.to_string()); *self.captured_system_prompt.lock().unwrap() = Some(system_prompt.map(String::from)); diff --git a/lib/crates/fabro-workflow/src/lifecycle/git.rs b/lib/crates/fabro-workflow/src/lifecycle/git.rs index 0e0e9f195..1e02808d0 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/git.rs @@ -292,6 +292,14 @@ impl RunLifecycle for GitLifecycle { error = %fabro_sandbox::display_for_log(&err), "git push from run lifecycle failed" ); + self.emitter.emit(&Event::RunNotice { + level: RunNoticeLevel::Warn, + code: "git_push_failed".to_string(), + message: format!( + "Failed to push run branch {branch}: {err}" + ), + exec_output_tail: exec_output_tail.clone(), + }); (false, exec_output_tail) } }; @@ -1062,4 +1070,76 @@ mod tests { Ok(None) } } + + #[expect( + clippy::disallowed_methods, + reason = "Test sets up a temporary git repo with a bogus remote to verify push-failure notice." + )] + #[tokio::test] + async fn checkpoint_push_failure_emits_git_push_failed_notice() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + // Add a bogus remote so git_push_ref actually attempts the push. + let add_remote = std::process::Command::new("git") + .args(["remote", "add", "origin", "https://127.0.0.1:1/nope.git"]) + .current_dir(repo_dir.path()) + .output() + .unwrap(); + assert!(add_remote.status.success()); + let branch = "fabro/metadata/run"; + let run_branch = "fabro/run/test"; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let events = record_events(&emitter); + let run_store = run_store(fixtures::RUN_1).await; + let handle = RunStoreHandle::local(run_store.clone()); + // run_options with a run_branch so the push path is entered + let options = Arc::new(RunOptions { + settings: WorkflowSettings::default(), + run_dir: repo_dir.path().to_path_buf(), + cancel_token: None, + run_id: fixtures::RUN_1, + labels: HashMap::new(), + workflow_slug: Some("metadata".to_string()), + github_app: None, + pre_run_git: None, + fork_source_ref: None, + base_branch: None, + display_base_sha: None, + git: Some(GitCheckpointOptions { + base_sha: None, + run_branch: Some(run_branch.to_string()), + meta_branch: Some(branch.to_string()), + }), + }); + let lifecycle = git_lifecycle( + repo_dir.path(), + emitter, + handle, + options, + Arc::new(RunMetadataRuntime::new()), + ); + let graph = workflow_graph(); + let node = graph.get_node("build").unwrap(); + let mut state = ExecutionState::new(&graph).unwrap(); + state.increment_visits("build"); + let result = WfNodeResult::new(Outcome::success(), Duration::from_millis(10), 1, 1); + + lifecycle + .on_checkpoint(&node, &result, Some("exit"), &state) + .await + .unwrap(); + + let events = events.lock().unwrap(); + let notice = events.iter().find(|e| { + matches!( + &e.body, + EventBody::RunNotice(props) if props.code == "git_push_failed" + ) + }); + assert!( + notice.is_some(), + "expected git_push_failed notice, got events: {:?}", + events.iter().map(RunEvent::event_name).collect::>() + ); + } } diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs index ce125827b..6691ec579 100644 --- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs @@ -239,11 +239,14 @@ async fn build_sandbox_env( Ok(token) => { env.insert("GITHUB_TOKEN".to_string(), token); } - Err(e) => emitter.notice( - RunNoticeLevel::Warn, - "github_token_failed", - format!("Failed to mint GitHub token: {e}"), - ), + Err(e) => { + tracing::warn!(error = %e, "Failed to mint GitHub token"); + emitter.notice( + RunNoticeLevel::Warn, + "github_token_failed", + format!("Failed to mint GitHub token: {e}"), + ); + } } } } @@ -504,6 +507,15 @@ pub async fn initialize( } } } else { + tracing::warn!( + worktree_mode = ?options.worktree_mode, + "worktree requested but cwd is not a git repository; running without a worktree" + ); + options.emitter.notice( + RunNoticeLevel::Warn, + "worktree_skipped_no_git", + "Worktree mode requested but no Git repository was found; running without a worktree.", + ); Arc::new(ReadBeforeWriteSandbox::new(inner)) } } else { @@ -619,7 +631,15 @@ pub async fn initialize( options.run_options.base_branch = info.base_branch; } } - Ok(None) => {} + Ok(None) => { + if sandbox.origin_url().is_some() { + options.emitter.notice( + RunNoticeLevel::Warn, + "sandbox_git_unavailable", + "Sandbox could not set up Git despite a configured origin; running without checkpointing or PR support.", + ); + } + } Err(e) => { return Err(Error::engine_with_source("Sandbox git setup failed", &e)); } @@ -1373,4 +1393,82 @@ mod tests { assert!(matches!(result, Err(Error::Cancelled))); } + + #[tokio::test] + async fn initialize_worktree_skipped_no_git_emits_notice() { + // Use a non-git temp dir so resolve_worktree_base_sha returns Ok(None) + let temp = tempfile::tempdir().unwrap(); + let non_git_dir = temp.path().to_path_buf(); + let run_dir = temp.path().join("run"); + std::fs::create_dir_all(&run_dir).unwrap(); + let (graph, source) = simple_graph(); + let persisted = test_persisted(graph, source, &run_dir); + let emitter = Arc::new(crate::event::Emitter::new(test_run_id())); + let seen = Arc::new(std::sync::Mutex::new(Vec::new())); + emitter.on_event({ + let seen = Arc::clone(&seen); + move |event| seen.lock().unwrap().push(event.clone()) + }); + + let _initialized = initialize(persisted, InitOptions { + run_id: test_run_id(), + run_store: { + let store = memory_store(); + let inner = store.create_run(&test_run_id()).await.unwrap(); + inner.into() + }, + dry_run: false, + emitter, + sandbox: SandboxSpec::Local { + working_directory: non_git_dir, + }, + llm: LlmSpec { + model: "test-model".to_string(), + provider: fabro_llm::Provider::Anthropic, + fallback_chain: Vec::new(), + mcp_servers: Vec::new(), + dry_run: true, + }, + interviewer: Arc::new(AutoApproveInterviewer::engine()), + lifecycle: crate::run_options::LifecycleOptions { + setup_commands: vec![], + setup_command_timeout_ms: 1_000, + devcontainer_phases: vec![], + }, + run_options: test_settings(&run_dir), + workflow_path: None, + workflow_bundle: None, + hooks: fabro_hooks::HookSettings { hooks: vec![] }, + sandbox_env: SandboxEnvSpec { + devcontainer_env: HashMap::new(), + toml_env: HashMap::new(), + github_permissions: None, + origin_url: None, + }, + vault: None, + devcontainer: None, + git: None, + worktree_mode: Some(WorktreeMode::Always), + run_control: None, + registry_override: None, + artifact_sink: None, + checkpoint: None, + seed_context: None, + }) + .await + .unwrap(); + + let events = seen.lock().unwrap(); + let notice = events.iter().find(|e| { + matches!( + &e.body, + EventBody::RunNotice(props) if props.code == "worktree_skipped_no_git" + ) + }); + assert!( + notice.is_some(), + "expected worktree_skipped_no_git notice, got events: {:?}", + events.iter().map(RunEvent::event_name).collect::>() + ); + } } diff --git a/lib/crates/fabro-workflow/tests/it/integration.rs b/lib/crates/fabro-workflow/tests/it/integration.rs index 002d473da..cd11db087 100644 --- a/lib/crates/fabro-workflow/tests/it/integration.rs +++ b/lib/crates/fabro-workflow/tests/it/integration.rs @@ -6193,6 +6193,7 @@ mod real_llm { use fabro_types::WorkflowSettings; use fabro_workflow::context::Context; use fabro_workflow::error::Error; + use fabro_workflow::event::StageScope; use fabro_workflow::handler::agent::{AgentHandler, CodergenBackend, CodergenResult}; struct LlmCodergenBackend { @@ -6221,6 +6222,8 @@ mod real_llm { _node: &Node, prompt: &str, _system_prompt: Option<&str>, + _emitter: &Arc, + _stage_scope: &StageScope, ) -> Result { self.complete(prompt).await }