From 392b8dd27b94818f1a1af61f72b33df1801e4d5c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 21:24:16 -0400 Subject: [PATCH 1/4] fix(agent): report real shell process outcomes The shell executor rendered every returned ExecResult and returned Ok(output), so nonzero exits, timeouts, and cancellations reached execute_one_tool() as successes. That false ToolResult propagated consistently: agent.tool.completed recorded is_error: false, the success post-tool hook ran, Anthropic saw is_error: false, OpenAI Responses saw a completed function-call output, and CLI/web rendered a successful tool call. ExecResult::is_success() is now the authoritative predicate. The executor runs through exec_command_streaming() with a sink callback, so it keeps the production providers' stream provenance and partial-output capture, and drops the exec 2>&1 prefix that merged stderr into stdout before Fabro could report it. Model-facing text labels termination, exit code, duration, and either separate stdout/stderr sections or one combined section when the provider cannot separate streams. Session-bound dispatch also emits a typed agent.tool.process.completed event carrying the process metadata, streams_separated, and bounded redacted output tails. It is subordinate diagnostic data: the following agent.tool.completed remains the one tool-protocol completion and the authoritative owner of is_error, so consumers need no new row. Nonzero, timed-out, and cancelled commands intentionally change from successful to failed tool results, and PostToolUseFailure replaces PostToolUse for them. On Docker the agent shell tool now uses the streaming path's bash -lc supervisor, which terminates the process group on timeout instead of leaving container-side processes running. The public shell schema is unchanged and pinned by an exact assertion. Co-Authored-By: Claude Opus 5 (1M context) --- docs/internal/events.md | 35 ++ .../fabro-agent/src/tool_execution.rs | 176 ++++++- lib/components/fabro-agent/src/tools.rs | 471 +++++++++++++++--- lib/components/fabro-agent/src/types.rs | 40 +- .../fabro-agent/tests/it/docker_shell.rs | 97 ++++ lib/components/fabro-agent/tests/it/main.rs | 2 + .../fabro-sandbox/src/test_support.rs | 52 +- .../fabro-workflow/src/event/convert.rs | 62 +++ .../fabro-workflow/src/event/names.rs | 1 + .../fabro-workflow/src/event/redaction.rs | 35 ++ .../fabro-workflow/src/event/stored_fields.rs | 3 +- .../fabro-types/src/run_event/agent.rs | 25 +- .../fabro-types/src/run_event/mod.rs | 92 +++- 13 files changed, 1009 insertions(+), 82 deletions(-) create mode 100644 lib/components/fabro-agent/tests/it/docker_shell.rs diff --git a/docs/internal/events.md b/docs/internal/events.md index eea2d6878..57bb9c115 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -1157,6 +1157,41 @@ Emitted when a tool call finishes. | `output` | any | Tool output (string or structured) | | `is_error` | boolean | Whether the tool returned an error | +### `agent.tool.process.completed` + +Subordinate diagnostic for a tool call that ran a process, emitted between +`agent.tool.started` and `agent.tool.completed`. It explains the underlying +process outcome; `agent.tool.completed.is_error` remains the protocol and UI +truth. Absent when the tool never produced a process result (setup, transport, +or launch failure) and when the tool ran without a session-bound emitter. + +```json +{ + "id": "...", "ts": "...", "run_id": "...", + "event": "agent.tool.process.completed", + "node_id": "code", "node_label": "code", + "session_id": "ses_abc", + "tool_call_id": "call_abc123", + "properties": { + "exit_code": 7, + "termination": "exited", + "duration_ms": 812, + "streams_separated": true, + "exec_output_tail": {"stdout": "...", "stderr": "..."}, + "visit": 1 + } +} +``` + +| Property | Type | Description | +|----------|------|-------------| +| `exit_code` | integer | Process exit code; omitted for timeout and cancellation | +| `termination` | string | `exited`, `timed_out`, or `cancelled` | +| `duration_ms` | integer | Process duration | +| `streams_separated` | boolean | `false` when the provider could not separate stdout from stderr; the combined output is then in `exec_output_tail.stdout` | +| `exec_output_tail` | object | Bounded, redacted output tails; omitted when both streams were empty | +| `visit` | integer | Stage visit | + ### `agent.error` Emitted when the agent encounters an error. diff --git a/lib/components/fabro-agent/src/tool_execution.rs b/lib/components/fabro-agent/src/tool_execution.rs index 2c77b3a43..259d04530 100644 --- a/lib/components/fabro-agent/src/tool_execution.rs +++ b/lib/components/fabro-agent/src/tool_execution.rs @@ -558,6 +558,7 @@ mod tests { use async_trait::async_trait; use fabro_llm::types::{ToolCall, ToolDefinition}; use fabro_model::AgentProfileKind; + use tokio::sync::broadcast; use super::*; use crate::config::{ @@ -570,11 +571,13 @@ mod tests { AgentToolRuntime, register_question_tools, }; use crate::read_before_write_sandbox::ReadBeforeWriteSandbox; - use crate::test_support::MutableMockSandbox; + use crate::test_support::{MockSandbox, MutableMockSandbox}; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource}; use crate::tools::{ - make_edit_file_tool, make_grep_tool, make_read_file_tool, make_write_file_tool, + make_edit_file_tool, make_grep_tool, make_read_file_tool, make_shell_tool, + make_write_file_tool, }; + use crate::types::SessionEvent; struct NamedPolicy { decisions: HashMap, @@ -1294,4 +1297,173 @@ mod tests { assert!(!result.is_error); } + + fn shell_sandbox(result: fabro_sandbox::ExecResult) -> Arc { + Arc::new(MockSandbox { + exec_result: result, + ..Default::default() + }) + } + + fn exited(exit_code: i32) -> fabro_sandbox::ExecResult { + fabro_sandbox::ExecResult { + stdout: "out".into(), + stderr: "err".into(), + exit_code: Some(exit_code), + termination: fabro_types::CommandTermination::Exited, + duration_ms: 12, + } + } + + fn cancelled() -> fabro_sandbox::ExecResult { + fabro_sandbox::ExecResult { + stdout: "out".into(), + stderr: String::new(), + exit_code: None, + termination: fabro_types::CommandTermination::Cancelled, + duration_ms: 12, + } + } + + async fn run_shell_tool( + exec_result: fabro_sandbox::ExecResult, + hooks: Option<&Arc>, + emitter: &Emitter, + ) -> ToolResult { + let mut registry = ToolRegistry::new(); + registry.register(make_shell_tool()); + let tc = make_tool_call( + "shell", + "call_1", + serde_json::json!({"command": "make test"}), + ); + + execute_and_emit_one_tool( + &tc, + ®istry, + shell_sandbox(exec_result), + hooks, + CancellationToken::new(), + &SessionOptions::default(), + emitter, + "test-session", + "test-session", + None, + ) + .await + } + + fn drain(receiver: &mut broadcast::Receiver) -> Vec { + let mut events = Vec::new(); + while let Ok(event) = receiver.try_recv() { + events.push(event); + } + events + } + + #[tokio::test] + async fn shell_nonzero_exit_becomes_an_error_tool_result() { + let emitter = Emitter::new(); + let result = run_shell_tool(exited(7), None, &emitter).await; + + assert!(result.is_error); + assert!( + result.content.as_str().unwrap().contains("Exit code: 7"), + "got: {}", + result.content + ); + } + + #[tokio::test] + async fn shell_exit_zero_remains_a_successful_tool_result() { + let emitter = Emitter::new(); + let result = run_shell_tool(exited(0), None, &emitter).await; + + assert!(!result.is_error); + } + + #[tokio::test] + async fn shell_failure_emits_started_then_process_then_completed() { + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); + run_shell_tool(exited(7), None, &emitter).await; + + let events = drain(&mut receiver); + let names: Vec<&str> = events + .iter() + .filter_map(|event| match &event.event { + AgentEvent::ToolCallStarted { .. } => Some("started"), + AgentEvent::ToolProcessCompleted { .. } => Some("process"), + AgentEvent::ToolCallCompleted { .. } => Some("completed"), + _ => None, + }) + .collect(); + assert_eq!(names, vec!["started", "process", "completed"]); + + for event in &events { + assert_eq!(event.session_id, "test-session"); + } + let process = events + .iter() + .find(|event| matches!(event.event, AgentEvent::ToolProcessCompleted { .. })) + .expect("process event"); + assert_eq!(process.tool_call_id.as_deref(), Some("call_1")); + match &process.event { + AgentEvent::ToolProcessCompleted { + exit_code, + termination, + .. + } => { + assert_eq!(*exit_code, Some(7)); + assert_eq!(*termination, fabro_types::CommandTermination::Exited); + } + other => panic!("expected a process event, got {other:?}"), + } + + let completed = events + .iter() + .find_map(|event| match &event.event { + AgentEvent::ToolCallCompleted { + tool_call_id, + is_error, + .. + } => Some((tool_call_id.clone(), *is_error)), + _ => None, + }) + .expect("tool completed event"); + assert_eq!(completed, ("call_1".to_string(), true)); + } + + #[tokio::test] + async fn shell_failure_runs_only_the_failure_hook() { + for exec_result in [exited(7), cancelled()] { + let mock = Arc::new(MockHookCallback::new(ToolHookDecision::Proceed)); + let hooks: Arc = mock.clone(); + run_shell_tool(exec_result, Some(&hooks), &Emitter::new()).await; + + assert_eq!(mock.post_failure_calls.lock().unwrap().len(), 1); + assert!(mock.post_calls.lock().unwrap().is_empty()); + } + } + + #[tokio::test] + async fn shell_success_runs_only_the_success_hook() { + let mock = Arc::new(MockHookCallback::new(ToolHookDecision::Proceed)); + let hooks: Arc = mock.clone(); + run_shell_tool(exited(0), Some(&hooks), &Emitter::new()).await; + + assert_eq!(mock.post_calls.lock().unwrap().len(), 1); + assert!(mock.post_failure_calls.lock().unwrap().is_empty()); + } + + #[test] + fn truncation_preserves_tool_call_id_and_error_state() { + let result = ToolResult::error("call_1", "x".repeat(60_000)); + + let truncated = truncate_tool_result(&result, "shell", &SessionOptions::default()); + + assert_eq!(truncated.tool_call_id, "call_1"); + assert!(truncated.is_error); + assert!(truncated.content.as_str().unwrap().len() < 60_000); + } } diff --git a/lib/components/fabro-agent/src/tools.rs b/lib/components/fabro-agent/src/tools.rs index 5eca7b5be..4043006fa 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -10,8 +10,9 @@ use fabro_static::EnvVars; use futures::{StreamExt, stream}; use crate::config::NativeToolOptions; -use crate::sandbox::GrepOptions; +use crate::sandbox::{CommandOutputCallback, ExecStreamingResult, GrepOptions}; use crate::tool_registry::{RegisteredTool, ToolRegistry, ToolSource}; +use crate::types::AgentEvent; const MAX_WEB_FETCH_BYTES: usize = 100 * 1024; const MAX_READ_MANY_FILES_CONCURRENCY: usize = 8; @@ -239,51 +240,92 @@ pub fn make_shell_tool_with_options(options: &NativeToolOptions) -> RegisteredTo executor: Arc::new(move |args, ctx| { Box::pin(async move { let command = required_str(&args, "command")?; - let command = format!("exec 2>&1\n{command}"); let timeout_ms = args .get("timeout_ms") .and_then(serde_json::Value::as_u64) .unwrap_or(default_timeout) .min(max_timeout); - let tool_env = ctx.resolve_tool_env().await.map_err(|e| format!("{e:#}"))?; + let tool_env = ctx + .resolve_tool_env() + .await + .map_err(|e| format!("{SHELL_NO_PROCESS_RESULT}: {e:#}"))?; tracing::debug!( env_var_count = tool_env.as_ref().map_or(0, std::collections::HashMap::len), "Injecting sandbox env vars into tool execution" ); - let result = ctx + let streaming = ctx .env - .exec_command( - &command, - timeout_ms, + .exec_command_streaming( + command, + Some(timeout_ms), None, tool_env.as_ref(), - Some(ctx.cancel), + Some(ctx.cancel.clone()), + discard_output_callback(), ) .await - .map_err(|e| e.display_with_causes())?; + .map_err(|e| { + format!("{SHELL_NO_PROCESS_RESULT}: {}", e.display_with_causes()) + })?; - let mut output = String::new(); - if result.is_timed_out() { - output.push_str("Command timed out.\n"); - } else if result.is_cancelled() { - output.push_str("Command cancelled.\n"); + let text = render_shell_result(&streaming); + ctx.emit_agent_event(AgentEvent::ToolProcessCompleted { + exit_code: streaming.result.exit_code, + termination: streaming.result.termination, + duration_ms: streaming.result.duration_ms, + streams_separated: streaming.streams_separated, + exec_output_tail: streaming.result.default_redacted_output_tail(), + }); + + if streaming.result.is_success() { + Ok(text) + } else { + Err(text) } - let _ = write!( - output, - "Exit code: {}\noutput:\n{}", - result - .exit_code - .map_or_else(|| "none".to_string(), |code| code.to_string()), - result.stdout - ); - Ok(output) }) }), source: ToolSource::Native, } } +/// Prefix for shell failures that never produced an `ExecResult`, so the model +/// can tell "the process never ran" apart from "the process ran and failed". +const SHELL_NO_PROCESS_RESULT: &str = "Shell command produced no process result"; + +/// The agent shell tool consumes the streaming exec path for its stream +/// provenance and partial-output capture, but does not forward live output +/// deltas onto the agent protocol. +fn discard_output_callback() -> CommandOutputCallback { + Arc::new(|_stream, _bytes| Box::pin(async { Ok(()) })) +} + +/// Renders the model-facing shell result: termination, exit code, duration, +/// and provider-honest output sections. Metadata stays at the head and +/// `stderr` at the tail so head/tail truncation preserves both. +fn render_shell_result(streaming: &ExecStreamingResult) -> String { + let result = &streaming.result; + let mut output = format!( + "Termination: {}\nExit code: {}\nDuration: {}ms\n", + result.termination.as_str(), + result + .exit_code + .map_or_else(|| "none".to_string(), |code| code.to_string()), + result.duration_ms, + ); + if streaming.streams_separated { + if !result.stdout.is_empty() { + let _ = write!(output, "stdout:\n{}\n", result.stdout); + } + if !result.stderr.is_empty() { + let _ = write!(output, "stderr:\n{}\n", result.stderr); + } + } else if !result.stdout.is_empty() { + let _ = write!(output, "output (combined):\n{}\n", result.stdout); + } + output +} + #[must_use] pub fn make_grep_tool() -> RegisteredTool { RegisteredTool { @@ -693,10 +735,12 @@ mod tests { use tokio_util::sync::CancellationToken; use super::*; - use crate::config::{NativeToolOptions, ToolSecrets}; + use crate::config::{NativeToolOptions, SessionOptions, ToolSecrets}; + use crate::local_sandbox::LocalSandbox; use crate::sandbox::*; use crate::test_support::MockSandbox; - use crate::tool_registry::ToolContext; + use crate::tool_registry::{AgentEventEmitter, ToolContext}; + use crate::truncation; #[test] fn core_tool_descriptions_include_actionable_guidance() { @@ -981,20 +1025,35 @@ mod tests { assert_eq!(written[0].1, "1 | keep this literal\ngoodbye"); } - #[tokio::test] - async fn shell_basic_command() { - let tool = make_shell_tool(); - let env: Arc = Arc::new(MockSandbox { - exec_result: ExecResult { - stdout: "hello".into(), - stderr: String::new(), - exit_code: Some(0), - termination: CommandTermination::Exited, - duration_ms: 10, - }, - ..Default::default() - }); - let result = (tool.executor)(serde_json::json!({"command": "echo hello"}), ToolContext { + /// Records the typed agent events a tool emits through its bound emitter. + #[derive(Default)] + struct RecordingAgentEmitter { + events: std::sync::Mutex>, + } + + impl RecordingAgentEmitter { + fn events(&self) -> Vec { + self.events.lock().expect("events lock poisoned").clone() + } + + fn only_process_event(&self) -> AgentEvent { + let events = self.events(); + assert_eq!(events.len(), 1, "expected one agent event, got {events:?}"); + events.into_iter().next().expect("one event") + } + } + + impl AgentEventEmitter for RecordingAgentEmitter { + fn emit(&self, event: AgentEvent) { + self.events + .lock() + .expect("events lock poisoned") + .push(event); + } + } + + fn shell_context(env: Arc) -> ToolContext { + ToolContext { env, cancel: CancellationToken::new(), tool_env_provider: None, @@ -1002,11 +1061,72 @@ mod tests { root_session_id: None, tool_call_id: None, agent_event_emitter: None, + } + } + + fn shell_context_with_emitter( + env: Arc, + emitter: Arc, + ) -> ToolContext { + ToolContext { + agent_event_emitter: Some(emitter), + ..shell_context(env) + } + } + + fn mock_sandbox_with(result: ExecResult) -> Arc { + Arc::new(MockSandbox { + exec_result: result, + ..Default::default() }) + } + + #[tokio::test] + async fn shell_success_returns_ok_with_metadata_and_separate_streams() { + let tool = make_shell_tool(); + let env: Arc = mock_sandbox_with(ExecResult { + stdout: "hello".into(), + stderr: "a warning".into(), + exit_code: Some(0), + termination: CommandTermination::Exited, + duration_ms: 10, + }); + let output = (tool.executor)( + serde_json::json!({"command": "echo hello"}), + shell_context(env), + ) + .await + .expect("exit 0 is a successful tool result"); + + assert_eq!( + output, + "Termination: exited\nExit code: 0\nDuration: 10ms\nstdout:\nhello\nstderr:\na \ + warning\n" + ); + } + + #[tokio::test] + async fn shell_forwards_command_without_stream_redirection_wrapper() { + let tool = make_shell_tool(); + let env = mock_sandbox_with(ExecResult { + stdout: String::new(), + stderr: String::new(), + exit_code: Some(0), + termination: CommandTermination::Exited, + duration_ms: 1, + }); + let _ = (tool.executor)( + serde_json::json!({"command": "make test"}), + shell_context(env.clone()), + ) .await; - let output = result.unwrap(); - assert!(output.contains("Exit code: 0")); - assert!(output.contains("hello")); + + let captured = env + .captured_command + .lock() + .expect("captured_command lock poisoned") + .clone(); + assert_eq!(captured.as_deref(), Some("make test")); } #[tokio::test] @@ -1043,46 +1163,253 @@ mod tests { }, ..Default::default() }); - let result = (tool.executor)(serde_json::json!({"command": "false"}), ToolContext { - env, - cancel: CancellationToken::new(), - tool_env_provider: None, - session_id: None, - root_session_id: None, - tool_call_id: None, - agent_event_emitter: None, - }) - .await; - let output = result.unwrap(); - assert!(output.contains("Exit code: 1")); - assert!(output.contains("error")); + let output = (tool.executor)(serde_json::json!({"command": "false"}), shell_context(env)) + .await + .expect_err("a nonzero exit is a failed tool result"); + assert!(output.contains("Termination: exited"), "got: {output}"); + assert!(output.contains("Exit code: 1"), "got: {output}"); + assert!(output.contains("stdout:\nerror"), "got: {output}"); + assert!(!output.contains("stderr:"), "got: {output}"); } #[tokio::test] - async fn shell_timeout_output() { + async fn shell_timeout_returns_error_with_partial_output() { + let tool = make_shell_tool(); + let env: Arc = mock_sandbox_with(ExecResult { + stdout: "partial".into(), + stderr: String::new(), + exit_code: None, + termination: CommandTermination::TimedOut, + duration_ms: 10000, + }); + let output = (tool.executor)( + serde_json::json!({"command": "sleep 100"}), + shell_context(env), + ) + .await + .expect_err("a timeout is a failed tool result"); + + assert!(output.contains("Termination: timed_out"), "got: {output}"); + assert!(output.contains("Exit code: none"), "got: {output}"); + assert!(output.contains("stdout:\npartial"), "got: {output}"); + } + + #[tokio::test] + async fn shell_cancellation_returns_error_with_partial_output() { + let tool = make_shell_tool(); + let env: Arc = mock_sandbox_with(ExecResult { + stdout: "partial".into(), + stderr: String::new(), + exit_code: None, + termination: CommandTermination::Cancelled, + duration_ms: 42, + }); + let output = (tool.executor)( + serde_json::json!({"command": "sleep 100"}), + shell_context(env), + ) + .await + .expect_err("a cancellation is a failed tool result"); + + assert!(output.contains("Termination: cancelled"), "got: {output}"); + assert!(output.contains("Exit code: none"), "got: {output}"); + assert!(output.contains("stdout:\npartial"), "got: {output}"); + } + + #[tokio::test] + async fn shell_sandbox_failure_returns_error_without_a_process_outcome() { + let tool = make_shell_tool(); + let env: Arc = Arc::new(MockSandbox { + exec_error: Some("sandbox transport is down".into()), + ..Default::default() + }); + let emitter = Arc::new(RecordingAgentEmitter::default()); + + let output = (tool.executor)( + serde_json::json!({"command": "make test"}), + shell_context_with_emitter(env, emitter.clone()), + ) + .await + .expect_err("a sandbox transport failure is a failed tool result"); + + assert!( + output.contains("Shell command produced no process result"), + "got: {output}" + ); + assert!( + output.contains("sandbox transport is down"), + "got: {output}" + ); + assert!(!output.contains("Exit code"), "got: {output}"); + assert!(emitter.events().is_empty(), "got: {:?}", emitter.events()); + } + + #[tokio::test] + async fn shell_emits_process_event_with_typed_outcome_and_redacted_tails() { + let tool = make_shell_tool(); + let env: Arc = mock_sandbox_with(ExecResult { + stdout: "out".into(), + stderr: "boom key=AKIAYRWQG5EJLPZLBYNP".into(), + exit_code: Some(7), + termination: CommandTermination::Exited, + duration_ms: 12, + }); + let emitter = Arc::new(RecordingAgentEmitter::default()); + + let _ = (tool.executor)( + serde_json::json!({"command": "printf out; printf err >&2; exit 7"}), + shell_context_with_emitter(env, emitter.clone()), + ) + .await; + + match emitter.only_process_event() { + AgentEvent::ToolProcessCompleted { + exit_code, + termination, + duration_ms, + streams_separated, + exec_output_tail, + } => { + assert_eq!(exit_code, Some(7)); + assert_eq!(termination, CommandTermination::Exited); + assert_eq!(duration_ms, 12); + assert!(streams_separated); + let tail = exec_output_tail.expect("output tail"); + assert_eq!(tail.stdout.as_deref(), Some("out")); + let stderr = tail.stderr.expect("stderr tail"); + assert!(stderr.contains("boom"), "got: {stderr}"); + assert!(!stderr.contains("AKIAYRWQG5EJLPZLBYNP"), "got: {stderr}"); + } + other => panic!("expected a process event, got {other:?}"), + } + } + + #[tokio::test] + async fn shell_renders_combined_output_when_streams_are_not_separated() { let tool = make_shell_tool(); let env: Arc = Arc::new(MockSandbox { exec_result: ExecResult { - stdout: String::new(), + stdout: "interleaved".into(), stderr: String::new(), - exit_code: None, - termination: CommandTermination::TimedOut, - duration_ms: 10000, + exit_code: Some(0), + termination: CommandTermination::Exited, + duration_ms: 5, }, + streams_separated: false, ..Default::default() }); - let result = (tool.executor)(serde_json::json!({"command": "sleep 100"}), ToolContext { - env, - cancel: CancellationToken::new(), - tool_env_provider: None, - session_id: None, - root_session_id: None, - tool_call_id: None, - agent_event_emitter: None, - }) - .await; - let output = result.unwrap(); - assert!(output.starts_with("Command timed out.\n")); + let emitter = Arc::new(RecordingAgentEmitter::default()); + + let output = (tool.executor)( + serde_json::json!({"command": "echo interleaved"}), + shell_context_with_emitter(env, emitter.clone()), + ) + .await + .expect("exit 0 is a successful tool result"); + + assert!( + output.contains("output (combined):\ninterleaved"), + "got: {output}" + ); + assert!(!output.contains("stderr:"), "got: {output}"); + match emitter.only_process_event() { + AgentEvent::ToolProcessCompleted { + streams_separated, .. + } => assert!(!streams_separated), + other => panic!("expected a process event, got {other:?}"), + } + } + + #[tokio::test] + async fn shell_truncation_preserves_exit_metadata_and_stderr_tail() { + let tool = make_shell_tool(); + let stdout = (0..400) + .map(|line| format!("{line}: {}", "x".repeat(100))) + .collect::>() + .join("\n"); + assert!(stdout.len() > 30_000); + let env: Arc = mock_sandbox_with(ExecResult { + stdout, + stderr: "the build failed".into(), + exit_code: Some(2), + termination: CommandTermination::Exited, + duration_ms: 900, + }); + + let output = (tool.executor)( + serde_json::json!({"command": "make build"}), + shell_context(env), + ) + .await + .expect_err("a nonzero exit is a failed tool result"); + let truncated = + truncation::truncate_tool_output(&output, "shell", &SessionOptions::default()); + + assert!(truncated.len() < output.len()); + assert!(truncated.starts_with("Termination: exited\nExit code: 2\n")); + assert!( + truncated.contains("stderr:\nthe build failed"), + "stderr tail did not survive truncation" + ); + } + + /// End-to-end against a real process: the local provider separates the + /// streams and reports the real exit code, and none of it is laundered + /// into a successful tool result. + #[tokio::test] + async fn shell_reports_real_local_process_outcome() { + let tool = make_shell_tool(); + let env: Arc = Arc::new(LocalSandbox::new( + std::env::current_dir().expect("current dir"), + )); + let emitter = Arc::new(RecordingAgentEmitter::default()); + + let output = (tool.executor)( + serde_json::json!({"command": "printf 'out'; printf 'err' >&2; exit 7"}), + shell_context_with_emitter(env, emitter.clone()), + ) + .await + .expect_err("exit 7 is a failed tool result"); + + assert!(output.contains("Termination: exited"), "got: {output}"); + assert!(output.contains("Exit code: 7"), "got: {output}"); + assert!(output.contains("stdout:\nout"), "got: {output}"); + assert!(output.contains("stderr:\nerr"), "got: {output}"); + + match emitter.only_process_event() { + AgentEvent::ToolProcessCompleted { + exit_code, + termination, + streams_separated, + exec_output_tail, + .. + } => { + assert_eq!(exit_code, Some(7)); + assert_eq!(termination, CommandTermination::Exited); + assert!(streams_separated); + let tail = exec_output_tail.expect("output tail"); + assert_eq!(tail.stdout.as_deref(), Some("out")); + assert_eq!(tail.stderr.as_deref(), Some("err")); + } + other => panic!("expected a process event, got {other:?}"), + } + } + + #[tokio::test] + async fn shell_public_schema_is_command_timeout_and_description() { + let tool = make_shell_tool(); + assert_eq!( + tool.definition.parameters, + serde_json::json!({ + "type": "object", + "properties": { + "command": {"type": "string", "description": "The shell command to execute"}, + "timeout_ms": {"type": "integer", "description": "Timeout in milliseconds"}, + "description": {"type": "string", "description": "Description of what this command does"} + }, + "required": ["command"] + }) + ); } #[tokio::test] diff --git a/lib/components/fabro-agent/src/types.rs b/lib/components/fabro-agent/src/types.rs index fc3581c57..129cfc8a0 100644 --- a/lib/components/fabro-agent/src/types.rs +++ b/lib/components/fabro-agent/src/types.rs @@ -4,7 +4,10 @@ use chrono::{DateTime, Utc}; use fabro_llm::Error as LlmError; use fabro_llm::types::{ContentPart, ThinkingData, TokenCounts, ToolCall, ToolResult}; use fabro_model::{CostSource, ModelRef}; -use fabro_types::{ReasoningOutput, SessionMessage, StageContextWindowProjection}; +use fabro_types::{ + CommandTermination, ExecOutputTail, ReasoningOutput, SessionMessage, + StageContextWindowProjection, +}; use serde::de::DeserializeOwned; use serde::{Deserialize, Serialize}; @@ -280,6 +283,19 @@ pub enum AgentEvent { output: serde_json::Value, is_error: bool, }, + /// Subordinate process outcome for a tool call that ran a command. + /// Emitted before the owning `ToolCallCompleted`, which stays the single + /// tool-protocol completion and the authoritative owner of `is_error`. + /// Session and tool-call identity come from the emitting envelope. + ToolProcessCompleted { + #[serde(default, skip_serializing_if = "Option::is_none")] + exit_code: Option, + termination: CommandTermination, + duration_ms: u64, + streams_separated: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + exec_output_tail: Option, + }, Error { error: Error, }, @@ -457,6 +473,28 @@ impl AgentEvent { "Tool call completed" ); } + Self::ToolProcessCompleted { + exit_code, + termination, + duration_ms, + streams_separated, + exec_output_tail, + } => { + let tail = ExecOutputTail::trace_summary(exec_output_tail.as_ref()); + debug!( + session_id, + exit_code = ?exit_code, + termination = termination.as_str(), + duration_ms, + streams_separated, + output_tail_present = tail.present, + stdout_bytes = tail.stdout_bytes, + stderr_bytes = tail.stderr_bytes, + stdout_truncated = tail.stdout_truncated, + stderr_truncated = tail.stderr_truncated, + "Tool process completed" + ); + } Self::Error { error } => { error!(session_id, error = %error, "Agent error"); } diff --git a/lib/components/fabro-agent/tests/it/docker_shell.rs b/lib/components/fabro-agent/tests/it/docker_shell.rs new file mode 100644 index 000000000..52d1dd081 --- /dev/null +++ b/lib/components/fabro-agent/tests/it/docker_shell.rs @@ -0,0 +1,97 @@ +//! Proves the agent shell tool reports real process outcomes through the +//! Docker provider's streaming path, which uses a `bash -lc` supervisor and +//! separate stdout/stderr channels. + +use std::sync::{Arc, Mutex}; + +use fabro_agent::sandbox::Sandbox; +use fabro_agent::tool_registry::{AgentEventEmitter, ToolContext}; +use fabro_agent::tools::make_shell_tool; +use fabro_agent::types::AgentEvent; +use fabro_agent::{DockerSandbox, DockerSandboxOptions}; +use fabro_types::CommandTermination; +use tokio_util::sync::CancellationToken; + +#[derive(Default)] +struct RecordingAgentEmitter { + events: Mutex>, +} + +impl AgentEventEmitter for RecordingAgentEmitter { + fn emit(&self, event: AgentEvent) { + self.events + .lock() + .expect("events lock poisoned") + .push(event); + } +} + +#[tokio::test] +#[ignore = "requires real Docker container lifecycle; run explicitly when changing shell tool exec integration"] +async fn shell_reports_real_docker_process_outcome() { + let Ok(sandbox) = DockerSandbox::new( + DockerSandboxOptions { + image: "buildpack-deps:noble".to_string(), + auto_pull: false, + skip_clone: true, + ..DockerSandboxOptions::default() + }, + None, + None, + None, + None, + ) else { + return; + }; + // No Docker daemon or no local image: nothing to prove here. + if sandbox.initialize().await.is_err() { + return; + } + + let sandbox = Arc::new(sandbox); + let emitter = Arc::new(RecordingAgentEmitter::default()); + let tool = make_shell_tool(); + let result = (tool.executor)( + serde_json::json!({"command": "printf 'out'; printf 'err' >&2; exit 7"}), + ToolContext { + env: sandbox.clone() as Arc, + cancel: CancellationToken::new(), + tool_env_provider: None, + session_id: None, + root_session_id: None, + tool_call_id: None, + agent_event_emitter: Some(emitter.clone()), + }, + ) + .await; + sandbox + .cleanup() + .await + .expect("docker cleanup should succeed"); + + let output = result.expect_err("exit 7 is a failed tool result"); + assert!(output.contains("Termination: exited"), "got: {output}"); + assert!(output.contains("Exit code: 7"), "got: {output}"); + assert!(output.contains("stdout:\nout"), "got: {output}"); + assert!(output.contains("stderr:\nerr"), "got: {output}"); + + let events = emitter.events.lock().expect("events lock poisoned").clone(); + assert_eq!(events.len(), 1, "expected one agent event, got {events:?}"); + match &events[0] { + AgentEvent::ToolProcessCompleted { + exit_code, + termination, + streams_separated, + exec_output_tail, + .. + } => { + assert_eq!(*exit_code, Some(7)); + assert_eq!(*termination, CommandTermination::Exited); + assert!(streams_separated); + let tail = exec_output_tail.as_ref().expect("output tail"); + assert_eq!(tail.stdout.as_deref(), Some("out")); + assert_eq!(tail.stderr.as_deref(), Some("err")); + } + other => panic!("expected a process event, got {other:?}"), + } +} diff --git a/lib/components/fabro-agent/tests/it/main.rs b/lib/components/fabro-agent/tests/it/main.rs index efb7dbdce..485fe7037 100644 --- a/lib/components/fabro-agent/tests/it/main.rs +++ b/lib/components/fabro-agent/tests/it/main.rs @@ -1,3 +1,5 @@ mod compaction; +#[cfg(feature = "docker")] +mod docker_shell; mod guardrails; mod parity_matrix; diff --git a/lib/components/fabro-sandbox/src/test_support.rs b/lib/components/fabro-sandbox/src/test_support.rs index 1593c9032..d48f552e2 100644 --- a/lib/components/fabro-sandbox/src/test_support.rs +++ b/lib/components/fabro-sandbox/src/test_support.rs @@ -44,6 +44,12 @@ pub struct MockSandbox { pub event_callback: Option, pub stdio_process_error: Option, pub stdio_process: Mutex>, + /// Fails `exec_command` and `exec_command_streaming` before any process + /// runs, so callers see a transport error rather than an `ExecResult`. + pub exec_error: Option, + /// Reported by `exec_command_streaming`. Set to `false` to model a + /// provider that cannot separate stdout from stderr. + pub streams_separated: bool, } impl MockSandbox { @@ -116,6 +122,8 @@ impl Default for MockSandbox { event_callback: None, stdio_process_error: None, stdio_process: Mutex::new(None), + exec_error: None, + streams_separated: true, } } } @@ -233,7 +241,49 @@ impl Sandbox for MockSandbox { .captured_env_vars .lock() .expect("captured_env_vars lock poisoned") = env_vars.cloned(); - Ok(self.exec_result.clone()) + match &self.exec_error { + Some(error) => Err(crate::Error::message(error.clone())), + None => Ok(self.exec_result.clone()), + } + } + + async fn exec_command_streaming( + &self, + command: &str, + timeout_ms: Option, + working_dir: Option<&str>, + env_vars: Option<&std::collections::HashMap>, + cancel_token: Option, + output_callback: crate::CommandOutputCallback, + ) -> crate::Result { + let result = self + .exec_command( + command, + timeout_ms.unwrap_or(u64::MAX), + working_dir, + env_vars, + cancel_token, + ) + .await?; + if !result.stdout.is_empty() { + output_callback( + fabro_types::CommandOutputStream::Stdout, + result.stdout.as_bytes().to_vec(), + ) + .await?; + } + if !result.stderr.is_empty() { + output_callback( + fabro_types::CommandOutputStream::Stderr, + result.stderr.as_bytes().to_vec(), + ) + .await?; + } + Ok(crate::ExecStreamingResult { + result, + streams_separated: self.streams_separated, + live_streaming: false, + }) } async fn spawn_stdio_process( diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index d86c17424..1731cd265 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -655,6 +655,22 @@ fn event_body_from_event(event: &Event) -> EventBody { tool_result: None, turn_id: None, }), + AgentEvent::ToolProcessCompleted { + exit_code, + termination, + duration_ms, + streams_separated, + exec_output_tail, + } => EventBody::AgentToolProcessCompleted( + fabro_types::AgentToolProcessCompletedProps { + exit_code: *exit_code, + termination: *termination, + duration_ms: *duration_ms, + streams_separated: *streams_separated, + exec_output_tail: exec_output_tail.clone(), + visit: *visit, + }, + ), AgentEvent::Error { error } => EventBody::AgentError(fabro_types::AgentErrorProps { error: serde_json::to_value(error).expect("agent Error derives Serialize with no custom logic that can fail"), visit: *visit, @@ -1517,6 +1533,52 @@ mod tests { assert_eq!(properties["visit"], 2); } + #[test] + fn run_event_agent_tool_process_completed_carries_stage_session_and_actor() { + let stored = to_run_event(&fixtures::RUN_4, &Event::Agent { + stage: "code".to_string(), + visit: 2, + event: AgentEvent::ToolProcessCompleted { + exit_code: Some(7), + termination: ::fabro_types::CommandTermination::Exited, + duration_ms: 12, + streams_separated: true, + exec_output_tail: Some(::fabro_types::ExecOutputTail { + stdout: Some("out".to_string()), + stderr: Some("err".to_string()), + stdout_truncated: false, + stderr_truncated: false, + }), + }, + session_id: Some("ses_child".to_string()), + parent_session_id: Some("ses_parent".to_string()), + tool_call_id: Some("call_1".to_string()), + }); + + assert_eq!(stored.event_name(), "agent.tool.process.completed"); + assert_eq!(stored.node_id.as_deref(), Some("code")); + assert_eq!(stored.stage_id, Some(StageId::new("code", 2))); + assert_eq!(stored.session_id.as_deref(), Some("ses_child")); + assert_eq!(stored.parent_session_id.as_deref(), Some("ses_parent")); + assert_eq!(stored.tool_call_id.as_deref(), Some("call_1")); + assert_eq!( + stored.actor, + Some(::fabro_types::Principal::Agent { + session_id: Some("ses_child".to_string()), + parent_session_id: Some("ses_parent".to_string()), + model: None, + }) + ); + + let properties = stored.properties().unwrap(); + assert_eq!(properties["exit_code"], 7); + assert_eq!(properties["termination"], "exited"); + assert_eq!(properties["duration_ms"], 12); + assert_eq!(properties["streams_separated"], true); + assert_eq!(properties["exec_output_tail"]["stdout"], "out"); + assert_eq!(properties["visit"], 2); + } + #[test] fn run_event_agent_tools_available_moves_session_and_stage_metadata_to_header() { let stored = to_run_event(&fixtures::RUN_4, &Event::AgentToolsAvailable { diff --git a/lib/components/fabro-workflow/src/event/names.rs b/lib/components/fabro-workflow/src/event/names.rs index 67be1ab3a..cd04b08b1 100644 --- a/lib/components/fabro-workflow/src/event/names.rs +++ b/lib/components/fabro-workflow/src/event/names.rs @@ -75,6 +75,7 @@ pub fn event_name(event: &Event) -> &'static str { AgentEvent::ToolCallStarted { .. } => "agent.tool.started", AgentEvent::ToolCallOutputDelta { .. } => "agent.tool.output.delta", AgentEvent::ToolCallCompleted { .. } => "agent.tool.completed", + AgentEvent::ToolProcessCompleted { .. } => "agent.tool.process.completed", AgentEvent::Error { .. } => "agent.error", AgentEvent::Warning { .. } => "agent.warning", AgentEvent::LoopDetected => "agent.loop.detected", diff --git a/lib/components/fabro-workflow/src/event/redaction.rs b/lib/components/fabro-workflow/src/event/redaction.rs index 16aabccab..49e86a2ab 100644 --- a/lib/components/fabro-workflow/src/event/redaction.rs +++ b/lib/components/fabro-workflow/src/event/redaction.rs @@ -76,6 +76,41 @@ mod tests { ); } + #[test] + fn build_redacted_event_payload_redacts_tool_process_output_tails() { + let secret = "sk-ant-api03-xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA"; + let stored = to_run_event(&fixtures::RUN_8, &Event::Agent { + stage: "code".to_string(), + visit: 1, + event: AgentEvent::ToolProcessCompleted { + exit_code: Some(7), + termination: ::fabro_types::CommandTermination::Exited, + duration_ms: 12, + streams_separated: true, + exec_output_tail: Some(fabro_types::ExecOutputTail { + stdout: Some(format!("stdout {secret}")), + stderr: Some("plain stderr".to_string()), + stdout_truncated: false, + stderr_truncated: false, + }), + }, + session_id: Some("ses_child".to_string()), + parent_session_id: None, + tool_call_id: Some("call_1".to_string()), + }); + + let payload = build_redacted_event_payload(&stored, &fixtures::RUN_8).unwrap(); + let payload_text = serde_json::to_string(payload.as_value()).unwrap(); + + assert!(!payload_text.contains(secret)); + assert!(payload_text.contains("REDACTED")); + assert_eq!(payload.as_value()["event"], "agent.tool.process.completed"); + assert_eq!( + payload.as_value()["properties"]["exec_output_tail"]["stderr"], + "plain stderr" + ); + } + /// Reasoning is model-authored text like any other, so it goes through /// the same canonical redaction pass as assistant output. #[test] diff --git a/lib/components/fabro-workflow/src/event/stored_fields.rs b/lib/components/fabro-workflow/src/event/stored_fields.rs index 285bd9c89..66129d757 100644 --- a/lib/components/fabro-workflow/src/event/stored_fields.rs +++ b/lib/components/fabro-workflow/src/event/stored_fields.rs @@ -328,7 +328,8 @@ fn agent_actor_for_event( }), AgentEvent::ToolCallStarted { .. } | AgentEvent::ToolCallOutputDelta { .. } - | AgentEvent::ToolCallCompleted { .. } => Some(Principal::Agent { + | AgentEvent::ToolCallCompleted { .. } + | AgentEvent::ToolProcessCompleted { .. } => Some(Principal::Agent { session_id: session_id.map(str::to_string), parent_session_id: parent_session_id.map(str::to_string), model: None, diff --git a/lib/foundation/fabro-types/src/run_event/agent.rs b/lib/foundation/fabro-types/src/run_event/agent.rs index 9cd852407..604843011 100644 --- a/lib/foundation/fabro-types/src/run_event/agent.rs +++ b/lib/foundation/fabro-types/src/run_event/agent.rs @@ -3,11 +3,11 @@ use serde::{Deserialize, Serialize}; use serde_json::Value; use strum::{Display, EnumString, IntoStaticStr}; -use super::BilledTokenCounts; +use super::{BilledTokenCounts, ExecOutputTail}; use crate::transcript::{ToolCall, ToolResult, TranscriptMessage}; use crate::{ - MessageId, ModelRef, PairId, PairMessageId, PairSystemMessageKind, PermissionLevel, - ReasoningOutput, StageContextWindowProjection, TurnId, + CommandTermination, MessageId, ModelRef, PairId, PairMessageId, PairSystemMessageKind, + PermissionLevel, ReasoningOutput, StageContextWindowProjection, TurnId, }; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -181,6 +181,25 @@ pub struct AgentToolCompletedProps { pub turn_id: Option, } +/// Subordinate diagnostic for a tool call that ran a process: the real +/// termination, exit code, duration, and bounded redacted output tails. +/// +/// This never replaces `agent.tool.completed`, which remains the single +/// tool-protocol completion and the authoritative owner of `is_error`. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct AgentToolProcessCompletedProps { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub exit_code: Option, + pub termination: CommandTermination, + pub duration_ms: u64, + /// `false` when the provider could not separate stdout from stderr. The + /// combined output is then carried in `exec_output_tail.stdout`. + pub streams_separated: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub exec_output_tail: Option, + pub visit: u32, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct AgentErrorProps { pub error: Value, diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index 7c8d1621a..4b937f682 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -208,6 +208,8 @@ pub enum EventBody { AgentToolStarted(AgentToolStartedProps), #[serde(rename = "agent.tool.completed")] AgentToolCompleted(AgentToolCompletedProps), + #[serde(rename = "agent.tool.process.completed")] + AgentToolProcessCompleted(AgentToolProcessCompletedProps), #[serde(rename = "agent.error")] AgentError(AgentErrorProps), #[serde(rename = "agent.warning")] @@ -487,6 +489,7 @@ impl EventBody { Self::AgentMessage(_) => "agent.message", Self::AgentToolStarted(_) => "agent.tool.started", Self::AgentToolCompleted(_) => "agent.tool.completed", + Self::AgentToolProcessCompleted(_) => "agent.tool.process.completed", Self::AgentError(_) => "agent.error", Self::AgentWarning(_) => "agent.warning", Self::AgentLoopDetected(_) => "agent.loop.detected", @@ -656,6 +659,7 @@ fn is_known_event_name(event: &str) -> bool { | "agent.message" | "agent.tool.started" | "agent.tool.completed" + | "agent.tool.process.completed" | "agent.error" | "agent.warning" | "agent.loop.detected" @@ -920,8 +924,8 @@ mod tests { use super::*; use crate::{ - AuthMethod, Edge, Graph, IdpIdentity, Node, PendingReason, RunBlobId, WorkflowSettings, - fixtures, test_support, + AuthMethod, CommandTermination, Edge, Graph, IdpIdentity, Node, PendingReason, RunBlobId, + WorkflowSettings, fixtures, test_support, }; fn user_principal(login: &str) -> Principal { @@ -2409,6 +2413,90 @@ mod tests { assert_eq!(parsed, body); } + #[test] + fn agent_tool_process_completed_round_trips_as_a_known_typed_event() { + let value = json!({ + "id": "evt_process", + "ts": "2026-04-08T16:21:11.106Z", + "run_id": fixtures::RUN_1, + "event": "agent.tool.process.completed", + "node_id": "code", + "session_id": "ses_child", + "tool_call_id": "call_1", + "properties": { + "exit_code": 7, + "termination": "exited", + "duration_ms": 12, + "streams_separated": true, + "exec_output_tail": {"stdout": "out", "stderr": "err"}, + "visit": 1 + } + }); + + let parsed = RunEvent::from_value(value.clone()).unwrap(); + assert_eq!(parsed.event_name(), "agent.tool.process.completed"); + assert_eq!(parsed.tool_call_id.as_deref(), Some("call_1")); + let EventBody::AgentToolProcessCompleted(props) = &parsed.body else { + panic!("expected a typed process event, got {:?}", parsed.body); + }; + assert_eq!(props.exit_code, Some(7)); + assert_eq!(props.termination, CommandTermination::Exited); + assert_eq!(props.duration_ms, 12); + assert!(props.streams_separated); + assert_eq!( + props.exec_output_tail.as_ref().unwrap().stdout.as_deref(), + Some("out") + ); + + assert_eq!(parsed.to_value().unwrap(), value); + } + + #[test] + fn agent_tool_process_completed_omits_absent_exit_code_and_output_tail() { + let body = EventBody::AgentToolProcessCompleted(AgentToolProcessCompletedProps { + exit_code: None, + termination: CommandTermination::TimedOut, + duration_ms: 10_000, + streams_separated: false, + exec_output_tail: None, + visit: 1, + }); + + let value = serde_json::to_value(&body).unwrap(); + + assert_eq!(value["event"], "agent.tool.process.completed"); + assert_eq!(value["properties"]["termination"], "timed_out"); + assert_eq!(value["properties"]["streams_separated"], false); + let properties = value["properties"].as_object().unwrap(); + assert!(!properties.contains_key("exit_code")); + assert!(!properties.contains_key("exec_output_tail")); + + let parsed: EventBody = serde_json::from_value(value).unwrap(); + assert_eq!(parsed, body); + } + + /// The trace summary is what `AgentEvent::trace()` expands into tracing + /// fields, so it must carry sizes and truncation flags only. + #[test] + fn exec_output_tail_trace_summary_exposes_sizes_not_content() { + let tail = ExecOutputTail { + stdout: Some("secret stdout bytes".to_string()), + stderr: Some("secret stderr".to_string()), + stdout_truncated: true, + stderr_truncated: false, + }; + + let summary = ExecOutputTail::trace_summary(Some(&tail)); + + assert!(summary.present); + assert_eq!(summary.stdout_bytes, 19); + assert_eq!(summary.stderr_bytes, 13); + assert!(summary.stdout_truncated); + assert!(!summary.stderr_truncated); + let rendered = format!("{summary:?}"); + assert!(!rendered.contains("secret"), "got: {rendered}"); + } + #[test] fn agent_tool_source_and_category_use_public_json_shape() { assert_eq!( From c803354309d657780c5a7255cf0c2cb2c0034308 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 21:50:07 -0400 Subject: [PATCH 2/4] refactor(agent): streamline shell outcome reporting --- docs/internal/events.md | 4 + lib/components/fabro-agent/src/tools.rs | 144 ++++++++++-------- .../fabro-agent/tests/it/docker_shell.rs | 69 ++++----- .../fabro-sandbox/src/daytona/mod.rs | 29 ++-- lib/components/fabro-sandbox/src/docker.rs | 14 +- lib/components/fabro-sandbox/src/local.rs | 8 +- lib/components/fabro-sandbox/src/sandbox.rs | 64 ++++---- .../fabro-sandbox/src/test_support.rs | 24 +-- .../tests/daytona_streaming_live.rs | 4 +- .../fabro-sandbox/tests/docker_streaming.rs | 2 +- .../fabro-workflow/src/event/convert.rs | 9 +- .../fabro-workflow/src/handler/command.rs | 2 +- .../fabro-types/src/run_event/infra.rs | 27 ++++ .../fabro-types/src/run_event/mod.rs | 22 --- 14 files changed, 222 insertions(+), 200 deletions(-) diff --git a/docs/internal/events.md b/docs/internal/events.md index 57bb9c115..0a51e4e13 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -1190,6 +1190,10 @@ or launch failure) and when the tool ran without a session-bound emitter. | `duration_ms` | integer | Process duration | | `streams_separated` | boolean | `false` when the provider could not separate stdout from stderr; the combined output is then in `exec_output_tail.stdout` | | `exec_output_tail` | object | Bounded, redacted output tails; omitted when both streams were empty | +| `exec_output_tail.stdout` | string | Bounded stdout tail, or combined-output tail when `streams_separated` is `false`; omitted when empty | +| `exec_output_tail.stderr` | string | Bounded stderr tail; omitted when empty | +| `exec_output_tail.stdout_truncated` | boolean | `true` when earlier stdout bytes were omitted; omitted when `false` | +| `exec_output_tail.stderr_truncated` | boolean | `true` when earlier stderr bytes were omitted; omitted when `false` | | `visit` | integer | Stage visit | ### `agent.error` diff --git a/lib/components/fabro-agent/src/tools.rs b/lib/components/fabro-agent/src/tools.rs index 4043006fa..3797350df 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -8,9 +8,10 @@ use fabro_model::ModelHandle; #[cfg(test)] use fabro_static::EnvVars; use futures::{StreamExt, stream}; +use tokio::task; use crate::config::NativeToolOptions; -use crate::sandbox::{CommandOutputCallback, ExecStreamingResult, GrepOptions}; +use crate::sandbox::{ExecStreamingResult, GrepOptions}; use crate::tool_registry::{RegisteredTool, ToolRegistry, ToolSource}; use crate::types::AgentEvent; @@ -262,7 +263,7 @@ pub fn make_shell_tool_with_options(options: &NativeToolOptions) -> RegisteredTo None, tool_env.as_ref(), Some(ctx.cancel.clone()), - discard_output_callback(), + None, ) .await .map_err(|e| { @@ -270,15 +271,37 @@ pub fn make_shell_tool_with_options(options: &NativeToolOptions) -> RegisteredTo })?; let text = render_shell_result(&streaming); - ctx.emit_agent_event(AgentEvent::ToolProcessCompleted { - exit_code: streaming.result.exit_code, - termination: streaming.result.termination, - duration_ms: streaming.result.duration_ms, - streams_separated: streaming.streams_separated, - exec_output_tail: streaming.result.default_redacted_output_tail(), - }); + let is_success = streaming.result.is_success(); + if let Some(emitter) = ctx.agent_event_emitter { + let exit_code = streaming.result.exit_code; + let termination = streaming.result.termination; + let duration_ms = streaming.result.duration_ms; + let streams_separated = streaming.streams_separated; + let result = streaming.result; + let exec_output_tail = match task::spawn_blocking(move || { + result.default_redacted_output_tail() + }) + .await + { + Ok(exec_output_tail) => exec_output_tail, + Err(err) => { + tracing::warn!( + error = ?err, + "Failed to redact shell process output tail" + ); + None + } + }; + emitter.emit(AgentEvent::ToolProcessCompleted { + exit_code, + termination, + duration_ms, + streams_separated, + exec_output_tail, + }); + } - if streaming.result.is_success() { + if is_success { Ok(text) } else { Err(text) @@ -290,16 +313,9 @@ pub fn make_shell_tool_with_options(options: &NativeToolOptions) -> RegisteredTo } /// Prefix for shell failures that never produced an `ExecResult`, so the model -/// can tell "the process never ran" apart from "the process ran and failed". +/// can distinguish missing process diagnostics from a reported process failure. const SHELL_NO_PROCESS_RESULT: &str = "Shell command produced no process result"; -/// The agent shell tool consumes the streaming exec path for its stream -/// provenance and partial-output capture, but does not forward live output -/// deltas onto the agent protocol. -fn discard_output_callback() -> CommandOutputCallback { - Arc::new(|_stream, _bytes| Box::pin(async { Ok(()) })) -} - /// Renders the model-facing shell result: termination, exit code, duration, /// and provider-honest output sections. Metadata stays at the head and /// `stderr` at the tail so head/tail truncation preserves both. @@ -732,15 +748,18 @@ mod tests { use fabro_llm::provider::ProviderAdapter; use fabro_model::ProviderId; use fabro_types::CommandTermination; + use tokio::sync::broadcast; use tokio_util::sync::CancellationToken; use super::*; use crate::config::{NativeToolOptions, SessionOptions, ToolSecrets}; + use crate::event::{Emitter, SessionBoundEmitter}; use crate::local_sandbox::LocalSandbox; use crate::sandbox::*; use crate::test_support::MockSandbox; - use crate::tool_registry::{AgentEventEmitter, ToolContext}; + use crate::tool_registry::ToolContext; use crate::truncation; + use crate::types::SessionEvent; #[test] fn core_tool_descriptions_include_actionable_guidance() { @@ -1025,33 +1044,6 @@ mod tests { assert_eq!(written[0].1, "1 | keep this literal\ngoodbye"); } - /// Records the typed agent events a tool emits through its bound emitter. - #[derive(Default)] - struct RecordingAgentEmitter { - events: std::sync::Mutex>, - } - - impl RecordingAgentEmitter { - fn events(&self) -> Vec { - self.events.lock().expect("events lock poisoned").clone() - } - - fn only_process_event(&self) -> AgentEvent { - let events = self.events(); - assert_eq!(events.len(), 1, "expected one agent event, got {events:?}"); - events.into_iter().next().expect("one event") - } - } - - impl AgentEventEmitter for RecordingAgentEmitter { - fn emit(&self, event: AgentEvent) { - self.events - .lock() - .expect("events lock poisoned") - .push(event); - } - } - fn shell_context(env: Arc) -> ToolContext { ToolContext { env, @@ -1064,16 +1056,31 @@ mod tests { } } - fn shell_context_with_emitter( - env: Arc, - emitter: Arc, - ) -> ToolContext { + fn shell_context_with_emitter(env: Arc, emitter: &Emitter) -> ToolContext { ToolContext { - agent_event_emitter: Some(emitter), + session_id: Some("test-session".to_string()), + root_session_id: Some("test-session".to_string()), + tool_call_id: Some("call_1".to_string()), + agent_event_emitter: Some(Arc::new(SessionBoundEmitter { + emitter: emitter.clone(), + session_id: "test-session".to_string(), + tool_call_id: Some("call_1".to_string()), + })), ..shell_context(env) } } + fn only_process_event(receiver: &mut broadcast::Receiver) -> AgentEvent { + let event = receiver.try_recv().expect("one process event"); + assert_eq!(event.session_id, "test-session"); + assert_eq!(event.tool_call_id.as_deref(), Some("call_1")); + assert!(matches!( + receiver.try_recv(), + Err(broadcast::error::TryRecvError::Empty) + )); + event.event + } + fn mock_sandbox_with(result: ExecResult) -> Arc { Arc::new(MockSandbox { exec_result: result, @@ -1223,11 +1230,12 @@ mod tests { exec_error: Some("sandbox transport is down".into()), ..Default::default() }); - let emitter = Arc::new(RecordingAgentEmitter::default()); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); let output = (tool.executor)( serde_json::json!({"command": "make test"}), - shell_context_with_emitter(env, emitter.clone()), + shell_context_with_emitter(env, &emitter), ) .await .expect_err("a sandbox transport failure is a failed tool result"); @@ -1241,7 +1249,10 @@ mod tests { "got: {output}" ); assert!(!output.contains("Exit code"), "got: {output}"); - assert!(emitter.events().is_empty(), "got: {:?}", emitter.events()); + assert!(matches!( + receiver.try_recv(), + Err(broadcast::error::TryRecvError::Empty) + )); } #[tokio::test] @@ -1254,15 +1265,16 @@ mod tests { termination: CommandTermination::Exited, duration_ms: 12, }); - let emitter = Arc::new(RecordingAgentEmitter::default()); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); let _ = (tool.executor)( serde_json::json!({"command": "printf out; printf err >&2; exit 7"}), - shell_context_with_emitter(env, emitter.clone()), + shell_context_with_emitter(env, &emitter), ) .await; - match emitter.only_process_event() { + match only_process_event(&mut receiver) { AgentEvent::ToolProcessCompleted { exit_code, termination, @@ -1298,11 +1310,12 @@ mod tests { streams_separated: false, ..Default::default() }); - let emitter = Arc::new(RecordingAgentEmitter::default()); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); let output = (tool.executor)( serde_json::json!({"command": "echo interleaved"}), - shell_context_with_emitter(env, emitter.clone()), + shell_context_with_emitter(env, &emitter), ) .await .expect("exit 0 is a successful tool result"); @@ -1312,7 +1325,7 @@ mod tests { "got: {output}" ); assert!(!output.contains("stderr:"), "got: {output}"); - match emitter.only_process_event() { + match only_process_event(&mut receiver) { AgentEvent::ToolProcessCompleted { streams_separated, .. } => assert!(!streams_separated), @@ -1362,11 +1375,12 @@ mod tests { let env: Arc = Arc::new(LocalSandbox::new( std::env::current_dir().expect("current dir"), )); - let emitter = Arc::new(RecordingAgentEmitter::default()); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); let output = (tool.executor)( serde_json::json!({"command": "printf 'out'; printf 'err' >&2; exit 7"}), - shell_context_with_emitter(env, emitter.clone()), + shell_context_with_emitter(env, &emitter), ) .await .expect_err("exit 7 is a failed tool result"); @@ -1376,7 +1390,7 @@ mod tests { assert!(output.contains("stdout:\nout"), "got: {output}"); assert!(output.contains("stderr:\nerr"), "got: {output}"); - match emitter.only_process_event() { + match only_process_event(&mut receiver) { AgentEvent::ToolProcessCompleted { exit_code, termination, @@ -1395,8 +1409,8 @@ mod tests { } } - #[tokio::test] - async fn shell_public_schema_is_command_timeout_and_description() { + #[test] + fn shell_public_schema_is_command_timeout_and_description() { let tool = make_shell_tool(); assert_eq!( tool.definition.parameters, diff --git a/lib/components/fabro-agent/tests/it/docker_shell.rs b/lib/components/fabro-agent/tests/it/docker_shell.rs index 52d1dd081..ac0692fba 100644 --- a/lib/components/fabro-agent/tests/it/docker_shell.rs +++ b/lib/components/fabro-agent/tests/it/docker_shell.rs @@ -2,34 +2,22 @@ //! Docker provider's streaming path, which uses a `bash -lc` supervisor and //! separate stdout/stderr channels. -use std::sync::{Arc, Mutex}; +use std::sync::Arc; +use fabro_agent::event::SessionBoundEmitter; use fabro_agent::sandbox::Sandbox; -use fabro_agent::tool_registry::{AgentEventEmitter, ToolContext}; +use fabro_agent::tool_registry::ToolContext; use fabro_agent::tools::make_shell_tool; use fabro_agent::types::AgentEvent; -use fabro_agent::{DockerSandbox, DockerSandboxOptions}; +use fabro_agent::{DockerSandbox, DockerSandboxOptions, Emitter}; use fabro_types::CommandTermination; +use tokio::sync::broadcast; use tokio_util::sync::CancellationToken; -#[derive(Default)] -struct RecordingAgentEmitter { - events: Mutex>, -} - -impl AgentEventEmitter for RecordingAgentEmitter { - fn emit(&self, event: AgentEvent) { - self.events - .lock() - .expect("events lock poisoned") - .push(event); - } -} - #[tokio::test] #[ignore = "requires real Docker container lifecycle; run explicitly when changing shell tool exec integration"] async fn shell_reports_real_docker_process_outcome() { - let Ok(sandbox) = DockerSandbox::new( + let sandbox = DockerSandbox::new( DockerSandboxOptions { image: "buildpack-deps:noble".to_string(), auto_pull: false, @@ -40,16 +28,16 @@ async fn shell_reports_real_docker_process_outcome() { None, None, None, - ) else { - return; - }; - // No Docker daemon or no local image: nothing to prove here. - if sandbox.initialize().await.is_err() { - return; - } + ) + .expect("docker sandbox should construct"); + sandbox + .initialize() + .await + .expect("docker sandbox should initialize"); let sandbox = Arc::new(sandbox); - let emitter = Arc::new(RecordingAgentEmitter::default()); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); let tool = make_shell_tool(); let result = (tool.executor)( serde_json::json!({"command": "printf 'out'; printf 'err' >&2; exit 7"}), @@ -57,10 +45,14 @@ async fn shell_reports_real_docker_process_outcome() { env: sandbox.clone() as Arc, cancel: CancellationToken::new(), tool_env_provider: None, - session_id: None, - root_session_id: None, - tool_call_id: None, - agent_event_emitter: Some(emitter.clone()), + session_id: Some("test-session".to_string()), + root_session_id: Some("test-session".to_string()), + tool_call_id: Some("call_1".to_string()), + agent_event_emitter: Some(Arc::new(SessionBoundEmitter { + emitter, + session_id: "test-session".to_string(), + tool_call_id: Some("call_1".to_string()), + })), }, ) .await; @@ -75,9 +67,14 @@ async fn shell_reports_real_docker_process_outcome() { assert!(output.contains("stdout:\nout"), "got: {output}"); assert!(output.contains("stderr:\nerr"), "got: {output}"); - let events = emitter.events.lock().expect("events lock poisoned").clone(); - assert_eq!(events.len(), 1, "expected one agent event, got {events:?}"); - match &events[0] { + let event = receiver.try_recv().expect("one process event"); + assert_eq!(event.session_id, "test-session"); + assert_eq!(event.tool_call_id.as_deref(), Some("call_1")); + assert!(matches!( + receiver.try_recv(), + Err(broadcast::error::TryRecvError::Empty) + )); + match event.event { AgentEvent::ToolProcessCompleted { exit_code, termination, @@ -85,10 +82,10 @@ async fn shell_reports_real_docker_process_outcome() { exec_output_tail, .. } => { - assert_eq!(*exit_code, Some(7)); - assert_eq!(*termination, CommandTermination::Exited); + assert_eq!(exit_code, Some(7)); + assert_eq!(termination, CommandTermination::Exited); assert!(streams_separated); - let tail = exec_output_tail.as_ref().expect("output tail"); + let tail = exec_output_tail.expect("output tail"); assert_eq!(tail.stdout.as_deref(), Some("out")); assert_eq!(tail.stderr.as_deref(), Some("err")); } diff --git a/lib/components/fabro-sandbox/src/daytona/mod.rs b/lib/components/fabro-sandbox/src/daytona/mod.rs index fa227d2c8..693a93307 100644 --- a/lib/components/fabro-sandbox/src/daytona/mod.rs +++ b/lib/components/fabro-sandbox/src/daytona/mod.rs @@ -1637,7 +1637,7 @@ impl Sandbox for DaytonaSandbox { working_dir: Option<&str>, env_vars: Option<&HashMap>, cancel_token: Option, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result { let sandbox = self.sandbox()?; let start = Instant::now(); @@ -1694,9 +1694,11 @@ impl Sandbox for DaytonaSandbox { if !bytes.is_empty() { saw_live_chunk.store(true, Ordering::Relaxed); stdout_seen.lock().await.extend_from_slice(&bytes); - callback(CommandOutputStream::Stdout, bytes) - .await - .map_err(|err| daytona_callback_error(&err))?; + if let Some(callback) = callback { + callback(CommandOutputStream::Stdout, bytes) + .await + .map_err(|err| daytona_callback_error(&err))?; + } } Ok(()) } @@ -1710,9 +1712,11 @@ impl Sandbox for DaytonaSandbox { if !bytes.is_empty() { saw_live_chunk.store(true, Ordering::Relaxed); stderr_seen.lock().await.extend_from_slice(&bytes); - callback(CommandOutputStream::Stderr, bytes) - .await - .map_err(|err| daytona_callback_error(&err))?; + if let Some(callback) = callback { + callback(CommandOutputStream::Stderr, bytes) + .await + .map_err(|err| daytona_callback_error(&err))?; + } } Ok(()) } @@ -1769,14 +1773,14 @@ impl Sandbox for DaytonaSandbox { CommandOutputStream::Stdout, logs.stdout.as_bytes(), &stdout_seen, - &output_callback, + output_callback.as_ref(), ) .await?; append_missing_log_suffix( CommandOutputStream::Stderr, logs.stderr.as_bytes(), &stderr_seen, - &output_callback, + output_callback.as_ref(), ) .await?; } @@ -2181,7 +2185,7 @@ async fn append_missing_log_suffix( stream: CommandOutputStream, final_bytes: &[u8], seen: &Arc>>, - output_callback: &CommandOutputCallback, + output_callback: Option<&CommandOutputCallback>, ) -> crate::Result<()> { if final_bytes.is_empty() { return Ok(()); @@ -2196,7 +2200,10 @@ async fn append_missing_log_suffix( let missing = final_bytes[offset..].to_vec(); seen.extend_from_slice(&missing); drop(seen); - output_callback(stream, missing).await + match output_callback { + Some(output_callback) => output_callback(stream, missing).await, + None => Ok(()), + } } fn missing_log_suffix_offset(seen: &[u8], final_bytes: &[u8]) -> usize { diff --git a/lib/components/fabro-sandbox/src/docker.rs b/lib/components/fabro-sandbox/src/docker.rs index 51f52a41c..07f39b10a 100644 --- a/lib/components/fabro-sandbox/src/docker.rs +++ b/lib/components/fabro-sandbox/src/docker.rs @@ -340,7 +340,7 @@ impl DockerSandbox { cmd: Vec, working_dir: Option, env: Option>, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result<(Vec, Vec, i32)> { let exec_opts = CreateExecOptions { cmd: Some(cmd), @@ -369,11 +369,15 @@ impl DockerSandbox { match chunk { Ok(LogOutput::StdOut { message }) => { stdout.extend_from_slice(&message); - output_callback(CommandOutputStream::Stdout, message.to_vec()).await?; + if let Some(output_callback) = output_callback.as_ref() { + output_callback(CommandOutputStream::Stdout, message.to_vec()).await?; + } } Ok(LogOutput::StdErr { message }) => { stderr.extend_from_slice(&message); - output_callback(CommandOutputStream::Stderr, message.to_vec()).await?; + if let Some(output_callback) = output_callback.as_ref() { + output_callback(CommandOutputStream::Stderr, message.to_vec()).await?; + } } Ok(_) => {} Err(e) => { @@ -460,7 +464,7 @@ impl DockerSandbox { working_dir: Option<&str>, env_vars: Option<&HashMap>, cancel_token: Option, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result { let start = Instant::now(); let effective_dir = working_dir @@ -1548,7 +1552,7 @@ impl Sandbox for DockerSandbox { working_dir: Option<&str>, env_vars: Option<&HashMap>, cancel_token: Option, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result { let dir = working_dir.map(|path| self.resolve_container_path(path)); self.docker_exec_shell_streaming( diff --git a/lib/components/fabro-sandbox/src/local.rs b/lib/components/fabro-sandbox/src/local.rs index ca21e5631..77fcd79f5 100644 --- a/lib/components/fabro-sandbox/src/local.rs +++ b/lib/components/fabro-sandbox/src/local.rs @@ -419,7 +419,7 @@ impl Sandbox for LocalSandbox { working_dir: Option<&str>, env_vars: Option<&std::collections::HashMap>, cancel_token: Option, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result { let start = Instant::now(); @@ -831,7 +831,7 @@ async fn sigterm_then_kill(child: &mut Child) { async fn drain_command_pipe( mut reader: Option, stream: CommandOutputStream, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result> where R: AsyncRead + Unpin, @@ -851,7 +851,9 @@ where return Ok(output); } output.extend_from_slice(&buf[..read]); - output_callback(stream, buf[..read].to_vec()).await?; + if let Some(output_callback) = output_callback.as_ref() { + output_callback(stream, buf[..read].to_vec()).await?; + } } } diff --git a/lib/components/fabro-sandbox/src/sandbox.rs b/lib/components/fabro-sandbox/src/sandbox.rs index fa2e9bacb..8bc9cb801 100644 --- a/lib/components/fabro-sandbox/src/sandbox.rs +++ b/lib/components/fabro-sandbox/src/sandbox.rs @@ -111,7 +111,7 @@ macro_rules! delegate_sandbox { working_dir: Option<&str>, env_vars: Option<&std::collections::HashMap>, cancel_token: Option, - output_callback: $crate::CommandOutputCallback, + output_callback: Option<$crate::CommandOutputCallback>, ) -> $crate::Result<$crate::ExecStreamingResult> { self.$field .exec_command_streaming( @@ -680,6 +680,34 @@ pub type CommandOutputCallback = Arc< + Sync, >; +pub(crate) async fn replay_exec_result( + result: ExecResult, + streams_separated: bool, + output_callback: Option<&CommandOutputCallback>, +) -> crate::Result { + if let Some(output_callback) = output_callback { + if !result.stdout.is_empty() { + output_callback( + CommandOutputStream::Stdout, + result.stdout.as_bytes().to_vec(), + ) + .await?; + } + if !result.stderr.is_empty() { + output_callback( + CommandOutputStream::Stderr, + result.stderr.as_bytes().to_vec(), + ) + .await?; + } + } + Ok(ExecStreamingResult { + result, + streams_separated, + live_streaming: false, + }) +} + pub struct StdioProcess { pub stdin: Pin>, pub stdout: Pin>, @@ -856,11 +884,13 @@ pub trait Sandbox: Send + Sync { /// /// **Production sandboxes must override this.** The default falls back to /// the non-streaming [`exec_command`](Self::exec_command) and replays its - /// output through `output_callback` at the end, marking - /// `live_streaming: false`. That's the right behavior for test mocks but - /// silently drops live output for any real sandbox that wraps another — - /// decorators in particular must forward to the inner sandbox's streaming - /// implementation rather than relying on this default. + /// output through `output_callback` at the end when one is supplied, + /// marking `live_streaming: false`. Passing `None` captures the final + /// result without paying per-chunk callback costs. That's the right + /// behavior for test mocks but silently drops live output for any real + /// sandbox that wraps another — decorators in particular must forward to + /// the inner sandbox's streaming implementation rather than relying on + /// this default. async fn exec_command_streaming( &self, command: &str, @@ -868,7 +898,7 @@ pub trait Sandbox: Send + Sync { working_dir: Option<&str>, env_vars: Option<&std::collections::HashMap>, cancel_token: Option, - output_callback: CommandOutputCallback, + output_callback: Option, ) -> crate::Result { let fallback_timeout_ms = timeout_ms.unwrap_or(u64::MAX); let result = self @@ -880,25 +910,7 @@ pub trait Sandbox: Send + Sync { cancel_token, ) .await?; - if !result.stdout.is_empty() { - output_callback( - CommandOutputStream::Stdout, - result.stdout.as_bytes().to_vec(), - ) - .await?; - } - if !result.stderr.is_empty() { - output_callback( - CommandOutputStream::Stderr, - result.stderr.as_bytes().to_vec(), - ) - .await?; - } - Ok(ExecStreamingResult { - result, - streams_separated: true, - live_streaming: false, - }) + replay_exec_result(result, true, output_callback.as_ref()).await } async fn spawn_stdio_process( diff --git a/lib/components/fabro-sandbox/src/test_support.rs b/lib/components/fabro-sandbox/src/test_support.rs index d48f552e2..f976f5069 100644 --- a/lib/components/fabro-sandbox/src/test_support.rs +++ b/lib/components/fabro-sandbox/src/test_support.rs @@ -9,7 +9,7 @@ use tokio::io::{DuplexStream, duplex}; use tokio::time::sleep; use tokio_util::sync::CancellationToken; -use crate::sandbox::StdioProcessControl; +use crate::sandbox::{self, StdioProcessControl}; use crate::{ DEFAULT_EXEC_OUTPUT_TAIL_BYTES, DirEntry, ExecResult, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback, StderrCollector, StdioProcess, StdioProcessHandle, @@ -254,7 +254,7 @@ impl Sandbox for MockSandbox { working_dir: Option<&str>, env_vars: Option<&std::collections::HashMap>, cancel_token: Option, - output_callback: crate::CommandOutputCallback, + output_callback: Option, ) -> crate::Result { let result = self .exec_command( @@ -265,25 +265,7 @@ impl Sandbox for MockSandbox { cancel_token, ) .await?; - if !result.stdout.is_empty() { - output_callback( - fabro_types::CommandOutputStream::Stdout, - result.stdout.as_bytes().to_vec(), - ) - .await?; - } - if !result.stderr.is_empty() { - output_callback( - fabro_types::CommandOutputStream::Stderr, - result.stderr.as_bytes().to_vec(), - ) - .await?; - } - Ok(crate::ExecStreamingResult { - result, - streams_separated: self.streams_separated, - live_streaming: false, - }) + sandbox::replay_exec_result(result, self.streams_separated, output_callback.as_ref()).await } async fn spawn_stdio_process( diff --git a/lib/components/fabro-sandbox/tests/daytona_streaming_live.rs b/lib/components/fabro-sandbox/tests/daytona_streaming_live.rs index 4b5dd6590..40f3a9b77 100644 --- a/lib/components/fabro-sandbox/tests/daytona_streaming_live.rs +++ b/lib/components/fabro-sandbox/tests/daytona_streaming_live.rs @@ -265,7 +265,7 @@ mod daytona_streaming_live { None, None, Some(cancel_for_exec), - callback, + Some(callback), ) .await }); @@ -390,7 +390,7 @@ mod daytona_streaming_live { None, None, cancel_token, - callback, + Some(callback), ) .await?; let chunks = chunks.lock().await.clone(); diff --git a/lib/components/fabro-sandbox/tests/docker_streaming.rs b/lib/components/fabro-sandbox/tests/docker_streaming.rs index f77674a30..6f4cd7562 100644 --- a/lib/components/fabro-sandbox/tests/docker_streaming.rs +++ b/lib/components/fabro-sandbox/tests/docker_streaming.rs @@ -53,7 +53,7 @@ async fn streaming_timeout_terminates_docker_exec_before_returning() { None, None, None, - callback, + Some(callback), ) .await .expect("streaming command should return a timeout result"); diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 1731cd265..f6e2de0f1 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -1543,12 +1543,7 @@ mod tests { termination: ::fabro_types::CommandTermination::Exited, duration_ms: 12, streams_separated: true, - exec_output_tail: Some(::fabro_types::ExecOutputTail { - stdout: Some("out".to_string()), - stderr: Some("err".to_string()), - stdout_truncated: false, - stderr_truncated: false, - }), + exec_output_tail: Some(exec_tail()), }, session_id: Some("ses_child".to_string()), parent_session_id: Some("ses_parent".to_string()), @@ -1575,7 +1570,7 @@ mod tests { assert_eq!(properties["termination"], "exited"); assert_eq!(properties["duration_ms"], 12); assert_eq!(properties["streams_separated"], true); - assert_eq!(properties["exec_output_tail"]["stdout"], "out"); + assert_eq!(properties["exec_output_tail"]["stdout"], "last stdout line"); assert_eq!(properties["visit"], 2); } diff --git a/lib/components/fabro-workflow/src/handler/command.rs b/lib/components/fabro-workflow/src/handler/command.rs index 54fe0ff20..74a32f0ac 100644 --- a/lib/components/fabro-workflow/src/handler/command.rs +++ b/lib/components/fabro-workflow/src/handler/command.rs @@ -128,7 +128,7 @@ impl Handler for CommandHandler { None, env_vars, Some(cancel_token.clone()), - output_callback, + Some(output_callback), ) .await; cancel_token.cancel(); diff --git a/lib/foundation/fabro-types/src/run_event/infra.rs b/lib/foundation/fabro-types/src/run_event/infra.rs index 8af295ffa..5348b4931 100644 --- a/lib/foundation/fabro-types/src/run_event/infra.rs +++ b/lib/foundation/fabro-types/src/run_event/infra.rs @@ -401,3 +401,30 @@ pub struct CliEnsureFailedProps { #[serde(default, skip_serializing_if = "Option::is_none")] pub exec_output_tail: Option, } + +#[cfg(test)] +mod tests { + use super::ExecOutputTail; + + /// The trace summary is expanded into tracing fields, so it must carry + /// sizes and truncation flags only. + #[test] + fn exec_output_tail_trace_summary_exposes_sizes_not_content() { + let tail = ExecOutputTail { + stdout: Some("secret stdout bytes".to_string()), + stderr: Some("secret stderr".to_string()), + stdout_truncated: true, + stderr_truncated: false, + }; + + let summary = ExecOutputTail::trace_summary(Some(&tail)); + + assert!(summary.present); + assert_eq!(summary.stdout_bytes, 19); + assert_eq!(summary.stderr_bytes, 13); + assert!(summary.stdout_truncated); + assert!(!summary.stderr_truncated); + let rendered = format!("{summary:?}"); + assert!(!rendered.contains("secret"), "got: {rendered}"); + } +} diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index 4b937f682..b79f7758b 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -2475,28 +2475,6 @@ mod tests { assert_eq!(parsed, body); } - /// The trace summary is what `AgentEvent::trace()` expands into tracing - /// fields, so it must carry sizes and truncation flags only. - #[test] - fn exec_output_tail_trace_summary_exposes_sizes_not_content() { - let tail = ExecOutputTail { - stdout: Some("secret stdout bytes".to_string()), - stderr: Some("secret stderr".to_string()), - stdout_truncated: true, - stderr_truncated: false, - }; - - let summary = ExecOutputTail::trace_summary(Some(&tail)); - - assert!(summary.present); - assert_eq!(summary.stdout_bytes, 19); - assert_eq!(summary.stderr_bytes, 13); - assert!(summary.stdout_truncated); - assert!(!summary.stderr_truncated); - let rendered = format!("{summary:?}"); - assert!(!rendered.contains("secret"), "got: {rendered}"); - } - #[test] fn agent_tool_source_and_category_use_public_json_shape() { assert_eq!( From 4666f51d98fa7ea320390bea020eb2cfc0dc5890 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 22:34:02 -0400 Subject: [PATCH 3/4] feat(model): add Claude Opus 5 to OpenRouter --- docs/public/changelog/2026-07-24.mdx | 2 +- docs/public/integrations/openrouter.mdx | 2 +- lib/foundation/fabro-model/src/catalog.rs | 20 +++++++++++++ .../src/catalog/providers/openrouter.toml | 28 ++++++++++++++++++- 4 files changed, 49 insertions(+), 3 deletions(-) diff --git a/docs/public/changelog/2026-07-24.mdx b/docs/public/changelog/2026-07-24.mdx index 35e35d846..9ae1d18ab 100644 --- a/docs/public/changelog/2026-07-24.mdx +++ b/docs/public/changelog/2026-07-24.mdx @@ -40,7 +40,7 @@ When a run resumes after a node was cancelled or lost mid-flight, the replay now -- Added `claude-opus-5` to the first-party Anthropic model catalog; the `opus` and `claude-opus` aliases now resolve to Opus 5 +- Added `claude-opus-5` to the first-party Anthropic and optional OpenRouter model catalogs; the `opus` and `claude-opus` aliases now resolve to Opus 5 - Added `gpt-sol`, `gpt-terra`, and `gpt-luna` aliases for GPT-5.6 offerings - Added portable `glm`, `glm52`, `glm5.2`, `deepseek`, and `deepseek-flash` aliases across direct and OpenRouter offerings diff --git a/docs/public/integrations/openrouter.mdx b/docs/public/integrations/openrouter.mdx index bd0c186dd..c3bf7ff2c 100644 --- a/docs/public/integrations/openrouter.mdx +++ b/docs/public/integrations/openrouter.mdx @@ -48,7 +48,7 @@ The built-in catalog gives OpenRouter offerings the same human-facing model slug | Fabro model slug | OpenRouter API ID / notes | | --- | --- | -| `claude-fable-5`, `claude-opus-4-8`, `claude-opus-4-7` | Matching `anthropic/...` API IDs; Anthropic-style cache billing | +| `claude-fable-5`, `claude-opus-5`, `claude-opus-4-8`, `claude-opus-4-7` | Matching `anthropic/...` API IDs; Anthropic-style cache billing | | `claude-sonnet-4-6` | `anthropic/claude-sonnet-4.6`; provider default | | `claude-haiku-4-5` | `anthropic/claude-haiku-4.5`; provider small default | | `gpt-5.6-sol`, `gpt-5.6-terra`, `gpt-5.6-luna`, `gpt-5.4`, `gpt-5.5` | Matching `openai/...` API IDs | diff --git a/lib/foundation/fabro-model/src/catalog.rs b/lib/foundation/fabro-model/src/catalog.rs index 7fa9eb161..f012a62ff 100644 --- a/lib/foundation/fabro-model/src/catalog.rs +++ b/lib/foundation/fabro-model/src/catalog.rs @@ -3083,6 +3083,19 @@ enabled = true false, BillingPolicy::OpenAi, ), + ( + "claude-opus-5", + "anthropic/claude-opus-5", + "claude-5", + 1_000_000, + 5.0, + 25.0, + 0.5, + ReasoningEffortFeature::Levels, + false, + true, + BillingPolicy::Anthropic, + ), ( "claude-opus-4-8", "anthropic/claude-opus-4.8", @@ -3161,6 +3174,13 @@ enabled = true "{id}" ); } + + for alias in ["opus", "claude-opus"] { + let model = catalog + .resolve_on_provider(&ProviderId::new("openrouter"), alias) + .unwrap_or_else(|error| panic!("{alias} should resolve on OpenRouter: {error}")); + assert_eq!(model.id, "claude-opus-5", "{alias}"); + } } #[test] diff --git a/lib/foundation/fabro-model/src/catalog/providers/openrouter.toml b/lib/foundation/fabro-model/src/catalog/providers/openrouter.toml index eea168785..2878ff33d 100644 --- a/lib/foundation/fabro-model/src/catalog/providers/openrouter.toml +++ b/lib/foundation/fabro-model/src/catalog/providers/openrouter.toml @@ -58,6 +58,33 @@ input_cost_per_mtok = 10.0 output_cost_per_mtok = 50.0 cache_input_cost_per_mtok = 1.0 +[providers.openrouter.models."claude-opus-5"] +api_id = "anthropic/claude-opus-5" +display_name = "Claude Opus 5 (via OpenRouter)" +family = "claude-5" +billing_policy = "anthropic" +training = "2026-05-01" +knowledge_cutoff = "May 2026" +aliases = ["opus", "claude-opus"] + +[providers.openrouter.models."claude-opus-5".limits] +context_window = 1000000 +max_output = 128000 + +[providers.openrouter.models."claude-opus-5".features] +tools = true +vision = true +reasoning = true +reasoning_effort = "levels" +prompt_cache = true +cache_control_breakpoints = true +sampling_params = false + +[providers.openrouter.models."claude-opus-5".costs] +input_cost_per_mtok = 5.0 +output_cost_per_mtok = 25.0 +cache_input_cost_per_mtok = 0.5 + [providers.openrouter.models."claude-opus-4-8"] api_id = "anthropic/claude-opus-4.8" display_name = "Claude Opus 4.8 (via OpenRouter)" @@ -65,7 +92,6 @@ family = "claude-4" billing_policy = "anthropic" training = "2026-01-01" knowledge_cutoff = "Jan 2026" -aliases = ["opus", "claude-opus"] [providers.openrouter.models."claude-opus-4-8".limits] context_window = 1000000 From 4d5458b64cea5d3acc81302ea2734fcfbba8e5bf Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 22:35:51 -0400 Subject: [PATCH 4/4] test(agent): honor Docker shell integration preconditions --- .../fabro-agent/tests/it/docker_shell.rs | 19 ++++++++++--------- 1 file changed, 10 insertions(+), 9 deletions(-) diff --git a/lib/components/fabro-agent/tests/it/docker_shell.rs b/lib/components/fabro-agent/tests/it/docker_shell.rs index ac0692fba..48c0d25b3 100644 --- a/lib/components/fabro-agent/tests/it/docker_shell.rs +++ b/lib/components/fabro-agent/tests/it/docker_shell.rs @@ -17,7 +17,7 @@ use tokio_util::sync::CancellationToken; #[tokio::test] #[ignore = "requires real Docker container lifecycle; run explicitly when changing shell tool exec integration"] async fn shell_reports_real_docker_process_outcome() { - let sandbox = DockerSandbox::new( + let Ok(sandbox) = DockerSandbox::new( DockerSandboxOptions { image: "buildpack-deps:noble".to_string(), auto_pull: false, @@ -28,12 +28,13 @@ async fn shell_reports_real_docker_process_outcome() { None, None, None, - ) - .expect("docker sandbox should construct"); - sandbox - .initialize() - .await - .expect("docker sandbox should initialize"); + ) else { + return; + }; + // No Docker daemon or no local image: the integration precondition is not met. + if sandbox.initialize().await.is_err() { + return; + } let sandbox = Arc::new(sandbox); let emitter = Emitter::new(); @@ -49,8 +50,8 @@ async fn shell_reports_real_docker_process_outcome() { root_session_id: Some("test-session".to_string()), tool_call_id: Some("call_1".to_string()), agent_event_emitter: Some(Arc::new(SessionBoundEmitter { - emitter, - session_id: "test-session".to_string(), + emitter: emitter.clone(), + session_id: "test-session".to_string(), tool_call_id: Some("call_1".to_string()), })), },