mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
Merge pull request #911 from fabro-sh/codex/json-human-gate-stream-race
Emit pending human-gate questions before JSON attach exits
This commit is contained in:
commit
e58bccee64
3 changed files with 217 additions and 2 deletions
|
|
@ -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)));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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."
|
||||
|
|
|
|||
|
|
@ -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!();
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue