refactor(run): remove file-backed attach fallback

Require attach to use SlateDB-backed run state and events instead of
falling back to progress.jsonl, status.json, and conclusion.json.
Also remove the unused disk progress logger and keep PR body plan text
store-backed so the remaining run_dir file writes can continue shrinking.
This commit is contained in:
Bryan Helmkamp 2026-04-03 11:19:23 -07:00
parent fdf7323948
commit a65941e418
3 changed files with 54 additions and 350 deletions

View file

@ -1,4 +1,4 @@
use std::io::{BufRead, BufReader, IsTerminal, Write};
use std::io::{IsTerminal, Write};
use std::path::{Path, PathBuf};
use std::process::ExitCode;
use std::sync::Arc;
@ -13,8 +13,7 @@ use fabro_interview::{AnswerValue, ConsoleInterviewer};
use fabro_store::{EventEnvelope, RuntimeState, SlateRunStore};
use fabro_util::terminal::Styles;
use fabro_workflow::outcome::StageStatus;
use fabro_workflow::records::{Conclusion, ConclusionExt};
use fabro_workflow::run_status::{RunStatus, RunStatusRecord, RunStatusRecordExt};
use fabro_workflow::run_status::RunStatus;
use serde_json::{Map, Value};
use tokio::signal::ctrl_c;
use tokio::time::{self, sleep};
@ -76,32 +75,22 @@ pub(crate) async fn attach_run(
.await;
}
Err(err) => {
tracing::warn!(
run_id = %run_id,
error = %err,
"Failed to list events from store; falling back to filesystem attach"
);
return Err(anyhow::anyhow!(
"Failed to list events from SlateDB for run {run_id}: {err}"
));
}
},
Err(err) => {
tracing::warn!(
run_id = %run_id,
error = %err,
"Failed to open store reader; falling back to filesystem attach"
);
return Err(anyhow::anyhow!(
"Failed to open SlateDB reader for run {run_id}: {err}"
));
}
}
}
attach_run_files(
run_dir,
false,
kill_on_detach,
styles,
engine_child,
json_output,
)
.await
Err(anyhow::anyhow!(
"Could not infer SlateDB storage location and run id for attach"
))
}
async fn attach_run_store(
@ -316,209 +305,6 @@ async fn flush_remaining_store_events(
Ok(())
}
async fn attach_run_files(
run_dir: &Path,
verbose: bool,
kill_on_detach: bool,
styles: &'static Styles,
engine_child: Option<std::process::Child>,
json_output: bool,
) -> Result<ExitCode> {
let progress_path = run_dir.join("progress.jsonl");
let conclusion_path = run_dir.join("conclusion.json");
let status_path = run_dir.join("status.json");
let runtime_state = RuntimeState::new(run_dir);
let runtime_interview_paths = InterviewPaths::from_runtime_state(&runtime_state);
let mut engine_guard = engine_child.map(EngineChildGuard::new);
let is_tty = std::io::stderr().is_terminal();
let mut progress_ui = run_progress::ProgressUI::new(is_tty, verbose);
let cancelled = Arc::new(AtomicBool::new(false));
{
let cancelled = Arc::clone(&cancelled);
tokio::spawn(async move {
let _ = ctrl_c().await;
cancelled.store(true, Ordering::Relaxed);
});
}
let mut reader = open_progress_reader(&progress_path)?;
let mut line = String::new();
let mut cached_pid: Option<u32> = None;
let attach_started = Instant::now();
loop {
if cancelled.load(Ordering::Relaxed) {
if kill_on_detach {
if let Some(guard) = engine_guard.as_mut() {
if let Some(child) = guard.inner() {
let _ = child.kill();
}
} else {
kill_engine(run_dir);
}
for _ in 0..20 {
if conclusion_path.exists()
|| read_status_record(&status_path)
.is_some_and(|record| record.status.is_terminal())
{
break;
}
sleep(Duration::from_millis(100)).await;
}
} else {
if let Some(guard) = engine_guard.as_mut() {
guard.defuse();
}
eprintln!("Detached from run (engine continues in background)");
}
break;
}
if reader.is_none() {
reader = open_progress_reader(&progress_path)?;
}
if let Some(reader) = reader.as_mut() {
loop {
line.clear();
let bytes_read = reader.read_line(&mut line)?;
if bytes_read == 0 {
break;
}
let trimmed = line.trim();
if !trimmed.is_empty() {
emit_progress_line(&mut progress_ui, trimmed, json_output)?;
}
}
}
if runtime_interview_paths.request_path.exists() {
let interview_paths = &runtime_interview_paths;
if !interview_paths.response_path.exists() {
if json_output {
defuse_engine_child(&mut engine_guard);
eprintln!("{JSON_INTERVIEW_MESSAGE}");
return Ok(ExitCode::from(1));
}
if let Some(_claim_guard) =
InterviewClaimGuard::acquire(&interview_paths.claim_path)
{
if let Ok(request_data) = std::fs::read_to_string(&interview_paths.request_path)
{
if let Ok(question) =
serde_json::from_str::<fabro_interview::Question>(&request_data)
{
hide_progress(&mut progress_ui, json_output);
let interviewer = ConsoleInterviewer::new(styles);
let answer =
fabro_interview::Interviewer::ask(&interviewer, question).await;
show_progress(&mut progress_ui, json_output);
if answer_requires_reattach(&answer) {
if let Some(guard) = engine_guard.as_mut() {
guard.defuse();
}
eprintln!("{INTERVIEW_UNANSWERED_MESSAGE}");
return Ok(ExitCode::from(1));
}
write_interview_response_atomically(
&interview_paths.response_path,
&answer,
)?;
}
}
}
}
}
let terminal_status = read_status_record(&status_path)
.map(|record| record.status)
.filter(|status| status.is_terminal());
let child_alive_via_handle = engine_guard.as_mut().and_then(|guard| {
guard.inner().map(|child| match child.try_wait() {
Ok(None) => true,
Ok(Some(_)) | Err(_) => false,
})
});
if let Some(child_alive) = child_alive_via_handle {
if !child_alive {
drain_remaining(reader.as_mut(), &mut line, &mut progress_ui, json_output)?;
break;
}
} else {
if terminal_status.is_some() {
drain_remaining(reader.as_mut(), &mut line, &mut progress_ui, json_output)?;
break;
}
let engine_alive = match cached_pid {
Some(pid) => process_alive(pid),
None => {
if let Some(pid) = read_launcher_pid(run_dir) {
cached_pid = Some(pid);
process_alive(pid)
} else {
attach_started.elapsed() < ATTACH_STARTUP_GRACE
}
}
};
if !engine_alive {
drain_remaining(reader.as_mut(), &mut line, &mut progress_ui, json_output)?;
break;
}
}
sleep(Duration::from_millis(100)).await;
}
finish_progress(&mut progress_ui, json_output);
Ok(determine_exit_code(
&conclusion_path,
read_status_record(&status_path),
))
}
fn open_progress_reader(progress_path: &Path) -> Result<Option<BufReader<std::fs::File>>> {
if !progress_path.exists() {
return Ok(None);
}
Ok(Some(BufReader::new(std::fs::File::open(progress_path)?)))
}
fn drain_remaining(
reader: Option<&mut BufReader<std::fs::File>>,
line: &mut String,
progress_ui: &mut run_progress::ProgressUI,
json_output: bool,
) -> Result<()> {
let Some(reader) = reader else {
return Ok(());
};
loop {
line.clear();
match reader.read_line(line) {
Ok(0) | Err(_) => break,
Ok(_) => {
let trimmed = line.trim();
if !trimmed.is_empty() {
emit_progress_line(progress_ui, trimmed, json_output)?;
}
}
}
}
Ok(())
}
fn emit_progress_line(
progress_ui: &mut run_progress::ProgressUI,
line: &str,
@ -573,10 +359,6 @@ fn normalize_json_value(value: Value) -> Value {
}
}
fn read_status_record(path: &Path) -> Option<RunStatusRecord> {
RunStatusRecord::load(path).ok()
}
fn read_launcher_pid(run_dir: &Path) -> Option<u32> {
super::launcher::active_launcher_record_for_run(run_dir).map(|record| record.pid)
}
@ -719,27 +501,6 @@ fn write_interview_response_atomically(
Ok(())
}
fn determine_exit_code(conclusion_path: &Path, status_record: Option<RunStatusRecord>) -> ExitCode {
if conclusion_path.exists() {
if let Ok(conclusion) = Conclusion::load(conclusion_path) {
let success = matches!(
conclusion.status,
StageStatus::Success | StageStatus::PartialSuccess
);
return if success {
ExitCode::from(0)
} else {
ExitCode::from(1)
};
}
}
match status_record.map(|record| record.status) {
Some(RunStatus::Succeeded) => ExitCode::from(0),
Some(_) | None => ExitCode::from(1),
}
}
async fn determine_exit_code_with_store(run_store: &SlateRunStore) -> ExitCode {
let deadline = Instant::now() + ATTACH_FINAL_STATUS_GRACE;
loop {
@ -789,75 +550,40 @@ mod tests {
use chrono::Utc;
use fabro_interview::{Answer, AnswerValue};
use fabro_util::terminal::Styles;
use fabro_workflow::outcome::StageStatus;
use fabro_workflow::records::Conclusion;
use fabro_workflow::run_status::{StatusReason, write_run_status};
fn no_color_styles() -> &'static Styles {
Box::leak(Box::new(Styles::new(false)))
}
fn sample_conclusion(status: StageStatus) -> Conclusion {
Conclusion {
timestamp: Utc::now(),
status,
duration_ms: 0,
failure_reason: None,
final_git_commit_sha: None,
stages: Vec::new(),
total_cost: None,
total_retries: 0,
total_input_tokens: 0,
total_output_tokens: 0,
total_cache_read_tokens: 0,
total_cache_write_tokens: 0,
total_reasoning_tokens: 0,
has_pricing: false,
}
}
#[tokio::test]
async fn attach_does_not_return_when_only_conclusion_exists() {
async fn attach_errors_without_store_context_even_if_conclusion_file_exists() {
let dir = tempfile::tempdir().unwrap();
sample_conclusion(StageStatus::Success)
.save(&dir.path().join("conclusion.json"))
.unwrap();
std::fs::write(dir.path().join("conclusion.json"), "{}").unwrap();
let child = std::process::Command::new("sh")
.args(["-c", "sleep 0.35"])
.spawn()
.unwrap();
let started = Instant::now();
let exit = attach_run(
let err = attach_run(
dir.path(),
None,
None,
false,
no_color_styles(),
Some(child),
None,
false,
)
.await
.unwrap();
.unwrap_err();
assert_eq!(exit, ExitCode::from(0));
assert!(
started.elapsed() >= Duration::from_millis(250),
"attach returned before the owned child exited"
err.to_string()
.contains("Could not infer SlateDB storage location and run id for attach")
);
}
#[tokio::test]
async fn attach_missing_pid_and_failed_status_is_not_alive() {
async fn attach_errors_without_store_context_even_if_status_file_exists() {
let dir = tempfile::tempdir().unwrap();
write_run_status(
dir.path(),
RunStatus::Failed,
Some(StatusReason::LaunchFailed),
);
std::fs::write(dir.path().join("status.json"), "{}").unwrap();
let exit = attach_run(
let err = attach_run(
dir.path(),
None,
None,
@ -867,9 +593,12 @@ mod tests {
false,
)
.await
.unwrap();
.unwrap_err();
assert_eq!(exit, ExitCode::from(1));
assert!(
err.to_string()
.contains("Could not infer SlateDB storage location and run id for attach")
);
}
#[test]

View file

@ -1,5 +1,3 @@
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicI64, Ordering};
@ -1453,27 +1451,6 @@ pub fn build_redacted_event_payload(
event_payload_from_redacted_json(&line, run_id)
}
pub fn append_progress_event(run_dir: &Path, envelope: &RunEventEnvelope) -> Result<()> {
let line = redacted_event_json(envelope)?;
append_progress_event_with_line(run_dir, &line)
}
pub fn append_progress_event_with_line(run_dir: &Path, line: &str) -> Result<()> {
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(run_dir.join("progress.jsonl"))
.with_context(|| {
format!(
"Failed to open {}",
run_dir.join("progress.jsonl").display()
)
})?;
writeln!(file, "{line}")?;
Ok(())
}
pub fn redacted_event_json(envelope: &RunEventEnvelope) -> Result<String> {
let line = serde_json::to_string(&normalized_envelope_value(envelope)?)?;
Ok(redact_jsonl_line(&line))
@ -1521,26 +1498,6 @@ pub async fn append_workflow_event(
.map_err(anyhow::Error::from)
}
pub struct ProgressLogger {
run_dir: PathBuf,
}
impl ProgressLogger {
#[must_use]
pub fn new(run_dir: impl Into<PathBuf>) -> Self {
Self {
run_dir: run_dir.into(),
}
}
pub fn register(self, emitter: &EventEmitter) {
let run_dir = self.run_dir;
emitter.on_event(move |event| {
let _ = append_progress_event(&run_dir, event);
});
}
}
enum StoreProgressCommand {
Event(EventPayload),
Flush(oneshot::Sender<()>),
@ -1927,9 +1884,17 @@ mod tests {
assert_eq!(envelope.properties["duration_ms"], 900);
}
#[test]
fn append_progress_event_writes_envelope_shape() {
let dir = tempfile::tempdir().unwrap();
#[tokio::test]
async fn append_workflow_event_writes_store_event_shape() {
let store = fabro_store::SlateStore::new(
std::sync::Arc::new(object_store::memory::InMemory::new()),
"",
std::time::Duration::from_millis(1),
);
let run_store = store
.create_run(&fixtures::RUN_7, Utc::now(), None)
.await
.unwrap();
let envelope = canonicalize_event(
&fixtures::RUN_7,
&WorkflowRunEvent::RunNotice {
@ -1939,10 +1904,15 @@ mod tests {
},
);
append_progress_event(dir.path(), &envelope).unwrap();
let payload = build_redacted_event_payload(&envelope, &fixtures::RUN_7).unwrap();
run_store.append_event(&payload).await.unwrap();
let progress = std::fs::read_to_string(dir.path().join("progress.jsonl")).unwrap();
let line: serde_json::Value = serde_json::from_str(progress.trim()).unwrap();
let events = run_store.list_events().await.unwrap();
let line = events
.into_iter()
.next()
.map(|event| event.payload.as_value().clone())
.unwrap();
assert!(line.get("id").is_some());
assert_eq!(line["event"], "run.notice");
assert_eq!(line["properties"]["code"], "example");

View file

@ -229,15 +229,20 @@ fn read_plan_text(state: &RunState) -> Option<String> {
.nodes
.iter()
.filter_map(|((node_id, visit), node)| {
node_id
.starts_with("plan")
.then_some((node_id.as_str(), *visit, node.response.as_deref()))
node_id.starts_with("plan").then_some((
node_id.as_str(),
*visit,
node.response.as_deref(),
))
})
.collect::<Vec<_>>();
plan_nodes.sort_by(|left, right| left.0.cmp(right.0).then(left.1.cmp(&right.1)));
for (node_id, visit, response) in plan_nodes {
if let Some(response) = response {
debug!(node_id, visit, "Found plan node response for PR body from run state");
debug!(
node_id,
visit, "Found plan node response for PR body from run state"
);
return Some(response.to_string());
}
}