Emit pending questions before JSON attach exits

This commit is contained in:
Scott Werner 2026-09-29 13:12:43 -04:00
parent 14963e7d06
commit 0537fbab8f
3 changed files with 217 additions and 2 deletions

View file

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

View file

@ -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<Value> = 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::<Vec<_>>(),
items
.iter()
.map(|item| &item["stream_seq"])
.collect::<Vec<_>>(),
"{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."

View file

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