diff --git a/lib/crates/fabro-retro/src/retro_agent.rs b/lib/crates/fabro-retro/src/retro_agent.rs index 97144bb43..90a384de0 100644 --- a/lib/crates/fabro-retro/src/retro_agent.rs +++ b/lib/crates/fabro-retro/src/retro_agent.rs @@ -1,4 +1,4 @@ -use std::path::{Path, PathBuf}; +use std::path::Path; use std::sync::{Arc, Mutex}; use std::time::Duration; @@ -11,8 +11,6 @@ use fabro_llm::client::Client; use fabro_llm::provider::Provider; use fabro_llm::types::ToolDefinition; use fabro_store::SlateRunStore; -use fabro_util::redact::redact_jsonl_line; -use tokio::sync::broadcast::Receiver; use tokio::task::JoinHandle; use crate::retro::{RetroNarrative, SmoothnessRating}; @@ -194,12 +192,6 @@ pub async fn run_retro_agent( None, ); - // Set up event writer before initialize (which emits SessionStarted) - let retro_dir = run_dir.join("retro"); - std::fs::create_dir_all(&retro_dir)?; - let rx = session.subscribe(); - let event_writer_handle = spawn_retro_event_writer(rx, retro_dir.join("retro_session.jsonl")); - // Optionally forward agent events via the callback let event_forwarder_handle = event_callback.map(|cb| spawn_retro_event_forwarder(&session, cb)); @@ -207,8 +199,6 @@ pub async fn run_retro_agent( let prompt = build_retro_prompt(RETRO_DATA_DIR); - write_retro_prompt(run_store, &retro_dir, &prompt)?; - let process_result = session .process_input(&prompt) .await @@ -228,7 +218,7 @@ pub async fn run_retro_agent( .to_string(); // Extract result / determine outcome - let (outcome, failure_reason, narrative_result) = match process_result { + let (_outcome, _failure_reason, narrative_result) = match process_result { Ok(()) => { let maybe_narrative = captured .lock() @@ -249,19 +239,8 @@ pub async fn run_retro_agent( } }; - // Write artifacts (on both success and failure) - write_retro_response(run_store, &retro_dir, &response_text)?; - write_retro_artifacts( - &retro_dir, - provider.as_str(), - model, - outcome, - failure_reason.as_deref(), - ); - - // Drop session to close the broadcast channel, then wait for event writer/forwarder + // Drop session to close the broadcast channel, then wait for event forwarder drop(session); - let _ = event_writer_handle.await; if let Some(handle) = event_forwarder_handle { let _ = handle.await; } @@ -285,72 +264,6 @@ pub fn dry_run_narrative() -> RetroNarrative { } } -fn write_retro_prompt( - _run_store: &SlateRunStore, - retro_dir: &Path, - prompt: &str, -) -> anyhow::Result<()> { - std::fs::write(retro_dir.join("prompt.md"), prompt)?; - Ok(()) -} - -fn write_retro_response( - _run_store: &SlateRunStore, - retro_dir: &Path, - response: &str, -) -> anyhow::Result<()> { - std::fs::write(retro_dir.join("response.md"), response)?; - Ok(()) -} - -/// Write retro artifact files (provider_used.json, status.json) into `retro_dir`. -/// Called on both success and failure paths so artifacts are always available for debugging. -fn write_retro_artifacts( - retro_dir: &Path, - provider: &str, - model: &str, - outcome: &str, - failure_reason: Option<&str>, -) { - let provider_used = serde_json::json!({ - "mode": "agent", - "provider": provider, - "model": model, - }); - if let Ok(json) = serde_json::to_string_pretty(&provider_used) { - let _ = std::fs::write(retro_dir.join("provider_used.json"), json); - } - - let status = serde_json::json!({ - "outcome": outcome, - "failure_reason": failure_reason, - "timestamp": chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true), - }); - if let Ok(json) = serde_json::to_string_pretty(&status) { - let _ = std::fs::write(retro_dir.join("status.json"), json); - } -} - -/// Spawn a background task that reads `SessionEvent`s from the broadcast receiver -/// and appends them as JSONL to the given path. -fn spawn_retro_event_writer(mut rx: Receiver, path: PathBuf) -> JoinHandle<()> { - tokio::spawn(async move { - use std::io::Write; - while let Ok(event) = rx.recv().await { - if let Ok(line) = serde_json::to_string(&event) { - let line = redact_jsonl_line(&line); - if let Ok(mut f) = std::fs::OpenOptions::new() - .create(true) - .append(true) - .open(&path) - { - let _ = writeln!(f, "{line}"); - } - } - } - }) -} - /// Spawn a background task that forwards session events via the provided callback. fn spawn_retro_event_forwarder( session: &Session, @@ -453,9 +366,6 @@ async fn upload_file( #[cfg(test)] mod tests { use super::*; - use fabro_agent::AgentEvent; - use std::time::SystemTime; - use tokio::sync::broadcast; #[test] fn submit_retro_schema_is_valid_json() { @@ -505,106 +415,4 @@ mod tests { assert!(narrative.friction_points.is_empty()); assert!(narrative.open_items.is_empty()); } - - #[test] - fn write_retro_artifacts_does_not_clobber_prompt_md() { - let dir = tempfile::tempdir().unwrap(); - let retro_dir = dir.path().join("retro"); - std::fs::create_dir_all(&retro_dir).unwrap(); - // prompt.md is written separately by run_retro_agent, not by write_retro_artifacts - std::fs::write(retro_dir.join("prompt.md"), "Analyze the run data").unwrap(); - write_retro_artifacts( - &retro_dir, - "anthropic", - "claude-sonnet-4-20250514", - "success", - None, - ); - let content = std::fs::read_to_string(retro_dir.join("prompt.md")).unwrap(); - assert_eq!(content, "Analyze the run data"); - } - - #[test] - fn writes_provider_used_json() { - let dir = tempfile::tempdir().unwrap(); - let retro_dir = dir.path().join("retro"); - std::fs::create_dir_all(&retro_dir).unwrap(); - write_retro_artifacts(&retro_dir, "openai", "gpt-4o", "success", None); - let content = std::fs::read_to_string(retro_dir.join("provider_used.json")).unwrap(); - let parsed: serde_json::Value = serde_json::from_str(&content).unwrap(); - assert_eq!(parsed["mode"], "agent"); - assert_eq!(parsed["provider"], "openai"); - assert_eq!(parsed["model"], "gpt-4o"); - } - - #[test] - fn writes_status_json_success() { - let dir = tempfile::tempdir().unwrap(); - let retro_dir = dir.path().join("retro"); - std::fs::create_dir_all(&retro_dir).unwrap(); - write_retro_artifacts( - &retro_dir, - "anthropic", - "claude-sonnet-4-20250514", - "success", - None, - ); - let content = std::fs::read_to_string(retro_dir.join("status.json")).unwrap(); - let parsed: serde_json::Value = serde_json::from_str(&content).unwrap(); - assert_eq!(parsed["outcome"], "success"); - assert!(parsed["failure_reason"].is_null()); - assert!(parsed["timestamp"].as_str().unwrap().contains('T')); - } - - #[test] - fn writes_status_json_failure() { - let dir = tempfile::tempdir().unwrap(); - let retro_dir = dir.path().join("retro"); - std::fs::create_dir_all(&retro_dir).unwrap(); - write_retro_artifacts( - &retro_dir, - "anthropic", - "claude-sonnet-4-20250514", - "error", - Some("Retro agent did not call submit_retro"), - ); - let content = std::fs::read_to_string(retro_dir.join("status.json")).unwrap(); - let parsed: serde_json::Value = serde_json::from_str(&content).unwrap(); - assert_eq!(parsed["outcome"], "error"); - assert_eq!( - parsed["failure_reason"], - "Retro agent did not call submit_retro" - ); - } - - #[tokio::test] - async fn event_writer_writes_session_events_to_jsonl() { - let dir = tempfile::tempdir().unwrap(); - let jsonl_path = dir.path().join("retro_session.jsonl"); - - let (tx, rx) = broadcast::channel::(16); - let handle = spawn_retro_event_writer(rx, jsonl_path.clone()); - - tx.send(SessionEvent { - event: AgentEvent::SessionStarted { - provider: Some("anthropic".into()), - model: Some("claude-opus".into()), - }, - timestamp: SystemTime::now(), - session_id: "retro-test".into(), - parent_session_id: None, - }) - .unwrap(); - - // Drop sender so the receiver loop ends - drop(tx); - handle.await.unwrap(); - - let content = std::fs::read_to_string(&jsonl_path).unwrap(); - let lines: Vec<&str> = content.lines().collect(); - assert_eq!(lines.len(), 1); - let parsed: serde_json::Value = serde_json::from_str(lines[0]).unwrap(); - assert_eq!(parsed["session_id"], "retro-test"); - assert!(lines[0].contains("SessionStarted")); - } } diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index be6cbac38..62d271889 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -1455,14 +1455,10 @@ pub fn build_redacted_event_payload( pub fn append_progress_event(run_dir: &Path, envelope: &RunEventEnvelope) -> Result<()> { let line = redacted_event_json(envelope)?; - append_progress_event_with_line(run_dir, envelope, &line) + append_progress_event_with_line(run_dir, &line) } -pub fn append_progress_event_with_line( - run_dir: &Path, - envelope: &RunEventEnvelope, - line: &str, -) -> Result<()> { +pub fn append_progress_event_with_line(run_dir: &Path, line: &str) -> Result<()> { let mut file = std::fs::OpenOptions::new() .create(true) .append(true) @@ -1475,11 +1471,6 @@ pub fn append_progress_event_with_line( })?; writeln!(file, "{line}")?; - let pretty = serde_json::to_string_pretty(&normalized_envelope_value(envelope)?)?; - let pretty = redact_jsonl_line(&pretty); - std::fs::write(run_dir.join("live.json"), pretty) - .with_context(|| format!("Failed to write {}", run_dir.join("live.json").display()))?; - Ok(()) } diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index ac511d751..2de148dca 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -252,11 +252,9 @@ impl Handler for AgentHandler { format!("{preamble}\n\n{expanded}") }; - // 2. Write prompt to logs let visit = visit_from_context(context); let stage_dir = node_dir(run_dir, &node.id, visit); fs::create_dir_all(&stage_dir).await?; - fs::write(stage_dir.join("prompt.md"), &prompt).await?; // 3. Call LLM backend (agent loop) let thread_id = context.thread_id(); @@ -313,9 +311,7 @@ impl Handler for AgentHandler { ) }; - // 4. Write response to logs - fs::write(stage_dir.join("response.md"), &response_text).await?; - // 7. Build and write status + // Build and write status let mut outcome = Outcome::success(); outcome.notes = Some(format!("Stage completed: {}", node.id)); outcome diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs index 832b6e113..6c8451fc2 100644 --- a/lib/crates/fabro-workflow/src/handler/command.rs +++ b/lib/crates/fabro-workflow/src/handler/command.rs @@ -5,10 +5,8 @@ use crate::context::keys; use crate::error::FabroError; use crate::event::WorkflowRunEvent; use crate::outcome::{Outcome, OutcomeExt}; -use crate::run_dir::{node_dir, visit_from_context}; use async_trait::async_trait; use fabro_graphviz::graph::{Graph, Node}; -use tokio::fs; use super::{EngineServices, Handler}; @@ -59,9 +57,9 @@ impl Handler for CommandHandler { async fn execute( &self, node: &Node, - context: &Context, + _context: &Context, _graph: &Graph, - run_dir: &Path, + _run_dir: &Path, services: &EngineServices, ) -> Result { let script = node @@ -87,20 +85,6 @@ impl Handler for CommandHandler { ))); } - let visit = visit_from_context(context); - let stage_dir = node_dir(run_dir, &node.id, visit); - fs::create_dir_all(&stage_dir).await?; - let invocation = serde_json::json!({ - "command": script, - "language": language, - "timeout_ms": timeout_ms(node), - }); - fs::write( - stage_dir.join("script_invocation.json"), - serde_json::to_string_pretty(&invocation).unwrap(), - ) - .await?; - services.emitter.emit(&WorkflowRunEvent::CommandStarted { node_id: node.id.clone(), script: script.to_string(), @@ -129,20 +113,6 @@ impl Handler for CommandHandler { .await .map_err(|e| FabroError::handler(format!("Failed to spawn script: {e}")))?; - fs::write(stage_dir.join("stdout.log"), &result.stdout).await?; - fs::write(stage_dir.join("stderr.log"), &result.stderr).await?; - - let timing = serde_json::json!({ - "duration_ms": result.duration_ms, - "exit_code": if result.timed_out { serde_json::Value::Null } else { serde_json::json!(result.exit_code) }, - "timed_out": result.timed_out, - }); - fs::write( - stage_dir.join("script_timing.json"), - serde_json::to_string_pretty(&timing).unwrap(), - ) - .await?; - services.emitter.emit(&WorkflowRunEvent::CommandCompleted { node_id: node.id.clone(), stdout: result.stdout.clone(), diff --git a/lib/crates/fabro-workflow/src/handler/fan_in.rs b/lib/crates/fabro-workflow/src/handler/fan_in.rs index 94a898633..43dd3ca91 100644 --- a/lib/crates/fabro-workflow/src/handler/fan_in.rs +++ b/lib/crates/fabro-workflow/src/handler/fan_in.rs @@ -236,7 +236,6 @@ async fn llm_evaluate( let visit_u32 = u32::try_from(visit).unwrap_or(u32::MAX); let stage_dir = node_dir(run_dir, node_id, visit); fs::create_dir_all(&stage_dir).await?; - fs::write(stage_dir.join("prompt.md"), &full_prompt).await?; emitter.emit(&WorkflowRunEvent::Prompt { stage: node_id.to_string(), @@ -282,7 +281,6 @@ async fn llm_evaluate( provider: String::new(), usage: None, }); - fs::write(stage_dir.join("response.md"), &response_text).await?; Ok(Candidate { id: best_id, status: outcome.status.to_string(), @@ -297,7 +295,6 @@ async fn llm_evaluate( provider: String::new(), usage: None, }); - fs::write(stage_dir.join("response.md"), &text).await?; // The LLM responded with text; try to find a matching candidate ID let text = text.trim().to_string(); diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 16f41e63d..8156720d6 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -327,7 +327,7 @@ impl CodergenBackend for AgentApiBackend { let default_provider = self.provider.as_str().to_string(); - let (response, actual_model, actual_provider) = match result { + let (response, actual_model, _actual_provider) = match result { Ok(resp) => ( resp, request.model.clone(), @@ -395,15 +395,6 @@ impl CodergenBackend for AgentApiBackend { let _ = fs::write(stage_dir.join("api_response.json"), json).await; } - let provider_used = serde_json::json!({ - "mode": "prompt", - "provider": &actual_provider, - "model": &actual_model, - }); - if let Ok(json) = serde_json::to_string_pretty(&provider_used) { - let _ = fs::write(stage_dir.join("provider_used.json"), json).await; - } - let mut stage_usage = StageUsage { model: actual_model, input_tokens: response.usage.input_tokens, diff --git a/lib/crates/fabro-workflow/src/handler/llm/cli.rs b/lib/crates/fabro-workflow/src/handler/llm/cli.rs index a72ddce4b..00408f9ba 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/cli.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/cli.rs @@ -509,15 +509,6 @@ impl CodergenBackend for AgentCliBackend { }); let _ = fs::create_dir_all(stage_dir).await; - let provider_used = serde_json::json!({ - "mode": "cli", - "provider": provider.as_str(), - "model": model, - "command": &command, - }); - if let Ok(json) = serde_json::to_string_pretty(&provider_used) { - let _ = fs::write(stage_dir.join("provider_used.json"), json).await; - } // Forward provider API key and custom env vars so the CLI tool can authenticate. // Build a HashMap to pass via exec_command's env_vars parameter — this diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 4c4872584..31004d2ef 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -15,11 +15,10 @@ use crate::git::sanitize_ref_component; use crate::hook_context::set_hook_node; use crate::millis_u64; use crate::outcome::{FailureCategory, FailureDetail, Outcome, OutcomeExt, StageStatus}; -use crate::run_dir::{node_dir, visit_from_context}; +use crate::run_dir::visit_from_context; use crate::sandbox_git::{GIT_REMOTE, git_checkpoint, git_merge_ff_only, git_remove_worktree}; use fabro_graphviz::graph::{AttrValue, Graph, Node}; use fabro_hooks::{HookContext, HookEvent}; -use tokio::fs; use super::{EngineServices, Handler}; @@ -477,13 +476,6 @@ impl Handler for ParallelHandler { context.set(keys::PARALLEL_RESULTS, serde_json::json!(results_json)); context.set(keys::PARALLEL_BRANCH_COUNT, serde_json::json!(total)); - let visit = visit_from_context(context); - let node_dir = node_dir(run_dir, &node.id, visit); - let _ = fs::create_dir_all(&node_dir).await; - if let Ok(json) = serde_json::to_string_pretty(&results_json) { - let _ = fs::write(node_dir.join("parallel_results.json"), json).await; - } - services.emitter.emit(&WorkflowRunEvent::ParallelCompleted { duration_ms: millis_u64(parallel_start.elapsed()), success_count, diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index 6c5de318f..7607156c5 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -87,11 +87,9 @@ impl Handler for PromptHandler { None }; - // 2. Write prompt to logs let visit = visit_from_context(context); let stage_dir = node_dir(run_dir, &node.id, visit); fs::create_dir_all(&stage_dir).await?; - fs::write(stage_dir.join("prompt.md"), &prompt).await?; let prompt_provider = node .provider() @@ -155,10 +153,7 @@ impl Handler for PromptHandler { usage: stage_usage.clone(), }); - // 4. Write response to logs - fs::write(stage_dir.join("response.md"), &response_text).await?; - - // 5. Build and write status + // 4. Build and write status let mut outcome = Outcome::success(); outcome.notes = Some(format!("Stage completed: {}", node.id)); outcome diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index 830d48471..9cedb99b8 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -4,7 +4,6 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -use chrono::Utc; use fabro_config::sandbox::WorktreeMode; use fabro_config::{project as project_config, run as run_config, sandbox as sandbox_config}; use fabro_interview::{AutoApproveInterviewer, Interviewer}; @@ -12,7 +11,6 @@ use fabro_model::{Catalog, FallbackTarget, Provider}; use fabro_sandbox::{SandboxProvider, SandboxSpec}; use fabro_store::{RunStoreHandle, SlateRunStore}; use fabro_types::{RunId, Settings}; -use serde::Serialize; use crate::context::Context; use crate::error::FabroError; @@ -775,32 +773,12 @@ impl Drop for DetachedRunCompletionGuard { async fn persist_detached_failure( run_id: RunId, run_store: &SlateRunStore, - run_dir: &Path, + _run_dir: &Path, phase: &'static str, reason: StatusReason, error: &FabroError, ) -> Result<(), FabroError> { - #[derive(Serialize)] - struct DetachedFailureRecord<'a> { - timestamp: chrono::DateTime, - phase: &'a str, - reason: StatusReason, - error: String, - } - let message = error.to_string(); - let record = DetachedFailureRecord { - timestamp: Utc::now(), - phase, - reason, - error: message.clone(), - }; - - std::fs::write( - run_dir.join("detached_failure.json"), - serde_json::to_string_pretty(&record).map_err(|err| FabroError::Io(err.to_string()))?, - ) - .map_err(|err| FabroError::Io(err.to_string()))?; if let Err(err) = append_workflow_event( run_store,