diff --git a/docs/internal/events.md b/docs/internal/events.md index eea2d6878..0a51e4e13 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -1157,6 +1157,45 @@ 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 | +| `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` Emitted when the agent encounters an error. diff --git a/docs/public/changelog/2026-07-24.mdx b/docs/public/changelog/2026-07-24.mdx index ace5ee3e6..1326013c6 100644 --- a/docs/public/changelog/2026-07-24.mdx +++ b/docs/public/changelog/2026-07-24.mdx @@ -54,7 +54,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/components/fabro-agent/src/profiles/kimi_tools.rs b/lib/components/fabro-agent/src/profiles/kimi_tools.rs index 3a7559fd6..9b9ce7c14 100644 --- a/lib/components/fabro-agent/src/profiles/kimi_tools.rs +++ b/lib/components/fabro-agent/src/profiles/kimi_tools.rs @@ -30,8 +30,8 @@ use crate::native_tool::NativeTool; use crate::sandbox::{GrepOptions, format_lines_numbered}; use crate::tool_registry::{RegisteredTool, ToolSource}; use crate::tools::{ - DEFAULT_READ_LINES, execute_grep, execute_shell_command, grep_result_path, make_edit_file_tool, - optional_usize_arg, required_str, + DEFAULT_READ_LINES, emit_shell_process_completed, execute_grep, execute_shell_command, + grep_result_path, make_edit_file_tool, optional_usize_arg, required_str, }; const DEFAULT_GREP_RESULTS: usize = 250; @@ -120,7 +120,8 @@ explicitly asked. Never run commands requiring superuser privileges unless expli None => default_timeout_ms, }; - let result = execute_shell_command(&ctx, command, timeout_ms, cwd).await?; + let streaming = execute_shell_command(&ctx, command, timeout_ms, cwd).await?; + let result = &streaming.result; let mut out = String::new(); if result.is_timed_out() { @@ -141,7 +142,9 @@ explicitly asked. Never run commands requiring superuser privileges unless expli } let _ = write!(out, "Command failed with exit code: {code}"); } - Ok(out) + let is_success = result.is_success(); + emit_shell_process_completed(&ctx, streaming).await; + if is_success { Ok(out) } else { Err(out) } }) }), source: ToolSource::Native, @@ -699,7 +702,7 @@ mod tests { tool_ctx, ) .await - .unwrap(); + .expect_err("a timeout is a failed tool result"); assert!(output.starts_with("Command timed out.\n"), "{output}"); assert_eq!(*env.captured_timeout.lock().unwrap(), Some(7_000)); @@ -707,12 +710,9 @@ mod tests { Some("/repo".to_string()) ]); assert_eq!(*env.captured_env_vars.lock().unwrap(), Some(tool_env)); - assert!( - env.captured_command - .lock() - .unwrap() - .as_deref() - .is_some_and(|command| command.starts_with("exec 2>&1\n")) + assert_eq!( + env.captured_command.lock().unwrap().as_deref(), + Some("echo $TOKEN") ); } } 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 5cdb9a1e4..7685a8f97 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -8,10 +8,12 @@ use fabro_model::ModelHandle; #[cfg(test)] use fabro_static::EnvVars; use futures::{StreamExt, stream}; +use tokio::task; use crate::config::NativeToolOptions; -use crate::sandbox::{ExecResult, GrepOptions}; +use crate::sandbox::{ExecStreamingResult, GrepOptions}; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource}; +use crate::types::AgentEvent; const MAX_WEB_FETCH_BYTES: usize = 100 * 1024; const MAX_READ_MANY_FILES_CONCURRENCY: usize = 8; @@ -259,29 +261,27 @@ pub fn make_shell_tool_with_options(options: &NativeToolOptions) -> RegisteredTo .unwrap_or(default_timeout) .min(max_timeout); - let result = execute_shell_command(&ctx, command, timeout_ms, None).await?; + let streaming = execute_shell_command(&ctx, command, timeout_ms, None).await?; - 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); + let is_success = streaming.result.is_success(); + emit_shell_process_completed(&ctx, streaming).await; + + if 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 distinguish missing process diagnostics from a reported process failure. +const SHELL_NO_PROCESS_RESULT: &str = "Shell command produced no process result"; + /// Execute a shell command with the session's environment and cancellation /// plumbing. Provider profiles can vary their wire schema and result /// rendering without accidentally bypassing those shared semantics. @@ -290,23 +290,88 @@ pub(crate) async fn execute_shell_command( command: &str, timeout_ms: u64, cwd: Option<&str>, -) -> Result { - let command = format!("exec 2>&1\n{command}"); - let tool_env = ctx.resolve_tool_env().await.map_err(|e| format!("{e:#}"))?; +) -> Result { + 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" ); ctx.env - .exec_command( - &command, - timeout_ms, + .exec_command_streaming( + command, + Some(timeout_ms), cwd, tool_env.as_ref(), Some(ctx.cancel.clone()), + None, ) .await - .map_err(|e| e.display_with_causes()) + .map_err(|e| format!("{SHELL_NO_PROCESS_RESULT}: {}", e.display_with_causes())) +} + +/// Emit the subordinate process outcome after model-facing output has been +/// rendered. Consumes the raw result so redaction does not require cloning +/// potentially large process output. +pub(crate) async fn emit_shell_process_completed( + ctx: &ToolContext, + streaming: ExecStreamingResult, +) { + if ctx.agent_event_emitter.is_none() { + return; + } + + 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 + } + }; + ctx.emit_agent_event(AgentEvent::ToolProcessCompleted { + exit_code, + termination, + duration_ms, + streams_separated, + exec_output_tail, + }); +} + +/// 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] @@ -746,13 +811,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, ToolSecrets}; + 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::ToolContext; + use crate::truncation; + use crate::types::SessionEvent; #[test] fn core_tool_descriptions_include_actionable_guidance() { @@ -1093,20 +1163,8 @@ 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 { + fn shell_context(env: Arc) -> ToolContext { + ToolContext { env, cancel: CancellationToken::new(), tool_env_provider: None, @@ -1114,11 +1172,87 @@ mod tests { root_session_id: None, tool_call_id: None, agent_event_emitter: None, + } + } + + fn shell_context_with_emitter(env: Arc, emitter: &Emitter) -> ToolContext { + ToolContext { + 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, + ..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] @@ -1155,46 +1289,243 @@ 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 = Emitter::new(); + let mut receiver = emitter.subscribe(); + + let output = (tool.executor)( + serde_json::json!({"command": "make test"}), + shell_context_with_emitter(env, &emitter), + ) + .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!(matches!( + receiver.try_recv(), + Err(broadcast::error::TryRecvError::Empty) + )); + } + + #[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 = 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), + ) + .await; + + match only_process_event(&mut receiver) { + 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 = Emitter::new(); + let mut receiver = emitter.subscribe(); + + let output = (tool.executor)( + serde_json::json!({"command": "echo interleaved"}), + shell_context_with_emitter(env, &emitter), + ) + .await + .expect("exit 0 is a successful tool result"); + + assert!( + output.contains("output (combined):\ninterleaved"), + "got: {output}" + ); + assert!(!output.contains("stderr:"), "got: {output}"); + match only_process_event(&mut receiver) { + 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 = 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), + ) + .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 only_process_event(&mut receiver) { + 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] 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..48c0d25b3 --- /dev/null +++ b/lib/components/fabro-agent/tests/it/docker_shell.rs @@ -0,0 +1,95 @@ +//! 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; + +use fabro_agent::event::SessionBoundEmitter; +use fabro_agent::sandbox::Sandbox; +use fabro_agent::tool_registry::ToolContext; +use fabro_agent::tools::make_shell_tool; +use fabro_agent::types::AgentEvent; +use fabro_agent::{DockerSandbox, DockerSandboxOptions, Emitter}; +use fabro_types::CommandTermination; +use tokio::sync::broadcast; +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 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: the integration precondition is not met. + if sandbox.initialize().await.is_err() { + return; + } + + let sandbox = Arc::new(sandbox); + 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"}), + ToolContext { + env: sandbox.clone() as Arc, + cancel: CancellationToken::new(), + tool_env_provider: None, + 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()), + })), + }, + ) + .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 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, + 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:?}"), + } +} 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/daytona/mod.rs b/lib/components/fabro-sandbox/src/daytona/mod.rs index 068789388..12a253e82 100644 --- a/lib/components/fabro-sandbox/src/daytona/mod.rs +++ b/lib/components/fabro-sandbox/src/daytona/mod.rs @@ -1707,7 +1707,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(); @@ -1765,9 +1765,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(()) } @@ -1781,9 +1783,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(()) } @@ -1840,14 +1844,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?; } @@ -2252,7 +2256,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(()); @@ -2267,7 +2271,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 e488c9fb8..992fb48a7 100644 --- a/lib/components/fabro-sandbox/src/docker.rs +++ b/lib/components/fabro-sandbox/src/docker.rs @@ -370,7 +370,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), @@ -399,11 +399,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) => { @@ -489,7 +493,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 @@ -1579,7 +1583,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 6eadb7424..919c58301 100644 --- a/lib/components/fabro-sandbox/src/local.rs +++ b/lib/components/fabro-sandbox/src/local.rs @@ -508,7 +508,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(); @@ -950,7 +950,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, @@ -970,7 +970,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?; + } } } @@ -1256,7 +1258,7 @@ mod tests { None, None, None, - Arc::new(|_, _| Box::pin(async { Ok(()) })), + Some(Arc::new(|_, _| Box::pin(async { Ok(()) }))), ) .await .unwrap(); @@ -1327,7 +1329,7 @@ mod tests { None, Some(&env_vars), None, - Arc::new(|_, _| Box::pin(async { Ok(()) })), + Some(Arc::new(|_, _| Box::pin(async { Ok(()) }))), ) .await .unwrap(); diff --git a/lib/components/fabro-sandbox/src/sandbox.rs b/lib/components/fabro-sandbox/src/sandbox.rs index 9dac4e499..a7c2fa1f5 100644 --- a/lib/components/fabro-sandbox/src/sandbox.rs +++ b/lib/components/fabro-sandbox/src/sandbox.rs @@ -184,7 +184,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( @@ -753,6 +753,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>, @@ -949,11 +977,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, @@ -961,7 +991,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 @@ -973,25 +1003,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 } /// Launch a long-lived process with bidirectional stdio attached. diff --git a/lib/components/fabro-sandbox/src/test_support.rs b/lib/components/fabro-sandbox/src/test_support.rs index 1593c9032..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, @@ -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,31 @@ 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: Option, + ) -> crate::Result { + let result = self + .exec_command( + command, + timeout_ms.unwrap_or(u64::MAX), + working_dir, + env_vars, + cancel_token, + ) + .await?; + 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 dab7bf4b6..c0c42ebd5 100644 --- a/lib/components/fabro-sandbox/tests/daytona_streaming_live.rs +++ b/lib/components/fabro-sandbox/tests/daytona_streaming_live.rs @@ -335,7 +335,7 @@ mod daytona_streaming_live { None, None, Some(cancel_for_exec), - callback, + Some(callback), ) .await }); @@ -460,7 +460,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 75677216a..f4833b1e7 100644 --- a/lib/components/fabro-sandbox/tests/docker_streaming.rs +++ b/lib/components/fabro-sandbox/tests/docker_streaming.rs @@ -55,7 +55,7 @@ async fn streaming_timeout_terminates_docker_exec_before_returning() { None, None, None, - capture_bytes(Arc::clone(&chunks)), + Some(capture_bytes(Arc::clone(&chunks))), ) .await .expect("streaming command should return a timeout result"); @@ -218,7 +218,7 @@ async fn docker_runs_clean_bash_through_both_command_paths() { None, None, None, - capture_bytes(Arc::clone(&chunks)), + Some(capture_bytes(Arc::clone(&chunks))), ) .await .expect("streaming command should run"); diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index d86c17424..f6e2de0f1 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,47 @@ 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(exec_tail()), + }, + 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"], "last stdout line"); + 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/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-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 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/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 7c8d1621a..b79f7758b 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,68 @@ 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); + } + #[test] fn agent_tool_source_and_category_use_public_json_shape() { assert_eq!(