Tighten non-interactive JSON mode

This commit is contained in:
Bryan Helmkamp 2026-03-31 09:39:10 -04:00
parent fb6f0eae1e
commit a6829c57a0
No known key found for this signature in database
11 changed files with 590 additions and 15 deletions

View file

@ -28,6 +28,7 @@ const ATTACH_STARTUP_GRACE: Duration = Duration::from_millis(200);
const ATTACH_STARTUP_GRACE: Duration = Duration::from_secs(3);
const INTERVIEW_UNANSWERED_MESSAGE: &str =
"Interview ended without an answer. The run is still waiting for input; reattach to answer it.";
const JSON_INTERVIEW_MESSAGE: &str = "This run is waiting for human input, but --json is non-interactive. Reattach without --json to answer it.";
/// Attach to a running (or finished) workflow run, rendering progress live.
///
@ -177,6 +178,11 @@ async fn attach_run_store(
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)
{
@ -390,6 +396,11 @@ async fn attach_run_files(
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)
{
@ -605,6 +616,12 @@ impl EngineChildGuard {
}
}
fn defuse_engine_child(engine_guard: &mut Option<EngineChildGuard>) {
if let Some(guard) = engine_guard.as_mut() {
guard.defuse();
}
}
impl Drop for EngineChildGuard {
fn drop(&mut self) {
if let Some(mut child) = self.child.take() {

View file

@ -3,7 +3,7 @@ use fabro_config::FabroSettingsExt;
use fabro_util::terminal::Styles;
use fabro_workflow::run_lookup::{resolve_run_combined, runs_base};
use crate::args::{GlobalArgs, RunCommands};
use crate::args::{GlobalArgs, RunArgs, RunCommands};
use crate::shared::print_json_pretty;
use crate::store;
use crate::user_config::{load_user_settings_with_globals, user_layer_with_globals};
@ -31,10 +31,20 @@ pub(super) fn short_run_id(id: &str) -> &str {
if id.len() > 12 { &id[..12] } else { id }
}
fn apply_json_defaults(args: &mut RunArgs, globals: &GlobalArgs) {
if globals.json {
args.auto_approve = true;
}
}
pub(crate) async fn dispatch(cmd: RunCommands, globals: &GlobalArgs) -> Result<()> {
match cmd {
RunCommands::Run(args) => command::execute(args, globals).await,
RunCommands::Create(args) => {
RunCommands::Run(mut args) => {
apply_json_defaults(&mut args, globals);
command::execute(args, globals).await
}
RunCommands::Create(mut args) => {
apply_json_defaults(&mut args, globals);
let styles: &'static Styles = Box::leak(Box::new(Styles::detect_stderr()));
let cli = user_layer_with_globals(globals)?;
let (run_id, _run_dir) = create::create_run(&args, cli, styles, true)?;

View file

@ -112,6 +112,10 @@ async fn remove_from(
had_errors = true;
continue;
}
removed.push(run_id.clone());
if !globals.json {
eprintln!("{}", short_run_id(&run_id));
}
if let Err(err) = store
.delete_run(&run.run_id)
.await
@ -127,10 +131,6 @@ async fn remove_from(
had_errors = true;
continue;
}
removed.push(run_id.clone());
if !globals.json {
eprintln!("{}", short_run_id(&run_id));
}
}
if globals.json {

View file

@ -21,6 +21,7 @@ pub(super) fn create_command(args: &WorkflowCreateArgs, globals: &GlobalArgs) ->
let created = write_workflow_scaffold(args, &fabro_root)?;
if globals.json {
let created: Vec<_> = created.iter().map(|path| relative_path(path)).collect();
print_json_pretty(&serde_json::json!({
"name": args.name,
"created": created,
@ -63,7 +64,10 @@ pub(super) fn create_command(args: &WorkflowCreateArgs, globals: &GlobalArgs) ->
Ok(())
}
fn write_workflow_scaffold(args: &WorkflowCreateArgs, fabro_root: &Path) -> Result<Vec<String>> {
fn write_workflow_scaffold(
args: &WorkflowCreateArgs,
fabro_root: &Path,
) -> Result<Vec<std::path::PathBuf>> {
let workflows_dir = fabro_root.join("workflows").join(&args.name);
if workflows_dir.exists() {
@ -103,10 +107,7 @@ fn write_workflow_scaffold(args: &WorkflowCreateArgs, fabro_root: &Path) -> Resu
std::fs::write(&toml_path, "version = 1\n")
.with_context(|| format!("failed to write {}", toml_path.display()))?;
Ok(vec![
format!("fabro/workflows/{}/workflow.fabro", args.name),
format!("fabro/workflows/{}/workflow.toml", args.name),
])
Ok(vec![dot_path, toml_path])
}
fn to_pascal_case(s: &str) -> String {

View file

@ -1,6 +1,9 @@
use fabro_test::{fabro_snapshot, test_context};
use serde_json::Value;
use crate::support::{example_fixture, run_output_filters};
use crate::support::{
compact_progress_event, example_fixture, fabro_json_snapshot, run_output_filters,
};
use super::support::{output_stdout, write_sleep_workflow};
@ -146,3 +149,149 @@ fn attach_before_completion_streams_to_finished_state() {
✓ exit [DURATION]
");
}
#[test]
fn attach_json_errors_without_prompting_for_human_input() {
let context = test_context!();
let workflow = context.temp_dir.join("human-gate.fabro");
context.write_temp(
"human-gate.fabro",
r#"digraph HumanGate {
graph [goal="Wait for approval"]
start [shape=Mdiamond, label="Start"]
exit [shape=Msquare, label="Exit"]
approve [shape=hexagon, label="Approve?"]
ship [shape=parallelogram, script="echo shipped"]
revise [shape=parallelogram, script="echo revised"]
start -> approve
approve -> ship [label="[A] Approve"]
approve -> revise [label="[R] Revise"]
ship -> exit
revise -> exit
}
"#,
);
let run_output = context
.command()
.current_dir(&context.temp_dir)
.env("OPENAI_API_KEY", "test")
.args([
"run",
"--detach",
"--no-retro",
"--sandbox",
"local",
"--provider",
"openai",
workflow.to_str().unwrap(),
])
.output()
.expect("detached run should execute");
assert!(
run_output.status.success(),
"detached run failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&run_output.stdout),
String::from_utf8_lossy(&run_output.stderr)
);
let run_id = output_stdout(&run_output).trim().to_string();
let cleanup_run_id = run_id.clone();
scopeguard::defer! {
let _ = context.command().args(["rm", "--force", &cleanup_run_id]).output();
}
let run_dir = context.find_run_dir(&run_id);
let request_path = run_dir.join("runtime/interview_request.json");
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
while !request_path.exists() {
assert!(
std::time::Instant::now() < deadline,
"timed out waiting for interview request for {run_id}"
);
std::thread::sleep(std::time::Duration::from_millis(50));
}
let output = context
.command()
.args(["--json", "attach", &run_id])
.timeout(std::time::Duration::from_secs(5))
.output()
.expect("attach should execute");
assert!(!output.status.success(), "attach --json should fail fast");
let stderr = String::from_utf8(output.stderr).expect("stderr should be UTF-8");
assert!(stderr.contains("--json is non-interactive"));
assert!(
!stderr.contains("Approve?"),
"attach should not prompt on stderr"
);
assert!(
request_path.exists(),
"the run should still be waiting on the interview request"
);
assert!(
!run_dir.join("runtime/interview_response.json").exists(),
"attach --json should not answer the interview"
);
let progress: Vec<Value> = String::from_utf8(output.stdout)
.expect("stdout should be UTF-8")
.lines()
.filter(|line| !line.trim().is_empty())
.map(|line| serde_json::from_str(line).expect("attach JSON output should be JSONL"))
.collect();
let progress_summary: Vec<_> = progress.iter().map(compact_progress_event).collect();
fabro_json_snapshot!(context, &progress_summary, @r#"
[
{
"event": "sandbox.initializing",
"provider": "local"
},
{
"event": "sandbox.ready",
"provider": "local"
},
{
"event": "sandbox.initialized"
},
{
"event": "run.started",
"name": "HumanGate",
"goal": "Wait for approval"
},
{
"event": "stage.started",
"node_id": "start",
"node_label": "Start",
"handler_type": "start",
"index": 0
},
{
"event": "stage.completed",
"node_id": "start",
"node_label": "Start",
"index": 0,
"status": "success"
},
{
"event": "edge.selected",
"from_node": "start",
"to_node": "approve",
"reason": "unconditional"
},
{
"event": "checkpoint.completed",
"node_id": "start",
"node_label": "start",
"status": "success"
},
{
"event": "stage.started",
"node_id": "approve",
"node_label": "Approve?",
"handler_type": "human",
"index": 1
}
]
"#);
}

View file

@ -251,6 +251,37 @@ fn create_persists_requested_overrides_into_run_json() {
"###);
}
#[test]
fn create_json_implies_auto_approve() {
let context = test_context!();
let workflow = fixture("simple.fabro");
let output = context
.command()
.args(["--json", "create", "--dry-run", workflow.to_str().unwrap()])
.output()
.expect("command should execute");
assert!(
output.status.success(),
"command failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
let value: serde_json::Value =
serde_json::from_slice(&output.stdout).expect("create JSON should parse");
let run_id = value["run_id"]
.as_str()
.expect("create JSON should include run_id");
let run = resolve_run(&context, run_id);
let run_json = read_json(run.run_dir.join("run.json"));
assert_eq!(
run_json.pointer("/settings/auto_approve"),
Some(&json!(true))
);
}
#[test]
fn create_invalid_workflow_fails_without_creating_run() {
let context = test_context!();

View file

@ -1,4 +1,8 @@
use fabro_test::{fabro_snapshot, test_context};
#[cfg(feature = "server")]
use httpmock::prelude::*;
#[cfg(feature = "server")]
use serde_json::Value;
#[test]
fn help() {
@ -33,3 +37,53 @@ fn help() {
----- stderr -----
");
}
#[test]
#[cfg(feature = "server")]
fn prompt_json_streaming_server_reports_resolved_model() {
let context = test_context!();
let server = MockServer::start();
let sse_body = "\
event: stream_event\n\
data: {\"type\":\"text_delta\",\"delta\":\"Hi\",\"text_id\":null}\n\
\n\
event: stream_event\n\
data: {\"type\":\"finish\",\"finish_reason\":\"stop\",\"usage\":{\"input_tokens\":5,\"output_tokens\":2,\"total_tokens\":7},\"response\":{\"id\":\"r1\",\"model\":\"resolved-model\",\"provider\":\"test-provider\",\"message\":{\"role\":\"assistant\",\"content\":[{\"kind\":\"text\",\"data\":\"Hi\"}],\"name\":null,\"tool_call_id\":null},\"finish_reason\":\"stop\",\"usage\":{\"input_tokens\":5,\"output_tokens\":2,\"total_tokens\":7},\"raw\":null,\"warnings\":[],\"rate_limit\":null}}\n\
\n";
let mock = server.mock(|when, then| {
when.method(POST).path("/completions");
then.status(200)
.header("content-type", "text/event-stream")
.body(sse_body);
});
let output = context
.command()
.env_remove("FABRO_STORAGE_DIR")
.args([
"--server-url",
&server.url(""),
"--json",
"llm",
"prompt",
"Hello",
])
.output()
.expect("command should run");
assert!(
output.status.success(),
"command failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
let value: Value =
serde_json::from_slice(&output.stdout).expect("llm prompt JSON should parse");
assert_eq!(value["response"], "Hi");
assert_eq!(value["model"], "resolved-model");
assert_eq!(value["usage"]["input_tokens"], 5);
assert_eq!(value["usage"]["output_tokens"], 2);
mock.assert();
}

View file

@ -2,6 +2,7 @@ use fabro_test::{fabro_snapshot, test_context};
use serde_json::Value;
use super::support::{setup_completed_dry_run, setup_created_dry_run};
use walkdir::WalkDir;
#[test]
fn help() {
@ -171,3 +172,71 @@ fn rm_partial_failure_json_includes_removed_and_errors() {
"existing run should still be removed"
);
}
#[test]
fn rm_json_reports_removed_run_when_store_delete_fails() {
let context = test_context!();
let run = setup_completed_dry_run(&context);
let by_id_path = find_store_catalog_entry(&context.storage_dir.join("store"), &run.run_id);
let original = std::fs::read(&by_id_path)
.unwrap_or_else(|err| panic!("failed to read {}: {err}", by_id_path.display()));
// This intentionally mutates the backing store metadata to force a store-only
// delete failure. Reproducing that failure through public commands is not
// practical, and the command contract under test is still `fabro rm`'s JSON
// partial-success reporting.
std::fs::write(&by_id_path, b"{not valid json")
.unwrap_or_else(|err| panic!("failed to corrupt {}: {err}", by_id_path.display()));
scopeguard::defer! {
let _ = std::fs::write(&by_id_path, &original);
}
let output = context
.command()
.args(["--json", "rm", &run.run_id])
.output()
.expect("command should run");
assert!(
!output.status.success(),
"rm should report the store failure"
);
let value: Value = serde_json::from_slice(&output.stdout).expect("rm JSON should parse");
assert_eq!(
value["removed"],
Value::Array(vec![Value::String(run.run_id.clone())])
);
assert_eq!(value["errors"][0]["identifier"], run.run_id);
assert!(
value["errors"][0]["error"]
.as_str()
.is_some_and(|error| error.contains("failed to delete store state"))
);
assert!(
!run.run_dir.exists(),
"run directory should still be deleted"
);
}
fn find_store_catalog_entry(root: &std::path::Path, run_id: &str) -> std::path::PathBuf {
let expected_name = format!("{run_id}.json");
WalkDir::new(root)
.into_iter()
.filter_map(Result::ok)
.map(|entry| entry.into_path())
.find(|path| {
path.is_file()
&& path
.file_name()
.is_some_and(|name| name.to_string_lossy() == expected_name)
&& path
.components()
.any(|component| component.as_os_str() == "by-id")
})
.unwrap_or_else(|| {
panic!(
"missing by-id catalog entry for {run_id} under {}",
root.display()
)
})
}

View file

@ -1,4 +1,5 @@
use fabro_test::{fabro_snapshot, test_context};
use serde_json::Value;
use crate::support::{
compact_progress_event, example_fixture, fabro_json_snapshot, read_json, read_jsonl,
@ -264,6 +265,193 @@ fn run_id_passthrough_uses_provided_ulid() {
assert_eq!(run_record["run_id"].as_str(), Some(run_id));
}
#[test]
fn json_run_implies_auto_approve_for_human_gates() {
let context = test_context!();
let workflow = context.temp_dir.join("human-gate.fabro");
context.write_temp(
"human-gate.fabro",
r#"digraph HumanGate {
graph [goal="Route through the default approval path"]
start [shape=Mdiamond, label="Start"]
exit [shape=Msquare, label="Exit"]
approve [shape=hexagon, label="Approve?"]
ship [shape=parallelogram, script="echo shipped"]
revise [shape=parallelogram, script="echo revised"]
start -> approve
approve -> ship [label="[A] Approve"]
approve -> revise [label="[R] Revise"]
ship -> exit
revise -> exit
}
"#,
);
let output = context
.command()
.current_dir(&context.temp_dir)
.args([
"--json",
"run",
"--sandbox",
"local",
"--no-retro",
workflow.to_str().unwrap(),
])
.output()
.expect("command should execute");
assert!(
output.status.success(),
"command failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
let progress: Vec<Value> = String::from_utf8(output.stdout)
.expect("stdout should be UTF-8")
.lines()
.filter(|line| !line.trim().is_empty())
.map(|line| serde_json::from_str(line).expect("run JSON output should be JSONL"))
.collect();
let progress_summary: Vec<_> = progress.iter().map(compact_progress_event).collect();
fabro_json_snapshot!(context, &progress_summary, @r#"
[
{
"event": "sandbox.initializing",
"provider": "local"
},
{
"event": "sandbox.ready",
"provider": "local"
},
{
"event": "sandbox.initialized"
},
{
"event": "run.notice"
},
{
"event": "run.started",
"name": "HumanGate",
"goal": "Route through the default approval path"
},
{
"event": "stage.started",
"node_id": "start",
"node_label": "Start",
"handler_type": "start",
"index": 0
},
{
"event": "stage.completed",
"node_id": "start",
"node_label": "Start",
"index": 0,
"status": "success"
},
{
"event": "edge.selected",
"from_node": "start",
"to_node": "approve",
"reason": "unconditional"
},
{
"event": "checkpoint.completed",
"node_id": "start",
"node_label": "start",
"status": "success"
},
{
"event": "stage.started",
"node_id": "approve",
"node_label": "Approve?",
"handler_type": "human",
"index": 1
},
{
"event": "stage.completed",
"node_id": "approve",
"node_label": "Approve?",
"index": 1,
"status": "success"
},
{
"event": "edge.selected",
"from_node": "approve",
"to_node": "ship",
"reason": "preferred_label"
},
{
"event": "checkpoint.completed",
"node_id": "approve",
"node_label": "approve",
"status": "success"
},
{
"event": "stage.started",
"node_id": "ship",
"node_label": "ship",
"handler_type": "command",
"index": 2
},
{
"event": "stage.completed",
"node_id": "ship",
"node_label": "ship",
"index": 2,
"status": "success"
},
{
"event": "edge.selected",
"from_node": "ship",
"to_node": "exit",
"reason": "unconditional"
},
{
"event": "checkpoint.completed",
"node_id": "ship",
"node_label": "ship",
"status": "success"
},
{
"event": "stage.started",
"node_id": "exit",
"node_label": "Exit",
"handler_type": "exit",
"index": 3
},
{
"event": "stage.completed",
"node_id": "exit",
"node_label": "Exit",
"index": 3,
"status": "success"
},
{
"event": "run.completed",
"status": "success",
"artifact_count": 0
},
{
"event": "sandbox.cleanup.started",
"provider": "local"
},
{
"event": "sandbox.cleanup.completed",
"provider": "local"
}
]
"#);
let run = context.single_run_dir();
let run_json = read_json(run.join("run.json"));
assert_eq!(
run_json.pointer("/settings/auto_approve"),
Some(&serde_json::json!(true))
);
}
#[test]
fn detach_prints_ulid_and_exits() {
let context = test_context!();

View file

@ -1,7 +1,10 @@
use insta::assert_snapshot;
use serde_json::Value;
use fabro_test::{fabro_snapshot, test_context};
use crate::support::fabro_json_snapshot;
use super::support::setup_project_fixture;
#[test]
@ -155,3 +158,50 @@ fn workflow_create_errors_without_project_config() {
error: No fabro.toml found in [TEMP_DIR] or any parent directory
");
}
#[test]
fn workflow_create_json_uses_resolved_custom_root_paths() {
let context = test_context!();
let project_dir = context.temp_dir.join("project");
context.write_temp(
"project/fabro.toml",
"version = 1\n[fabro]\nroot = \"custom/fabro-data\"\n",
);
let output = context
.command()
.current_dir(&project_dir)
.args(["--json", "workflow", "create", "hello-world"])
.output()
.expect("command should run");
assert!(
output.status.success(),
"command failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
let value: Value =
serde_json::from_slice(&output.stdout).expect("workflow create JSON should parse");
fabro_json_snapshot!(context, &value, @r#"
{
"name": "hello-world",
"created": [
"custom/fabro-data/workflows/hello-world/workflow.fabro",
"custom/fabro-data/workflows/hello-world/workflow.toml"
]
}
"#);
assert!(
project_dir
.join("custom/fabro-data/workflows/hello-world/workflow.fabro")
.exists()
);
assert!(
project_dir
.join("custom/fabro-data/workflows/hello-world/workflow.toml")
.exists()
);
}

View file

@ -557,6 +557,7 @@ pub async fn run_prompt_via_server(
let show_usage = args.usage;
let mut output_usage: Option<Usage> = None;
let mut output_model = args.model.clone();
let mut full_text = String::new();
parse_sse_frames(response, |event_type, data| {
@ -571,8 +572,13 @@ pub async fn run_prompt_via_server(
let _ = io::stdout().flush();
}
}
StreamEvent::Finish { usage, .. } => {
StreamEvent::Finish {
usage, response, ..
} => {
output_usage = Some(usage);
if output_model.is_none() {
output_model = Some(response.model.clone());
}
}
StreamEvent::Error { error, .. } => {
bail!("Server error: {error}");
@ -587,7 +593,7 @@ pub async fn run_prompt_via_server(
if json_output {
let mut value = serde_json::Map::new();
value.insert("response".to_string(), full_text.into());
if let Some(model) = args.model {
if let Some(model) = output_model {
value.insert("model".to_string(), model.into());
}
if let Some(usage) = output_usage {