From 0537fbab8f99dddfaaa95f7aad85555fe0c0b965 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Tue, 29 Sep 2026 13:12:43 -0400 Subject: [PATCH] Emit pending questions before JSON attach exits --- lib/apps/fabro-cli/src/commands/run/attach.rs | 8 + lib/apps/fabro-cli/tests/it/cmd/attach.rs | 167 +++++++++++++++++- lib/apps/fabro-cli/tests/it/cmd/run.rs | 44 +++++ 3 files changed, 217 insertions(+), 2 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/attach.rs b/lib/apps/fabro-cli/src/commands/run/attach.rs index 379deb3b2..b7c51ea1a 100644 --- a/lib/apps/fabro-cli/src/commands/run/attach.rs +++ b/lib/apps/fabro-cli/src/commands/run/attach.rs @@ -322,6 +322,14 @@ async fn handle_pending_petri_interview( }; if json_pending_interview_requires_manual_input(opts.json_output, opts.auto_approve) { + // The pending-question projection may be ahead of our replay/live + // cursor. It commits together with the stream records, so catch up + // before exiting to include the question and everything preceding it. + // We return immediately; buffered live items cannot be emitted twice. + for item in client.list_run_stream(run_id, *cursor).await? { + emit_stream_item(progress_ui, &item, opts.json_output)?; + *cursor = item.stream_seq; + } fabro_util::printerr!(printer, "{JSON_INTERVIEW_MESSAGE}"); return Ok(Some(ExitCode::from(1))); } diff --git a/lib/apps/fabro-cli/tests/it/cmd/attach.rs b/lib/apps/fabro-cli/tests/it/cmd/attach.rs index f08238db6..8c03ba3ca 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/attach.rs @@ -3,8 +3,10 @@ reason = "integration tests: read child-process stdout line-by-line via std::io::BufReader" )] +use std::fmt::Write as _; use std::io::{BufRead, BufReader, Read, Write}; use std::process::{Output, Stdio}; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::mpsc; use std::time::{Duration, Instant}; @@ -12,11 +14,12 @@ use fabro_test::{ apply_filters, assert_reqwest_status, expect_reqwest_json, fabro_json_snapshot, fabro_snapshot, test_context, }; +use httpmock::{HttpMockResponse, MockServer}; use serde_json::Value; use super::support::{ - created_run_id, output_stdout, resolve_run, server_endpoint, wait_for_status, - write_gated_workflow, + created_run_id, output_stdout, remote_run_summary_json, resolve_run, server_endpoint, + wait_for_status, write_gated_workflow, }; use crate::support::run_output_filters; @@ -227,6 +230,166 @@ fn wait_for_pending_question(context: &fabro_test::TestContext, run_id: &str) { }); } +#[test] +fn attach_json_emits_pending_question_once_across_replay_and_live_boundaries() { + let context = test_context!(); + let run_id = start_detached_human_run( + &context, + "stream-boundary.fabro", + r#"digraph HumanGate { + start [shape=Mdiamond] + approve [shape=hexagon, label="Approve?"] + exit [shape=Msquare] + start -> approve + approve -> exit [label="[A] Approve"] +} +"#, + ); + scopeguard::defer! { + let _ = context.command().args(["rm", "--force", &run_id]).output(); + } + + // Capture real public API responses. Only their delivery boundaries are + // scripted below; no run records or runtime files are fabricated. + let (client, base_url) = server_endpoint(&context.storage_dir).unwrap(); + let (question, state, history) = tokio::runtime::Runtime::new().unwrap().block_on(async { + let question = wait_for_server_question(&client, &base_url, &run_id).await; + let mut responses = Vec::new(); + for endpoint in ["state", "events"] { + let url = format!("{base_url}/api/v1/runs/{run_id}/{endpoint}"); + let response = client.get(&url).send().await.unwrap(); + responses.push(expect_reqwest_json(response, fabro_http::StatusCode::OK, &url).await); + } + (question, responses.remove(0), responses.remove(0)) + }); + assert_eq!(history["meta"]["has_more"], false); + let mut items = history["data"].as_array().unwrap().clone(); + let question_index = items + .iter() + .position(|item| { + item.pointer("/item/derived/parsed/kind") == Some(&Value::from("question")) + }) + .expect("a pending question must already have a committed stream record"); + items.truncate(question_index + 1); + assert!(question_index > 4); + + for (boundary, replay_len, initially_pending) in [ + ("between replay and pending check", 4, true), + ("in replay", items.len(), true), + ("in live stream", 4, false), + ] { + let server = MockServer::start(); + let resolve = server.mock(|when, then| { + when.method("GET").path("/api/v1/runs/resolve"); + then.status(200).json_body(remote_run_summary_json( + &run_id, + "HumanGate", + "human-gate", + "Approve?", + &state["status"], + "2026-09-29T09:00:00Z", + )); + }); + server.mock(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/state")); + then.status(200).json_body(state.clone()); + }); + let mut replay = history.clone(); + replay["data"] = serde_json::json!(&items[..replay_len]); + let replay_mock = server.mock(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/events")) + .query_param("after", "0"); + then.status(200).json_body(replay); + }); + // Catch-up reads are deliberately paginated. Replaying from zero, + // skipping an intervening item, or repeating the question changes + // the exact output comparison below. + for index in replay_len..=items.len() { + let mut page = history.clone(); + let end = (index + 2).min(items.len()); + page["data"] = serde_json::json!(&items[index..end]); + page["meta"]["has_more"] = Value::from(end < items.len()); + server.mock(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/events")) + .query_param("after", items[index - 1]["stream_seq"].to_string()); + then.status(200).json_body(page); + }); + } + let mut live_body = String::new(); + for item in &items[replay_len..] { + writeln!(live_body, "data: {item}\n").unwrap(); + } + let stream = server.mock(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/attach")) + .query_param("after", items[replay_len - 1]["stream_seq"].to_string()); + then.status(200) + .header("Content-Type", "text/event-stream") + .body(live_body); + }); + let questions = server.mock(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/questions")); + let calls = AtomicUsize::new(0); + let question = question.clone(); + then.respond_with(move |_| { + let pending = initially_pending || calls.fetch_add(1, Ordering::SeqCst) > 0; + HttpMockResponse::builder() + .status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "data": if pending { vec![question.clone()] } else { vec![] }, + "meta": { "has_more": false } + }) + .to_string(), + ) + .build() + }); + }); + + let output = context + .command() + .args(["--json", "attach", "--server", &server.base_url(), &run_id]) + .timeout(SHARED_DAEMON_TIMEOUT) + .output() + .unwrap(); + assert_eq!(output.status.code(), Some(1), "{boundary}"); + assert_eq!( + String::from_utf8_lossy(&output.stderr).trim(), + "This run is waiting for human input, but --json is non-interactive. Reattach without --json to answer it.", + "{boundary}" + ); + let actual: Vec = String::from_utf8(output.stdout) + .unwrap() + .lines() + .map(|line| serde_json::from_str(line).unwrap()) + .collect(); + assert_eq!( + actual + .iter() + .map(|item| &item["stream_seq"]) + .collect::>(), + items + .iter() + .map(|item| &item["stream_seq"]) + .collect::>(), + "{boundary}: stdout must include every record through the question exactly once" + ); + assert_eq!( + actual, items, + "{boundary}: preserve the original stream envelopes" + ); + resolve.assert_calls(1); + replay_mock.assert_calls(1); + stream.assert_calls(1); + questions.assert_calls(if initially_pending { 1 } else { 2 }); + } +} + #[expect( clippy::disallowed_methods, reason = "This sync integration helper writes scripted answers to an attach child process." diff --git a/lib/apps/fabro-cli/tests/it/cmd/run.rs b/lib/apps/fabro-cli/tests/it/cmd/run.rs index 2931f27f9..f86e2f3ab 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/run.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/run.rs @@ -1024,6 +1024,50 @@ fn run_rejects_removed_run_id_flag() { ); } +#[test] +fn json_run_auto_approves_human_gates() { + let context = test_context!(); + context.ensure_home_server_auth_methods(); + context.write_temp( + "auto-approve.fabro", + r#"digraph HumanGate { + start [shape=Mdiamond] + approve [shape=hexagon, label="Approve?"] + exit [shape=Msquare] + start -> approve + approve -> exit [label="[A] Approve"] +} +"#, + ); + let output = context + .command() + .args([ + "--json", + "run", + "--auto-approve", + "--environment", + "local", + "auto-approve.fabro", + ]) + .assert() + .success(); + let items: Vec = std::str::from_utf8(&output.get_output().stdout) + .unwrap() + .lines() + .map(|line| serde_json::from_str(line).unwrap()) + .collect(); + assert!(items.iter().any(|item| { + item.pointer("/item/derived/parsed/kind") == Some(&Value::from("question")) + })); + assert!( + items.iter().any(|item| { + item.pointer("/item/record/body/event") == Some(&Value::from("run.finished")) + && item.pointer("/item/record/body/status") == Some(&Value::from("success")) + }), + "JSON mode must continue through an automatically approved gate to completion" + ); +} + #[test] fn json_run_requires_manual_input_for_human_gates_without_auto_approve() { let context = test_context!();