Remove redundant run_dir file writes that duplicate event-sourced data

All data in these files is already stored in SlateDB via events and
projected into RunState. No production code reads them from disk.

Removed writes: prompt.md, response.md, stdout.log, stderr.log,
script_invocation.json, script_timing.json, parallel_results.json,
provider_used.json, retro/{prompt,response,status,session}, live.json,
detached_failure.json.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-04-03 11:04:55 -07:00
parent 7ba75e4743
commit 1cc389a485
No known key found for this signature in database
10 changed files with 12 additions and 303 deletions

View file

@ -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<SessionEvent>, 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::<SessionEvent>(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"));
}
}

View file

@ -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(())
}

View file

@ -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

View file

@ -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<Outcome, FabroError> {
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(),

View file

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

View file

@ -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,

View file

@ -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

View file

@ -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,

View file

@ -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

View file

@ -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<Utc>,
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,