From ddcdafa06b49098444a117b18d91a5116d0b897b Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 12:34:43 -0400 Subject: [PATCH 1/6] Bound agent tool output capture --- lib/components/fabro-agent/src/lib.rs | 2 +- .../fabro-agent/src/profiles/kimi_tools.rs | 7 + lib/components/fabro-agent/src/sandbox.rs | 2 +- .../fabro-agent/src/tool_execution.rs | 93 ++++++-- lib/components/fabro-agent/src/tools.rs | 12 +- lib/components/fabro-agent/src/truncation.rs | 82 +++++++ .../fabro-sandbox/src/daytona/mod.rs | 88 ++++++- lib/components/fabro-sandbox/src/docker.rs | 23 +- lib/components/fabro-sandbox/src/lib.rs | 4 +- lib/components/fabro-sandbox/src/local.rs | 75 +++++- lib/components/fabro-sandbox/src/sandbox.rs | 216 +++++++++++++++++- .../fabro-sandbox/src/test_support.rs | 9 +- 12 files changed, 548 insertions(+), 65 deletions(-) diff --git a/lib/components/fabro-agent/src/lib.rs b/lib/components/fabro-agent/src/lib.rs index 9f982f47a..f1c9c5f17 100644 --- a/lib/components/fabro-agent/src/lib.rs +++ b/lib/components/fabro-agent/src/lib.rs @@ -60,7 +60,7 @@ pub use question_tools::{ }; pub use sandbox::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingRequest, ExecStreamingResult, - GrepOptions, RefreshOutcome, RemoteCredentialAction, Sandbox, SandboxEvent, + GrepOptions, OutputCaptureStats, RefreshOutcome, RemoteCredentialAction, Sandbox, SandboxEvent, SandboxEventCallback, StderrCollector, StdioProcess, StdioProcessHandle, TokenProvenance, TokenSnapshot, format_lines_numbered, shell_quote, }; diff --git a/lib/components/fabro-agent/src/profiles/kimi_tools.rs b/lib/components/fabro-agent/src/profiles/kimi_tools.rs index a379569a8..6ab61f1dd 100644 --- a/lib/components/fabro-agent/src/profiles/kimi_tools.rs +++ b/lib/components/fabro-agent/src/profiles/kimi_tools.rs @@ -33,6 +33,7 @@ use crate::tools::{ DEFAULT_READ_LINES, emit_shell_process_completed, execute_grep, execute_shell_command, grep_result_path, make_edit_file_tool, optional_usize_arg, required_str, }; +use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output}; const DEFAULT_GREP_RESULTS: usize = 250; const MAX_GREP_RESULTS: usize = 2000; @@ -143,6 +144,12 @@ explicitly asked. Never run commands requiring superuser privileges unless expli let _ = write!(out, "Command failed with exit code: {code}"); } let is_success = result.is_success(); + let out = retain_tool_output( + &out, + MAX_RETAINED_TOOL_OUTPUT_BYTES, + streaming.output_capture().omitted_bytes, + ) + .output; emit_shell_process_completed(&ctx, streaming).await; if is_success { Ok(out) } else { Err(out) } }) diff --git a/lib/components/fabro-agent/src/sandbox.rs b/lib/components/fabro-agent/src/sandbox.rs index 2f1ade3a4..3194dd104 100644 --- a/lib/components/fabro-agent/src/sandbox.rs +++ b/lib/components/fabro-agent/src/sandbox.rs @@ -3,7 +3,7 @@ // `crate::delegate_sandbox!` invocations continue to work. pub use fabro_sandbox::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingRequest, ExecStreamingResult, - GrepOptions, RefreshOutcome, RemoteCredentialAction, Sandbox, SandboxEvent, + GrepOptions, OutputCaptureStats, RefreshOutcome, RemoteCredentialAction, Sandbox, SandboxEvent, SandboxEventCallback, SandboxFile, StderrCollector, StdioProcess, StdioProcessHandle, StdioProcessTermination, TokenProvenance, TokenSnapshot, WalkOptions, delegate_sandbox, format_lines_numbered, shell_quote, diff --git a/lib/components/fabro-agent/src/tool_execution.rs b/lib/components/fabro-agent/src/tool_execution.rs index 2a2c78b1c..228799654 100644 --- a/lib/components/fabro-agent/src/tool_execution.rs +++ b/lib/components/fabro-agent/src/tool_execution.rs @@ -11,7 +11,7 @@ use crate::question_tools::{self, AgentToolRuntime, is_question_tool}; use crate::sandbox::Sandbox; use crate::session::ToolEnvProvider; use crate::tool_registry::{AgentEventEmitter, RegisteredTool, ToolContext, ToolRegistry}; -use crate::truncation::truncate_tool_output; +use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output, truncate_tool_output}; use crate::types::AgentEvent; /// Execute tool calls, choosing parallel or sequential based on `parallel` @@ -265,7 +265,7 @@ fn error_tool_result_with_events( message: &str, ) -> ToolResult { emit_tool_call_started(emitter, session_id, tc); - let result = ToolResult::error(&tc.id, message); + let result = retain_tool_result(&ToolResult::error(&tc.id, message)); emit_tool_call_result(emitter, session_id, tc, &result); truncate_tool_result(&result, &tc.name, config) } @@ -386,7 +386,7 @@ async fn execute_and_emit_one_tool_with_lookup( emit_tool_call_started(emitter, session_id, tc); if let Some(reason) = access_denial { - let result = ToolResult::error(&tc.id, &reason); + let result = retain_tool_result(&ToolResult::error(&tc.id, &reason)); emit_tool_call_result(emitter, session_id, tc, &result); return truncate_tool_result(&result, &tc.name, config); } @@ -400,24 +400,26 @@ async fn execute_and_emit_one_tool_with_lookup( debug!(tool = %tc.name, hook_event = "pre_tool_use", ?decision, duration_ms = elapsed, "Tool hook complete"); if let ToolHookDecision::Block { reason } = decision { - let result = ToolResult::error(&tc.id, &reason); + let result = retain_tool_result(&ToolResult::error(&tc.id, &reason)); emit_tool_call_result(emitter, session_id, tc, &result); return truncate_tool_result(&result, &tc.name, config); } } - let result = execute_one_tool( - tc, - registered_tool, - env, - cancel_token, - emitter, - session_id, - root_session_id, - tool_env_provider, - agent_tool_runtime, - ) - .await; + let result = retain_tool_result( + &execute_one_tool( + tc, + registered_tool, + env, + cancel_token, + emitter, + session_id, + root_session_id, + tool_env_provider, + agent_tool_runtime, + ) + .await, + ); emit_tool_call_result(emitter, session_id, tc, &result); @@ -446,6 +448,24 @@ async fn execute_and_emit_one_tool_with_lookup( truncate_tool_result(&result, &tc.name, config) } +/// Bound model-native tool output before it reaches hooks, events, or history. +fn retain_tool_result(result: &ToolResult) -> ToolResult { + let retained_content = match &result.content { + serde_json::Value::String(output) => serde_json::Value::String( + retain_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, 0).output, + ), + other => other.clone(), + }; + + ToolResult { + tool_call_id: result.tool_call_id.clone(), + content: retained_content, + is_error: result.is_error, + image_data: result.image_data.clone(), + image_media_type: result.image_media_type.clone(), + } +} + /// Execute a single tool call: argument validation and execution. #[allow( clippy::too_many_arguments, @@ -877,6 +897,47 @@ mod tests { assert!(content.contains("echo: hello")); } + #[tokio::test] + async fn tool_output_is_bounded_before_events_and_history() { + let mut registry = ToolRegistry::new(); + registry.register(make_echo_tool()); + let text = "x".repeat(MAX_RETAINED_TOOL_OUTPUT_BYTES + 100); + let tc = make_tool_call("echo", "call_large", serde_json::json!({"text": text})); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); + + let result = execute_and_emit_one_tool( + &tc, + ®istry, + make_sandbox(), + None, + CancellationToken::new(), + &SessionOptions::default(), + &emitter, + "test-session", + "test-session", + None, + ) + .await; + + let result_output = result.content.as_str().expect("string tool output"); + assert_eq!(result_output.len(), MAX_RETAINED_TOOL_OUTPUT_BYTES); + + let completed_output = loop { + let event = receiver.try_recv().expect("tool completion event"); + if let AgentEvent::ToolCallCompleted { output, .. } = event.event { + break output; + } + }; + assert_eq!( + completed_output + .as_str() + .expect("string event output") + .len(), + MAX_RETAINED_TOOL_OUTPUT_BYTES + ); + } + #[tokio::test] async fn post_tool_use_hook_fires_on_success() { let mut registry = ToolRegistry::new(); diff --git a/lib/components/fabro-agent/src/tools.rs b/lib/components/fabro-agent/src/tools.rs index 6d9f673e1..d4d429146 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -13,6 +13,7 @@ use tokio::task; use crate::config::NativeToolOptions; use crate::sandbox::{ExecStreamingResult, GrepOptions}; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource}; +use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, RetainedToolOutput, retain_tool_output}; use crate::types::AgentEvent; use crate::web_search::{SearchBackend, make_web_search_tool}; @@ -303,6 +304,7 @@ pub(crate) async fn execute_shell_command( working_dir: cwd, env_vars: tool_env.as_ref(), cancel_token: Some(ctx.cancel.clone()), + stream_output_bytes_cap: Some(MAX_RETAINED_TOOL_OUTPUT_BYTES), ..crate::ExecStreamingRequest::new(command) }) .await @@ -318,7 +320,7 @@ pub(crate) async fn run_shell_command( cwd: Option<&str>, ) -> Result { let streaming = execute_shell_command(ctx, command, timeout_ms, cwd).await?; - let text = render_shell_result(&streaming); + let text = render_shell_result(&streaming).output; let is_success = streaming.result.is_success(); emit_shell_process_completed(ctx, streaming).await; @@ -364,7 +366,7 @@ pub(crate) async fn emit_shell_process_completed( /// 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 { +fn render_shell_result(streaming: &ExecStreamingResult) -> RetainedToolOutput { let result = &streaming.result; let mut output = format!( "Termination: {}\nExit code: {}\nDuration: {}ms\n", @@ -384,7 +386,11 @@ fn render_shell_result(streaming: &ExecStreamingResult) -> String { } else if !result.stdout.is_empty() { let _ = write!(output, "output (combined):\n{}\n", result.stdout); } - output + retain_tool_output( + &output, + MAX_RETAINED_TOOL_OUTPUT_BYTES, + streaming.output_capture().omitted_bytes, + ) } #[must_use] diff --git a/lib/components/fabro-agent/src/truncation.rs b/lib/components/fabro-agent/src/truncation.rs index d7fbe3ffc..edc3122b3 100644 --- a/lib/components/fabro-agent/src/truncation.rs +++ b/lib/components/fabro-agent/src/truncation.rs @@ -1,6 +1,64 @@ use crate::config::SessionOptions; +use crate::sandbox::OutputCaptureStats; use crate::tool_permissions::canonical_tool_name; +pub(crate) const MAX_RETAINED_TOOL_OUTPUT_BYTES: usize = 1024 * 1024; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct RetainedToolOutput { + pub output: String, + pub stats: OutputCaptureStats, +} + +/// Keep an equal-sized UTF-8 prefix and suffix within a byte budget. +/// +/// `previously_omitted_bytes` accounts for output a streaming provider +/// discarded before the rendered result was assembled. +#[must_use] +pub(crate) fn retain_tool_output( + output: &str, + max_bytes: usize, + previously_omitted_bytes: usize, +) -> RetainedToolOutput { + let observed_bytes = output.len().saturating_add(previously_omitted_bytes); + if output.len() <= max_bytes { + return RetainedToolOutput { + output: output.to_string(), + stats: OutputCaptureStats { + observed_bytes, + retained_bytes: output.len(), + omitted_bytes: previously_omitted_bytes, + }, + }; + } + + let head_budget = max_bytes / 2; + let tail_budget = max_bytes.saturating_sub(head_budget); + let head_end = output.floor_char_boundary(head_budget); + let tail_start = ceil_char_boundary(output, output.len().saturating_sub(tail_budget)); + let retained_bytes = head_end.saturating_add(output.len().saturating_sub(tail_start)); + let mut retained = String::with_capacity(retained_bytes); + retained.push_str(&output[..head_end]); + retained.push_str(&output[tail_start..]); + + RetainedToolOutput { + output: retained, + stats: OutputCaptureStats { + observed_bytes, + retained_bytes, + omitted_bytes: observed_bytes.saturating_sub(retained_bytes), + }, + } +} + +fn ceil_char_boundary(output: &str, index: usize) -> usize { + let mut index = index.min(output.len()); + while index < output.len() && !output.is_char_boundary(index) { + index += 1; + } + index +} + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum TruncationMode { HeadTail, @@ -121,6 +179,30 @@ pub fn truncate_tool_output(output: &str, tool_name: &str, config: &SessionOptio mod tests { use super::*; + #[test] + fn retained_tool_output_keeps_equal_head_and_tail() { + let retained = retain_tool_output("abcdefghijkl", 8, 0); + + assert_eq!(retained.output, "abcdijkl"); + assert_eq!(retained.stats.observed_bytes, 12); + assert_eq!(retained.stats.retained_bytes, 8); + assert_eq!(retained.stats.omitted_bytes, 4); + } + + #[test] + fn retained_tool_output_stays_within_budget_at_utf8_boundaries() { + let retained = retain_tool_output("aa😀😀zz", 7, 3); + + assert!(retained.output.len() <= 7, "{}", retained.output.len()); + assert!(retained.output.starts_with("aa")); + assert!(retained.output.ends_with("zz")); + assert_eq!(retained.stats.observed_bytes, "aa😀😀zz".len() + 3); + assert_eq!( + retained.stats.omitted_bytes, + retained.stats.observed_bytes - retained.output.len() + ); + } + #[test] fn under_limit_passthrough_chars() { let output = "short output"; diff --git a/lib/components/fabro-sandbox/src/daytona/mod.rs b/lib/components/fabro-sandbox/src/daytona/mod.rs index 14af37bd7..6da9c6e72 100644 --- a/lib/components/fabro-sandbox/src/daytona/mod.rs +++ b/lib/components/fabro-sandbox/src/daytona/mod.rs @@ -32,8 +32,9 @@ use crate::git_retry::{self, CredentialContext, GitRetryReason}; use crate::push_credentials::{self, PushCredentialState}; use crate::redact::redact_auth_url; use crate::sandbox::{ - self, BASH_ENV_VAR, BASH_PROBE_MARKER, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH, - REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, optional_timeout, resolve_path, validate_bash_probe, + self, BASH_ENV_VAR, BASH_PROBE_MARKER, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, + OutputCaptureBuffer, REMOTE_BASH, REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, optional_timeout, + resolve_path, validate_bash_probe, }; use crate::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingRequest, ExecStreamingResult, @@ -2212,6 +2213,7 @@ impl Sandbox for DaytonaSandbox { cancel_token, stdin, output_callback, + stream_output_bytes_cap, } = request; let sandbox = self.sandbox()?; let start = Instant::now(); @@ -2252,8 +2254,12 @@ impl Sandbox for DaytonaSandbox { return Err(crate::Error::context("Failed to get process service", err)); } }; - let stdout_seen = Arc::new(Mutex::new(Vec::new())); - let stderr_seen = Arc::new(Mutex::new(Vec::new())); + let stdout_seen = Arc::new(Mutex::new(OutputCaptureBuffer::new( + stream_output_bytes_cap, + ))); + let stderr_seen = Arc::new(Mutex::new(OutputCaptureBuffer::new( + stream_output_bytes_cap, + ))); let saw_live_chunk = Arc::new(AtomicBool::new(false)); let stream_session_id = session.id().to_string(); @@ -2277,7 +2283,7 @@ impl Sandbox for DaytonaSandbox { let bytes = chunk.into_bytes(); if !bytes.is_empty() { saw_live_chunk.store(true, Ordering::Relaxed); - stdout_seen.lock().await.extend_from_slice(&bytes); + stdout_seen.lock().await.push(&bytes); if let Some(callback) = callback { callback(CommandOutputStream::Stdout, bytes) .await @@ -2295,7 +2301,7 @@ impl Sandbox for DaytonaSandbox { let bytes = chunk.into_bytes(); if !bytes.is_empty() { saw_live_chunk.store(true, Ordering::Relaxed); - stderr_seen.lock().await.extend_from_slice(&bytes); + stderr_seen.lock().await.push(&bytes); if let Some(callback) = callback { callback(CommandOutputStream::Stderr, bytes) .await @@ -2369,8 +2375,20 @@ impl Sandbox for DaytonaSandbox { .await?; } - let stdout = String::from_utf8_lossy(&stdout_seen.lock().await).into_owned(); - let stderr = String::from_utf8_lossy(&stderr_seen.lock().await).into_owned(); + let (stdout, stdout_capture) = { + let stdout_seen = stdout_seen.lock().await; + ( + String::from_utf8_lossy(&stdout_seen.to_bytes()).into_owned(), + stdout_seen.stats(), + ) + }; + let (stderr, stderr_capture) = { + let stderr_seen = stderr_seen.lock().await; + ( + String::from_utf8_lossy(&stderr_seen.to_bytes()).into_owned(), + stderr_seen.stats(), + ) + }; let result = ExecStreamingResult { result: ExecResult { @@ -2384,6 +2402,8 @@ impl Sandbox for DaytonaSandbox { }, streams_separated, live_streaming: saw_live_chunk.load(Ordering::Relaxed), + stdout_capture, + stderr_capture, }; if let Some(stdin_file) = stdin_file.as_mut() { stdin_file.close().await; @@ -2886,7 +2906,7 @@ async fn fetch_daytona_session_logs( async fn append_missing_log_suffix( stream: CommandOutputStream, final_bytes: &[u8], - seen: &Arc>>, + seen: &Arc>, output_callback: Option<&CommandOutputCallback>, ) -> crate::Result<()> { if final_bytes.is_empty() { @@ -2894,13 +2914,13 @@ async fn append_missing_log_suffix( } let mut seen = seen.lock().await; - let offset = missing_log_suffix_offset(&seen, final_bytes); + let offset = captured_log_suffix_offset(&seen, final_bytes); if offset >= final_bytes.len() { return Ok(()); } let missing = final_bytes[offset..].to_vec(); - seen.extend_from_slice(&missing); + seen.push(&missing); drop(seen); match output_callback { Some(output_callback) => output_callback(stream, missing).await, @@ -2908,6 +2928,42 @@ async fn append_missing_log_suffix( } } +fn captured_log_suffix_offset(seen: &OutputCaptureBuffer, final_bytes: &[u8]) -> usize { + let stats = seen.stats(); + if stats.omitted_bytes == 0 { + return missing_log_suffix_offset(&seen.to_bytes(), final_bytes); + } + + let observed_bytes = seen.observed_bytes(); + let head = seen.retained_head(); + let tail = seen.retained_tail(); + if final_bytes.len() >= observed_bytes + && final_bytes.starts_with(head) + && tail.iter().copied().eq(final_bytes + [observed_bytes.saturating_sub(tail.len())..observed_bytes] + .iter() + .copied()) + { + return observed_bytes; + } + if final_bytes.len() <= observed_bytes && final_bytes.starts_with(head) { + return final_bytes.len(); + } + + let max_overlap = tail.len().min(final_bytes.len()); + for overlap in (1..=max_overlap).rev() { + if tail + .iter() + .skip(tail.len() - overlap) + .copied() + .eq(final_bytes[..overlap].iter().copied()) + { + return overlap; + } + } + 0 +} + fn missing_log_suffix_offset(seen: &[u8], final_bytes: &[u8]) -> usize { if final_bytes.starts_with(seen) { return seen.len(); @@ -4730,6 +4786,16 @@ mod tests { assert_eq!(missing_log_suffix_offset(b"abc", b"def"), 0); } + #[test] + fn captured_log_suffix_offset_uses_observed_length_after_truncation() { + let mut seen = OutputCaptureBuffer::new(Some(6)); + seen.push(b"abcdefgh"); + + assert_eq!(captured_log_suffix_offset(&seen, b"abcdefghij"), 8); + assert_eq!(captured_log_suffix_offset(&seen, b"abcdefgh"), 8); + assert_eq!(captured_log_suffix_offset(&seen, b"abcd"), 4); + } + #[test] fn detect_git_remote_from_repo() { let dir = tempfile::tempdir().unwrap(); diff --git a/lib/components/fabro-sandbox/src/docker.rs b/lib/components/fabro-sandbox/src/docker.rs index 83b0dc362..0b08b97b4 100644 --- a/lib/components/fabro-sandbox/src/docker.rs +++ b/lib/components/fabro-sandbox/src/docker.rs @@ -33,7 +33,7 @@ use crate::managed_labels::{self, MANAGED_LABEL, RUN_ID_LABEL}; use crate::push_credentials::{self, PushCredentialState}; use crate::redact::redact_auth_url; use crate::sandbox::{ - self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH, + self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, OutputCaptureBuffer, REMOTE_BASH, REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, StdioProcessControl, optional_timeout, resolve_path, validate_bash_probe, write_process_stdin, }; @@ -438,7 +438,8 @@ impl DockerSandbox { env: Option>, stdin: Option>, output_callback: Option, - ) -> crate::Result<(Vec, Vec, i32)> { + stream_output_bytes_cap: Option, + ) -> crate::Result<(OutputCaptureBuffer, OutputCaptureBuffer, i32)> { let exec_opts = CreateExecOptions { cmd: Some(cmd), attach_stdin: Some(stdin.is_some()), @@ -460,8 +461,8 @@ impl DockerSandbox { ) .await?; - let mut stdout = Vec::new(); - let mut stderr = Vec::new(); + let mut stdout = OutputCaptureBuffer::new(stream_output_bytes_cap); + let mut stderr = OutputCaptureBuffer::new(stream_output_bytes_cap); let StartExecResults::Attached { mut output, input } = start_result else { return Err(crate::Error::message( @@ -478,13 +479,13 @@ impl DockerSandbox { while let Some(chunk) = output.next().await { match chunk { Ok(LogOutput::StdOut { message }) => { - stdout.extend_from_slice(&message); + stdout.push(&message); 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); + stderr.push(&message); if let Some(output_callback) = output_callback.as_ref() { output_callback(CommandOutputStream::Stderr, message.to_vec()).await?; } @@ -580,6 +581,7 @@ impl DockerSandbox { cancel_token, stdin, output_callback, + stream_output_bytes_cap, } = request; let start = Instant::now(); let effective_dir = working_dir @@ -607,6 +609,7 @@ impl DockerSandbox { env, stdin, output_callback, + stream_output_bytes_cap, )); let mut termination = CommandTermination::Exited; @@ -632,9 +635,11 @@ impl DockerSandbox { }; let (stdout, stderr, exit_code) = output; + let (stdout, stdout_capture) = stdout.into_parts(); + let (stderr, stderr_capture) = stderr.into_parts(); let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX); Ok(ExecStreamingResult { - result: ExecResult { + result: ExecResult { stdout: String::from_utf8_lossy(&stdout).into_owned(), stderr: String::from_utf8_lossy(&stderr).into_owned(), exit_code: (termination == CommandTermination::Exited).then_some(exit_code), @@ -642,7 +647,9 @@ impl DockerSandbox { duration_ms, }, streams_separated: true, - live_streaming: true, + live_streaming: true, + stdout_capture, + stderr_capture, }) } diff --git a/lib/components/fabro-sandbox/src/lib.rs b/lib/components/fabro-sandbox/src/lib.rs index 3594aa0c0..988eacf4f 100644 --- a/lib/components/fabro-sandbox/src/lib.rs +++ b/lib/components/fabro-sandbox/src/lib.rs @@ -60,8 +60,8 @@ pub use reconnect::{reconnect, reconnect_for_run, reconnect_for_run_with_callbac pub use sandbox::{ CommandOutputCallback, DEFAULT_EXEC_OUTPUT_TAIL_BYTES, DirEntry, ExecResult, ExecStreamingRequest, ExecStreamingResult, GitRunInfo, GitSetupIntent, GrepOptions, - PushAttempt, PushError, PushReport, RefreshOutcome, RemoteCredentialAction, Sandbox, - SandboxEvent, SandboxEventCallback, SandboxFile, StderrCollector, StdioProcess, + OutputCaptureStats, PushAttempt, PushError, PushReport, RefreshOutcome, RemoteCredentialAction, + Sandbox, SandboxEvent, SandboxEventCallback, SandboxFile, StderrCollector, StdioProcess, StdioProcessHandle, StdioProcessTermination, WalkOptions, format_lines_numbered, redacted_output_tail, setup_git_via_exec, shell_quote, }; diff --git a/lib/components/fabro-sandbox/src/local.rs b/lib/components/fabro-sandbox/src/local.rs index 650d57c47..d2174af6c 100644 --- a/lib/components/fabro-sandbox/src/local.rs +++ b/lib/components/fabro-sandbox/src/local.rs @@ -13,8 +13,8 @@ use tokio::{fs, time}; use tokio_util::sync::CancellationToken; use crate::sandbox::{ - self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, StdioProcessControl, - optional_timeout, validate_bash_probe, write_process_stdin, + self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, OutputCaptureBuffer, + StdioProcessControl, optional_timeout, validate_bash_probe, write_process_stdin, }; use crate::{ CommandOutputCallback, DEFAULT_EXEC_OUTPUT_TAIL_BYTES, DirEntry, ExecResult, @@ -525,6 +525,7 @@ impl Sandbox for LocalSandbox { cancel_token, stdin, output_callback, + stream_output_bytes_cap, } = request; let start = Instant::now(); @@ -573,10 +574,22 @@ impl Sandbox for LocalSandbox { let stdout_callback = output_callback.clone(); let stderr_callback = output_callback; let stdout_task = tokio::spawn(async move { - drain_command_pipe(stdout_pipe, CommandOutputStream::Stdout, stdout_callback).await + drain_command_pipe( + stdout_pipe, + CommandOutputStream::Stdout, + stdout_callback, + stream_output_bytes_cap, + ) + .await }); let stderr_task = tokio::spawn(async move { - drain_command_pipe(stderr_pipe, CommandOutputStream::Stderr, stderr_callback).await + drain_command_pipe( + stderr_pipe, + CommandOutputStream::Stderr, + stderr_callback, + stream_output_bytes_cap, + ) + .await }); let (termination, exit_code) = tokio::select! { @@ -613,15 +626,17 @@ impl Sandbox for LocalSandbox { } } } - let stdout_bytes = stdout_task + let stdout_capture = stdout_task .await .map_err(|e| crate::Error::context("stdout stream task failed", e))??; - let stderr_bytes = stderr_task + let stderr_capture = stderr_task .await .map_err(|e| crate::Error::context("stderr stream task failed", e))??; + let (stdout_bytes, stdout_capture) = stdout_capture.into_parts(); + let (stderr_bytes, stderr_capture) = stderr_capture.into_parts(); Ok(ExecStreamingResult { - result: ExecResult { + result: ExecResult { stdout: String::from_utf8_lossy(&stdout_bytes).into_owned(), stderr: String::from_utf8_lossy(&stderr_bytes).into_owned(), exit_code, @@ -629,7 +644,9 @@ impl Sandbox for LocalSandbox { duration_ms, }, streams_separated: true, - live_streaming: true, + live_streaming: true, + stdout_capture, + stderr_capture, }) } @@ -1002,11 +1019,12 @@ async fn drain_command_pipe( mut reader: Option, stream: CommandOutputStream, output_callback: Option, -) -> crate::Result> + stream_output_bytes_cap: Option, +) -> crate::Result where R: AsyncRead + Unpin, { - let mut output = Vec::new(); + let mut output = OutputCaptureBuffer::new(stream_output_bytes_cap); let Some(reader) = reader.as_mut() else { return Ok(output); }; @@ -1020,7 +1038,7 @@ where if read == 0 { return Ok(output); } - output.extend_from_slice(&buf[..read]); + output.push(&buf[..read]); if let Some(output_callback) = output_callback.as_ref() { output_callback(stream, buf[..read].to_vec()).await?; } @@ -1171,7 +1189,7 @@ mod tests { use std::io; use std::path::PathBuf; use std::pin::Pin; - use std::sync::Arc; + use std::sync::{Arc, Mutex}; use std::task::{Context as TaskContext, Poll}; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader, ReadBuf}; @@ -1408,6 +1426,39 @@ mod tests { std::fs::remove_dir_all(&dir).unwrap(); } + #[tokio::test] + async fn exec_command_streaming_drains_full_output_after_capture_cap() { + let dir = temp_dir(); + let sandbox = LocalSandbox::new(dir.clone()); + let callback_bytes = Arc::new(Mutex::new(Vec::new())); + let callback_bytes_for_stream = Arc::clone(&callback_bytes); + + let result = sandbox + .exec_command_streaming(ExecStreamingRequest { + timeout_ms: Some(5000), + output_callback: Some(Arc::new(move |stream, bytes| { + let callback_bytes = Arc::clone(&callback_bytes_for_stream); + Box::pin(async move { + if stream == CommandOutputStream::Stdout { + callback_bytes.lock().unwrap().extend(bytes); + } + Ok(()) + }) + })), + stream_output_bytes_cap: Some(8), + ..ExecStreamingRequest::new("printf 'abcdefghijklmnopqrst'") + }) + .await + .unwrap(); + + assert_eq!(result.result.stdout, "abcdqrst"); + assert_eq!(result.stdout_capture.observed_bytes, 20); + assert_eq!(result.stdout_capture.retained_bytes, 8); + assert_eq!(result.stdout_capture.omitted_bytes, 12); + assert_eq!(&*callback_bytes.lock().unwrap(), b"abcdefghijklmnopqrst"); + std::fs::remove_dir_all(&dir).unwrap(); + } + #[tokio::test] async fn exec_command_streaming_writes_exact_stdin_and_closes_it() { let dir = temp_dir(); diff --git a/lib/components/fabro-sandbox/src/sandbox.rs b/lib/components/fabro-sandbox/src/sandbox.rs index f9d153b2e..473cfa71b 100644 --- a/lib/components/fabro-sandbox/src/sandbox.rs +++ b/lib/components/fabro-sandbox/src/sandbox.rs @@ -1,4 +1,4 @@ -use std::collections::HashMap; +use std::collections::{HashMap, VecDeque}; use std::fmt::Write; use std::future::Future; use std::path::Path; @@ -766,6 +766,143 @@ pub struct ExecStreamingResult { pub result: ExecResult, pub streams_separated: bool, pub live_streaming: bool, + pub stdout_capture: OutputCaptureStats, + pub stderr_capture: OutputCaptureStats, +} + +impl ExecStreamingResult { + #[must_use] + pub fn output_capture(&self) -> OutputCaptureStats { + self.stdout_capture.combine(self.stderr_capture) + } +} + +/// Byte counts for output observed and retained while draining a process. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct OutputCaptureStats { + pub observed_bytes: usize, + pub retained_bytes: usize, + pub omitted_bytes: usize, +} + +impl OutputCaptureStats { + #[must_use] + pub fn complete(byte_count: usize) -> Self { + Self { + observed_bytes: byte_count, + retained_bytes: byte_count, + omitted_bytes: 0, + } + } + + #[must_use] + pub fn combine(self, other: Self) -> Self { + Self { + observed_bytes: self.observed_bytes.saturating_add(other.observed_bytes), + retained_bytes: self.retained_bytes.saturating_add(other.retained_bytes), + omitted_bytes: self.omitted_bytes.saturating_add(other.omitted_bytes), + } + } +} + +/// A byte buffer that keeps an equal-sized stable prefix and rolling suffix. +/// +/// The process is always drained. Once the optional cap is full, bytes from +/// the middle are discarded while the newest suffix replaces the old tail. +#[derive(Debug)] +pub(crate) struct OutputCaptureBuffer { + max_bytes: Option, + head: Vec, + tail: VecDeque, + observed_bytes: usize, +} + +impl OutputCaptureBuffer { + #[must_use] + pub(crate) fn new(max_bytes: Option) -> Self { + Self { + max_bytes, + head: Vec::new(), + tail: VecDeque::new(), + observed_bytes: 0, + } + } + + pub(crate) fn push(&mut self, bytes: &[u8]) { + self.observed_bytes = self.observed_bytes.saturating_add(bytes.len()); + let Some(max_bytes) = self.max_bytes else { + self.head.extend_from_slice(bytes); + return; + }; + + let head_budget = max_bytes / 2; + let tail_budget = max_bytes.saturating_sub(head_budget); + let head_remaining = head_budget.saturating_sub(self.head.len()); + let head_take = head_remaining.min(bytes.len()); + self.head.extend_from_slice(&bytes[..head_take]); + + let tail_bytes = &bytes[head_take..]; + let overflow = self + .tail + .len() + .saturating_add(tail_bytes.len()) + .saturating_sub(tail_budget); + if overflow >= self.tail.len() { + let skip = overflow.saturating_sub(self.tail.len()); + self.tail.clear(); + self.tail.extend(&tail_bytes[skip..]); + } else { + self.tail.drain(..overflow); + self.tail.extend(tail_bytes); + } + } + + #[must_use] + pub(crate) fn stats(&self) -> OutputCaptureStats { + let retained_bytes = self.head.len().saturating_add(self.tail.len()); + OutputCaptureStats { + observed_bytes: self.observed_bytes, + retained_bytes, + omitted_bytes: self.observed_bytes.saturating_sub(retained_bytes), + } + } + + #[cfg(feature = "daytona")] + #[must_use] + pub(crate) fn to_bytes(&self) -> Vec { + let stats = self.stats(); + let mut bytes = Vec::with_capacity(stats.retained_bytes); + bytes.extend_from_slice(&self.head); + bytes.extend(self.tail.iter().copied()); + bytes + } + + #[must_use] + pub(crate) fn into_parts(self) -> (Vec, OutputCaptureStats) { + let stats = self.stats(); + let mut bytes = Vec::with_capacity(stats.retained_bytes); + bytes.extend(self.head); + bytes.extend(self.tail); + (bytes, stats) + } + + #[cfg(feature = "daytona")] + #[must_use] + pub(crate) fn observed_bytes(&self) -> usize { + self.observed_bytes + } + + #[cfg(feature = "daytona")] + #[must_use] + pub(crate) fn retained_head(&self) -> &[u8] { + &self.head + } + + #[cfg(feature = "daytona")] + #[must_use] + pub(crate) fn retained_tail(&self) -> &VecDeque { + &self.tail + } } pub type CommandOutputCallback = Arc< @@ -785,13 +922,16 @@ pub type CommandOutputCallback = Arc< /// type does not implement `Debug` because standard input can contain /// sensitive workflow data. pub struct ExecStreamingRequest<'a> { - pub command: &'a str, - pub timeout_ms: Option, - pub working_dir: Option<&'a str>, - pub env_vars: Option<&'a HashMap>, - pub cancel_token: Option, - pub stdin: Option>, - pub output_callback: Option, + pub command: &'a str, + pub timeout_ms: Option, + pub working_dir: Option<&'a str>, + pub env_vars: Option<&'a HashMap>, + pub cancel_token: Option, + pub stdin: Option>, + pub output_callback: Option, + /// Maximum bytes retained from each stream. Providers continue draining + /// stdout and stderr after the cap is reached. + pub stream_output_bytes_cap: Option, } impl<'a> ExecStreamingRequest<'a> { @@ -805,6 +945,7 @@ impl<'a> ExecStreamingRequest<'a> { cancel_token: None, stdin: None, output_callback: None, + stream_output_bytes_cap: None, } } } @@ -846,9 +987,10 @@ where } pub(crate) async fn replay_exec_result( - result: ExecResult, + mut result: ExecResult, streams_separated: bool, output_callback: Option<&CommandOutputCallback>, + stream_output_bytes_cap: Option, ) -> crate::Result { if let Some(output_callback) = output_callback { if !result.stdout.is_empty() { @@ -866,10 +1008,21 @@ pub(crate) async fn replay_exec_result( .await?; } } + let mut stdout_capture = OutputCaptureBuffer::new(stream_output_bytes_cap); + stdout_capture.push(result.stdout.as_bytes()); + let mut stderr_capture = OutputCaptureBuffer::new(stream_output_bytes_cap); + stderr_capture.push(result.stderr.as_bytes()); + let (stdout, stdout_capture) = stdout_capture.into_parts(); + let (stderr, stderr_capture) = stderr_capture.into_parts(); + result.stdout = String::from_utf8_lossy(&stdout).into_owned(); + result.stderr = String::from_utf8_lossy(&stderr).into_owned(); + Ok(ExecStreamingResult { result, streams_separated, live_streaming: false, + stdout_capture, + stderr_capture, }) } @@ -1179,7 +1332,13 @@ pub trait Sandbox: Send + Sync { request.cancel_token, ) .await?; - replay_exec_result(result, true, request.output_callback.as_ref()).await + replay_exec_result( + result, + true, + request.output_callback.as_ref(), + request.stream_output_bytes_cap, + ) + .await } /// Launch a long-lived process with bidirectional stdio attached. @@ -2517,6 +2676,43 @@ mod push_tests { mod tests { use super::*; + #[test] + fn output_capture_buffer_keeps_stable_head_and_rolling_tail() { + let mut buffer = OutputCaptureBuffer::new(Some(8)); + buffer.push(b"abc"); + buffer.push(b"defghi"); + buffer.push(b"jkl"); + + let (bytes, stats) = buffer.into_parts(); + assert_eq!(bytes, b"abcdijkl"); + assert_eq!(stats.observed_bytes, 12); + assert_eq!(stats.retained_bytes, 8); + assert_eq!(stats.omitted_bytes, 4); + } + + #[test] + fn output_capture_buffer_without_cap_retains_everything() { + let mut buffer = OutputCaptureBuffer::new(None); + buffer.push(b"abc"); + buffer.push(b"def"); + + let (bytes, stats) = buffer.into_parts(); + assert_eq!(bytes, b"abcdef"); + assert_eq!(stats, OutputCaptureStats::complete(6)); + } + + #[test] + fn zero_byte_output_capture_buffer_still_counts_drained_bytes() { + let mut buffer = OutputCaptureBuffer::new(Some(0)); + buffer.push(b"abcdef"); + + let (bytes, stats) = buffer.into_parts(); + assert!(bytes.is_empty()); + assert_eq!(stats.observed_bytes, 6); + assert_eq!(stats.retained_bytes, 0); + assert_eq!(stats.omitted_bytes, 6); + } + #[test] fn exec_result_fields() { let result = ExecResult { diff --git a/lib/components/fabro-sandbox/src/test_support.rs b/lib/components/fabro-sandbox/src/test_support.rs index d25d9bbf4..3b97d4b4c 100644 --- a/lib/components/fabro-sandbox/src/test_support.rs +++ b/lib/components/fabro-sandbox/src/test_support.rs @@ -324,6 +324,7 @@ impl Sandbox for MockSandbox { cancel_token, stdin, output_callback, + stream_output_bytes_cap, } = request; *self .captured_stdin @@ -338,7 +339,13 @@ impl Sandbox for MockSandbox { cancel_token, ) .await?; - sandbox::replay_exec_result(result, self.streams_separated, output_callback.as_ref()).await + sandbox::replay_exec_result( + result, + self.streams_separated, + output_callback.as_ref(), + stream_output_bytes_cap, + ) + .await } async fn spawn_stdio_process( From 401acb6cdff13c1c6aec67153e108e743a23bef3 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 12:46:27 -0400 Subject: [PATCH 2/6] Record tool output byte counts --- .../src/commands/run/run_progress/mod.rs | 22 +- lib/apps/fabro-server/src/demo/mod.rs | 6 + .../src/server/handler/sessions.rs | 10 + lib/components/fabro-agent/src/event.rs | 35 ++- .../fabro-agent/src/profiles/kimi_tools.rs | 7 +- .../fabro-agent/src/tool_execution.rs | 263 ++++++++++++++---- .../fabro-agent/src/tool_registry.rs | 13 +- lib/components/fabro-agent/src/tools.rs | 23 +- lib/components/fabro-agent/src/types.rs | 57 +++- .../fabro-agent/tests/it/docker_shell.rs | 10 +- .../fabro-store/src/run_sessions.rs | 3 + lib/components/fabro-store/src/run_state.rs | 17 +- .../fabro-workflow/src/event/convert.rs | 32 ++- .../fabro-workflow/src/event/redaction.rs | 13 +- .../fabro-workflow/src/handler/llm/api.rs | 44 +-- .../fabro-types/src/run_event/agent.rs | 67 +++-- .../fabro-types/src/run_event/mod.rs | 24 +- .../fabro-types/src/run_event/session.rs | 35 ++- 18 files changed, 522 insertions(+), 159 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index 53d336d49..62bb37cd6 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -1029,10 +1029,13 @@ mod tests { emit( &mut ui, agent_event("plan", AgentEvent::ToolCallCompleted { - tool_name: "read_file".into(), - tool_call_id: "tc1".into(), - output: serde_json::json!({"ok": true}), - is_error: false, + tool_name: "read_file".into(), + tool_call_id: "tc1".into(), + output: serde_json::json!({"ok": true}), + is_error: false, + output_bytes_observed: 11, + output_bytes_retained: 11, + output_bytes_omitted: 0, }), ); emit(&mut ui, stage_completed("plan", "Plan")); @@ -1570,10 +1573,13 @@ mod tests { let tool_completed = serde_json::to_string(&to_run_event_at( &fixtures::RUN_1, &agent_event("code", AgentEvent::ToolCallCompleted { - tool_name: "read_file".into(), - tool_call_id: "tc1".into(), - output: serde_json::json!({"ok": true}), - is_error: false, + tool_name: "read_file".into(), + tool_call_id: "tc1".into(), + output: serde_json::json!({"ok": true}), + is_error: false, + output_bytes_observed: 11, + output_bytes_retained: 11, + output_bytes_omitted: 0, }), completed_ts, None, diff --git a/lib/apps/fabro-server/src/demo/mod.rs b/lib/apps/fabro-server/src/demo/mod.rs index 18c51e8a4..7a8b1b347 100644 --- a/lib/apps/fabro-server/src/demo/mod.rs +++ b/lib/apps/fabro-server/src/demo/mod.rs @@ -1531,6 +1531,9 @@ mod runs { output: serde_json::json!("[redis]\nhost = \"redis-prod.internal\"\nport = 6379"), is_error: false, visit: 1, + output_bytes_observed: None, + output_bytes_retained: None, + output_bytes_omitted: None, tool_result: None, turn_id: None, }), @@ -1557,6 +1560,9 @@ mod runs { output: serde_json::json!("[redis]\nhost = \"redis-staging.internal\"\nport = 6379"), is_error: false, visit: 1, + output_bytes_observed: None, + output_bytes_retained: None, + output_bytes_omitted: None, tool_result: None, turn_id: None, }), diff --git a/lib/apps/fabro-server/src/server/handler/sessions.rs b/lib/apps/fabro-server/src/server/handler/sessions.rs index 835b9ed51..8475b01bf 100644 --- a/lib/apps/fabro-server/src/server/handler/sessions.rs +++ b/lib/apps/fabro-server/src/server/handler/sessions.rs @@ -1271,6 +1271,9 @@ fn agent_event_payload(event_turn_id: TurnId, event: AgentEvent) -> Option Some(EventBody::RunSessionToolCallCompleted( RunSessionToolCallCompletedProps { turn_id: event_turn_id, @@ -1278,6 +1281,13 @@ fn agent_event_payload(event_turn_id: TurnId, event: AgentEvent) -> Option None, diff --git a/lib/components/fabro-agent/src/event.rs b/lib/components/fabro-agent/src/event.rs index f2dafe507..eaea439bd 100644 --- a/lib/components/fabro-agent/src/event.rs +++ b/lib/components/fabro-agent/src/event.rs @@ -1,7 +1,9 @@ +use std::sync::{Arc, Mutex}; use std::time::SystemTime; use tokio::sync::broadcast; +use crate::sandbox::OutputCaptureStats; use crate::tool_registry::AgentEventEmitter; use crate::types::{AgentEvent, SessionEvent}; @@ -62,9 +64,29 @@ impl Default for Emitter { /// when a subagent's events are forwarded through its parent. #[derive(Clone)] pub struct SessionBoundEmitter { - pub emitter: Emitter, - pub session_id: String, - pub tool_call_id: Option, + pub emitter: Emitter, + pub session_id: String, + pub tool_call_id: Option, + tool_output_stats: Arc>>, +} + +impl SessionBoundEmitter { + #[must_use] + pub fn new(emitter: Emitter, session_id: String, tool_call_id: Option) -> Self { + Self { + emitter, + session_id, + tool_call_id, + tool_output_stats: Arc::new(Mutex::new(None)), + } + } + + pub fn take_tool_output_stats(&self) -> Option { + self.tool_output_stats + .lock() + .expect("tool output stats lock poisoned") + .take() + } } impl AgentEventEmitter for SessionBoundEmitter { @@ -75,6 +97,13 @@ impl AgentEventEmitter for SessionBoundEmitter { self.tool_call_id.clone(), ); } + + fn record_tool_output_stats(&self, stats: OutputCaptureStats) { + *self + .tool_output_stats + .lock() + .expect("tool output stats lock poisoned") = Some(stats); + } } #[cfg(test)] diff --git a/lib/components/fabro-agent/src/profiles/kimi_tools.rs b/lib/components/fabro-agent/src/profiles/kimi_tools.rs index 6ab61f1dd..6f13f2b57 100644 --- a/lib/components/fabro-agent/src/profiles/kimi_tools.rs +++ b/lib/components/fabro-agent/src/profiles/kimi_tools.rs @@ -144,12 +144,13 @@ explicitly asked. Never run commands requiring superuser privileges unless expli let _ = write!(out, "Command failed with exit code: {code}"); } let is_success = result.is_success(); - let out = retain_tool_output( + let retained = retain_tool_output( &out, MAX_RETAINED_TOOL_OUTPUT_BYTES, streaming.output_capture().omitted_bytes, - ) - .output; + ); + ctx.record_tool_output_stats(retained.stats); + let out = retained.output; emit_shell_process_completed(&ctx, streaming).await; if is_success { Ok(out) } else { Err(out) } }) diff --git a/lib/components/fabro-agent/src/tool_execution.rs b/lib/components/fabro-agent/src/tool_execution.rs index 228799654..a2f14d1d9 100644 --- a/lib/components/fabro-agent/src/tool_execution.rs +++ b/lib/components/fabro-agent/src/tool_execution.rs @@ -8,7 +8,7 @@ use tracing::debug; use crate::config::{SessionOptions, ToolHookCallback, ToolHookDecision}; use crate::event::{Emitter, SessionBoundEmitter}; use crate::question_tools::{self, AgentToolRuntime, is_question_tool}; -use crate::sandbox::Sandbox; +use crate::sandbox::{OutputCaptureStats, Sandbox}; use crate::session::ToolEnvProvider; use crate::tool_registry::{AgentEventEmitter, RegisteredTool, ToolContext, ToolRegistry}; use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output, truncate_tool_output}; @@ -265,9 +265,15 @@ fn error_tool_result_with_events( message: &str, ) -> ToolResult { emit_tool_call_started(emitter, session_id, tc); - let result = retain_tool_result(&ToolResult::error(&tc.id, message)); - emit_tool_call_result(emitter, session_id, tc, &result); - truncate_tool_result(&result, &tc.name, config) + let retained = retain_tool_result(&ToolResult::error(&tc.id, message), None); + emit_tool_call_result( + emitter, + session_id, + tc, + &retained.result, + retained.output_stats, + ); + truncate_tool_result(&retained.result, &tc.name, config) } fn emit_tool_call_started(emitter: &Emitter, session_id: &str, tc: &ToolCall) { @@ -278,15 +284,24 @@ fn emit_tool_call_started(emitter: &Emitter, session_id: &str, tc: &ToolCall) { }); } -fn emit_tool_call_result(emitter: &Emitter, session_id: &str, tc: &ToolCall, result: &ToolResult) { +fn emit_tool_call_result( + emitter: &Emitter, + session_id: &str, + tc: &ToolCall, + result: &ToolResult, + output_stats: OutputCaptureStats, +) { emitter.emit(session_id.to_owned(), AgentEvent::ToolCallOutputDelta { delta: result.content.to_string(), }); emitter.emit(session_id.to_owned(), AgentEvent::ToolCallCompleted { - tool_name: tc.name.clone(), - tool_call_id: tc.id.clone(), - output: result.content.clone(), - is_error: result.is_error, + tool_name: tc.name.clone(), + tool_call_id: tc.id.clone(), + output: result.content.clone(), + is_error: result.is_error, + output_bytes_observed: output_stats.observed_bytes, + output_bytes_retained: output_stats.retained_bytes, + output_bytes_omitted: output_stats.omitted_bytes, }); } @@ -386,9 +401,15 @@ async fn execute_and_emit_one_tool_with_lookup( emit_tool_call_started(emitter, session_id, tc); if let Some(reason) = access_denial { - let result = retain_tool_result(&ToolResult::error(&tc.id, &reason)); - emit_tool_call_result(emitter, session_id, tc, &result); - return truncate_tool_result(&result, &tc.name, config); + let retained = retain_tool_result(&ToolResult::error(&tc.id, &reason), None); + emit_tool_call_result( + emitter, + session_id, + tc, + &retained.result, + retained.output_stats, + ); + return truncate_tool_result(&retained.result, &tc.name, config); } // Pre-tool-use hook @@ -400,28 +421,34 @@ async fn execute_and_emit_one_tool_with_lookup( debug!(tool = %tc.name, hook_event = "pre_tool_use", ?decision, duration_ms = elapsed, "Tool hook complete"); if let ToolHookDecision::Block { reason } = decision { - let result = retain_tool_result(&ToolResult::error(&tc.id, &reason)); - emit_tool_call_result(emitter, session_id, tc, &result); - return truncate_tool_result(&result, &tc.name, config); + let retained = retain_tool_result(&ToolResult::error(&tc.id, &reason), None); + emit_tool_call_result( + emitter, + session_id, + tc, + &retained.result, + retained.output_stats, + ); + return truncate_tool_result(&retained.result, &tc.name, config); } } - let result = retain_tool_result( - &execute_one_tool( - tc, - registered_tool, - env, - cancel_token, - emitter, - session_id, - root_session_id, - tool_env_provider, - agent_tool_runtime, - ) - .await, - ); + let executed = execute_one_tool( + tc, + registered_tool, + env, + cancel_token, + emitter, + session_id, + root_session_id, + tool_env_provider, + agent_tool_runtime, + ) + .await; + let retained = retain_tool_result(&executed.result, executed.output_stats); + let result = retained.result; - emit_tool_call_result(emitter, session_id, tc, &result); + emit_tool_call_result(emitter, session_id, tc, &result, retained.output_stats); // Post-tool-use hooks if let Some(hooks) = tool_hooks { @@ -448,24 +475,48 @@ async fn execute_and_emit_one_tool_with_lookup( truncate_tool_result(&result, &tc.name, config) } +struct RetainedToolResult { + result: ToolResult, + output_stats: OutputCaptureStats, +} + /// Bound model-native tool output before it reaches hooks, events, or history. -fn retain_tool_result(result: &ToolResult) -> ToolResult { - let retained_content = match &result.content { - serde_json::Value::String(output) => serde_json::Value::String( - retain_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, 0).output, - ), - other => other.clone(), +fn retain_tool_result( + result: &ToolResult, + previous_stats: Option, +) -> RetainedToolResult { + let (retained_content, output_stats) = match &result.content { + serde_json::Value::String(output) => { + let previously_omitted = previous_stats.map_or(0, |stats| stats.omitted_bytes); + let retained = + retain_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, previously_omitted); + (serde_json::Value::String(retained.output), retained.stats) + } + other => { + let byte_count = serde_json::to_vec(other) + .expect("serde_json::Value always serializes") + .len(); + (other.clone(), OutputCaptureStats::complete(byte_count)) + } }; - ToolResult { - tool_call_id: result.tool_call_id.clone(), - content: retained_content, - is_error: result.is_error, - image_data: result.image_data.clone(), - image_media_type: result.image_media_type.clone(), + RetainedToolResult { + result: ToolResult { + tool_call_id: result.tool_call_id.clone(), + content: retained_content, + is_error: result.is_error, + image_data: result.image_data.clone(), + image_media_type: result.image_media_type.clone(), + }, + output_stats, } } +struct ExecutedToolResult { + result: ToolResult, + output_stats: Option, +} + /// Execute a single tool call: argument validation and execution. #[allow( clippy::too_many_arguments, @@ -481,23 +532,27 @@ async fn execute_one_tool( root_session_id: &str, tool_env_provider: Option<&Arc>, agent_tool_runtime: &AgentToolRuntime, -) -> ToolResult { +) -> ExecutedToolResult { match registered_tool { Some(tool) => { if tc.tool_type != "custom" { if let Err(validation_error) = validate_tool_args(&tool.definition.parameters, &tc.arguments) { - return ToolResult::error(&tc.id, validation_error); + return ExecutedToolResult { + result: ToolResult::error(&tc.id, validation_error), + output_stats: None, + }; } } + let session_emitter = Arc::new(SessionBoundEmitter::new( + emitter.clone(), + session_id.to_owned(), + Some(tc.id.clone()), + )); let agent_event_emitter: Option> = - Some(Arc::new(SessionBoundEmitter { - emitter: emitter.clone(), - session_id: session_id.to_owned(), - tool_call_id: Some(tc.id.clone()), - })); + Some(session_emitter.clone()); let ctx = ToolContext { env, cancel: cancel_token, @@ -508,14 +563,24 @@ async fn execute_one_tool( agent_event_emitter, }; let execution = (tool.executor)(tc.arguments.clone(), ctx); - match question_tools::scope_agent_tool_runtime(agent_tool_runtime.clone(), execution) - .await + let result = match question_tools::scope_agent_tool_runtime( + agent_tool_runtime.clone(), + execution, + ) + .await { Ok(output) => ToolResult::success(&tc.id, serde_json::json!(output)), Err(err) => ToolResult::error(&tc.id, err), + }; + ExecutedToolResult { + result, + output_stats: session_emitter.take_tool_output_stats(), } } - None => ToolResult::error(&tc.id, format!("Unknown tool: {}", tc.name)), + None => ExecutedToolResult { + result: ToolResult::error(&tc.id, format!("Unknown tool: {}", tc.name)), + output_stats: None, + }, } } @@ -923,19 +988,37 @@ mod tests { let result_output = result.content.as_str().expect("string tool output"); assert_eq!(result_output.len(), MAX_RETAINED_TOOL_OUTPUT_BYTES); - let completed_output = loop { + let completed = loop { let event = receiver.try_recv().expect("tool completion event"); - if let AgentEvent::ToolCallCompleted { output, .. } = event.event { - break output; + if let AgentEvent::ToolCallCompleted { + output, + is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, + .. + } = event.event + { + break ( + output, + is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, + ); } }; assert_eq!( - completed_output - .as_str() - .expect("string event output") - .len(), + completed.0.as_str().expect("string event output").len(), MAX_RETAINED_TOOL_OUTPUT_BYTES ); + assert!(!completed.1, "truncation must not make the tool an error"); + assert_eq!( + completed.2, + MAX_RETAINED_TOOL_OUTPUT_BYTES + 100 + "echo: ".len() + ); + assert_eq!(completed.3, MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert_eq!(completed.4, 100 + "echo: ".len()); } #[tokio::test] @@ -1232,6 +1315,68 @@ mod tests { assert!(!result.is_error); } + #[tokio::test] + async fn shell_events_record_process_and_rendered_output_byte_counts() { + let output_len = MAX_RETAINED_TOOL_OUTPUT_BYTES + 1_000; + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); + let result = run_shell_tool( + fabro_sandbox::ExecResult { + stdout: "x".repeat(output_len), + stderr: String::new(), + exit_code: Some(0), + termination: fabro_types::CommandTermination::Exited, + duration_ms: 12, + }, + None, + &emitter, + ) + .await; + + assert!(!result.is_error); + let events = drain(&mut receiver); + let process = events + .iter() + .find_map(|event| match &event.event { + AgentEvent::ToolProcessCompleted { + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, + .. + } => Some(( + *output_bytes_observed, + *output_bytes_retained, + *output_bytes_omitted, + )), + _ => None, + }) + .expect("process event"); + assert_eq!(process, (output_len, MAX_RETAINED_TOOL_OUTPUT_BYTES, 1_000)); + + let completed = events + .iter() + .find_map(|event| match &event.event { + AgentEvent::ToolCallCompleted { + is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, + .. + } => Some(( + *is_error, + *output_bytes_observed, + *output_bytes_retained, + *output_bytes_omitted, + )), + _ => None, + }) + .expect("tool completion event"); + assert!(!completed.0, "truncation must not make the tool an error"); + assert!(completed.1 > output_len); + assert_eq!(completed.2, MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert_eq!(completed.3, completed.1 - completed.2); + } + #[tokio::test] async fn shell_failure_emits_started_then_process_then_completed() { let emitter = Emitter::new(); diff --git a/lib/components/fabro-agent/src/tool_registry.rs b/lib/components/fabro-agent/src/tool_registry.rs index 9ee56522f..7e955aed0 100644 --- a/lib/components/fabro-agent/src/tool_registry.rs +++ b/lib/components/fabro-agent/src/tool_registry.rs @@ -9,7 +9,7 @@ use tokio_util::sync::CancellationToken; use crate::config::{ToolAccessPolicy, ToolExposureMode}; use crate::native_tool::{NativeTool, ToolVocabulary}; -use crate::sandbox::Sandbox; +use crate::sandbox::{OutputCaptureStats, Sandbox}; use crate::session::ToolEnvProvider; use crate::tool_permissions; use crate::types::AgentEvent; @@ -20,6 +20,10 @@ use crate::types::AgentEvent; /// the session is using. pub trait AgentEventEmitter: Send + Sync { fn emit(&self, event: AgentEvent); + + /// Record byte counts for the model-facing output produced by this tool. + /// Emitters without a tool-execution owner may ignore this side channel. + fn record_tool_output_stats(&self, _stats: OutputCaptureStats) {} } pub struct ToolContext { @@ -53,6 +57,13 @@ impl ToolContext { emitter.emit(event); } } + + /// Record model-facing output byte counts for the owning tool call. + pub fn record_tool_output_stats(&self, stats: OutputCaptureStats) { + if let Some(emitter) = self.agent_event_emitter.as_ref() { + emitter.record_tool_output_stats(stats); + } + } } pub type ToolExecutor = Arc< diff --git a/lib/components/fabro-agent/src/tools.rs b/lib/components/fabro-agent/src/tools.rs index d4d429146..a6060bb9f 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -320,7 +320,9 @@ pub(crate) async fn run_shell_command( cwd: Option<&str>, ) -> Result { let streaming = execute_shell_command(ctx, command, timeout_ms, cwd).await?; - let text = render_shell_result(&streaming).output; + let retained = render_shell_result(&streaming); + ctx.record_tool_output_stats(retained.stats); + let text = retained.output; let is_success = streaming.result.is_success(); emit_shell_process_completed(ctx, streaming).await; @@ -342,6 +344,7 @@ pub(crate) async fn emit_shell_process_completed( let termination = streaming.result.termination; let duration_ms = streaming.result.duration_ms; let streams_separated = streaming.streams_separated; + let output_stats = streaming.output_capture(); let result = streaming.result; let exec_output_tail = match task::spawn_blocking(move || result.default_redacted_output_tail()).await { @@ -360,6 +363,9 @@ pub(crate) async fn emit_shell_process_completed( duration_ms, streams_separated, exec_output_tail, + output_bytes_observed: output_stats.observed_bytes, + output_bytes_retained: output_stats.retained_bytes, + output_bytes_omitted: output_stats.omitted_bytes, }); } @@ -1099,11 +1105,11 @@ mod tests { 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()), - })), + agent_event_emitter: Some(Arc::new(SessionBoundEmitter::new( + emitter.clone(), + "test-session".to_string(), + Some("call_1".to_string()), + ))), ..shell_context(env) } } @@ -1319,11 +1325,16 @@ mod tests { duration_ms, streams_separated, exec_output_tail, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, } => { assert_eq!(exit_code, Some(7)); assert_eq!(termination, CommandTermination::Exited); assert_eq!(duration_ms, 12); assert!(streams_separated); + assert_eq!(output_bytes_observed, output_bytes_retained); + assert_eq!(output_bytes_omitted, 0); let tail = exec_output_tail.expect("output tail"); assert_eq!(tail.stdout.as_deref(), Some("out")); let stderr = tail.stderr.expect("stderr tail"); diff --git a/lib/components/fabro-agent/src/types.rs b/lib/components/fabro-agent/src/types.rs index de3d0434a..eba037472 100644 --- a/lib/components/fabro-agent/src/types.rs +++ b/lib/components/fabro-agent/src/types.rs @@ -348,10 +348,16 @@ pub enum AgentEvent { delta: String, }, ToolCallCompleted { - tool_name: String, - tool_call_id: String, - output: serde_json::Value, - is_error: bool, + tool_name: String, + tool_call_id: String, + output: serde_json::Value, + is_error: bool, + #[serde(default)] + output_bytes_observed: usize, + #[serde(default)] + output_bytes_retained: usize, + #[serde(default)] + output_bytes_omitted: usize, }, /// Subordinate process outcome for a tool call that ran a command. /// Emitted before the owning `ToolCallCompleted`, which stays the single @@ -359,12 +365,18 @@ pub enum AgentEvent { /// 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, + exit_code: Option, + termination: CommandTermination, + duration_ms: u64, + streams_separated: bool, #[serde(default, skip_serializing_if = "Option::is_none")] - exec_output_tail: Option, + exec_output_tail: Option, + #[serde(default)] + output_bytes_observed: usize, + #[serde(default)] + output_bytes_retained: usize, + #[serde(default)] + output_bytes_omitted: usize, }, Error { error: Error, @@ -558,6 +570,9 @@ impl AgentEvent { tool_name, tool_call_id, is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, .. } => { info!( @@ -565,6 +580,9 @@ impl AgentEvent { tool = tool_name.as_str(), tool_call_id, is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, "Tool call completed" ); } @@ -574,6 +592,9 @@ impl AgentEvent { duration_ms, streams_separated, exec_output_tail, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, } => { let tail = ExecOutputTail::trace_summary(exec_output_tail.as_ref()); debug!( @@ -587,6 +608,9 @@ impl AgentEvent { stderr_bytes = tail.stderr_bytes, stdout_truncated = tail.stdout_truncated, stderr_truncated = tail.stderr_truncated, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, "Tool process completed" ); } @@ -958,6 +982,21 @@ mod tests { })); } + #[test] + fn legacy_tool_completion_defaults_output_byte_counts() { + let event: AgentEvent = serde_json::from_str( + r#"{"ToolCallCompleted":{"tool_name":"shell","tool_call_id":"call_1","output":"ok","is_error":false}}"#, + ) + .unwrap(); + + assert!(matches!(event, AgentEvent::ToolCallCompleted { + output_bytes_observed: 0, + output_bytes_retained: 0, + output_bytes_omitted: 0, + .. + })); + } + #[test] fn session_event_serde_round_trip_without_parent_session_id() { let event = SessionEvent { diff --git a/lib/components/fabro-agent/tests/it/docker_shell.rs b/lib/components/fabro-agent/tests/it/docker_shell.rs index 6f50423ba..19f0f00f8 100644 --- a/lib/components/fabro-agent/tests/it/docker_shell.rs +++ b/lib/components/fabro-agent/tests/it/docker_shell.rs @@ -50,11 +50,11 @@ async fn shell_reports_real_docker_process_outcome() { 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()), - })), + agent_event_emitter: Some(Arc::new(SessionBoundEmitter::new( + emitter.clone(), + "test-session".to_string(), + Some("call_1".to_string()), + ))), }, ) .await; diff --git a/lib/components/fabro-store/src/run_sessions.rs b/lib/components/fabro-store/src/run_sessions.rs index a5998baed..a8aafc079 100644 --- a/lib/components/fabro-store/src/run_sessions.rs +++ b/lib/components/fabro-store/src/run_sessions.rs @@ -354,6 +354,9 @@ mod tests { tool_call_id: "call_1".to_string(), output: json!("contents"), is_error: false, + output_bytes_observed: None, + output_bytes_retained: None, + output_bytes_omitted: None, }), ), ]; diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index da54a2d29..ec25688df 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -1761,13 +1761,16 @@ mod tests { fn tool_completed(tool_call_id: &str) -> EventBody { EventBody::AgentToolCompleted(AgentToolCompletedProps { - tool_name: "Bash".to_string(), - tool_call_id: tool_call_id.to_string(), - output: json!("ok"), - is_error: false, - visit: 1, - tool_result: None, - turn_id: None, + tool_name: "Bash".to_string(), + tool_call_id: tool_call_id.to_string(), + output: json!("ok"), + is_error: false, + visit: 1, + output_bytes_observed: None, + output_bytes_retained: None, + output_bytes_omitted: None, + tool_result: None, + turn_id: None, }) } diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 7d9af8f79..f9ab23d4e 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -23,6 +23,10 @@ fn stage_status_from_string(status: &str) -> StageOutcome { }) } +fn output_byte_count(value: usize) -> u64 { + u64::try_from(value).unwrap_or(u64::MAX) +} + /// Project the sandbox layer's runtime push attempts into the durable /// `git.push` attempt shape. /// @@ -690,12 +694,18 @@ fn event_body_from_event(event: &Event) -> EventBody { tool_call_id, output, is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, } => EventBody::AgentToolCompleted(fabro_types::AgentToolCompletedProps { tool_name: tool_name.clone(), tool_call_id: tool_call_id.clone(), output: output.clone(), is_error: *is_error, visit: *visit, + output_bytes_observed: Some(output_byte_count(*output_bytes_observed)), + output_bytes_retained: Some(output_byte_count(*output_bytes_retained)), + output_bytes_omitted: Some(output_byte_count(*output_bytes_omitted)), tool_result: None, turn_id: None, }), @@ -705,6 +715,9 @@ fn event_body_from_event(event: &Event) -> EventBody { duration_ms, streams_separated, exec_output_tail, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, } => EventBody::AgentToolProcessCompleted( fabro_types::AgentToolProcessCompletedProps { exit_code: *exit_code, @@ -712,6 +725,9 @@ fn event_body_from_event(event: &Event) -> EventBody { duration_ms: *duration_ms, streams_separated: *streams_separated, exec_output_tail: exec_output_tail.clone(), + output_bytes_observed: Some(output_byte_count(*output_bytes_observed)), + output_bytes_retained: Some(output_byte_count(*output_bytes_retained)), + output_bytes_omitted: Some(output_byte_count(*output_bytes_omitted)), visit: *visit, }, ), @@ -1620,11 +1636,14 @@ mod tests { 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()), + exit_code: Some(7), + termination: ::fabro_types::CommandTermination::Exited, + duration_ms: 12, + streams_separated: true, + output_bytes_observed: 120, + output_bytes_retained: 100, + output_bytes_omitted: 20, + exec_output_tail: Some(exec_tail()), }, session_id: Some("ses_child".to_string()), parent_session_id: Some("ses_parent".to_string()), @@ -1651,6 +1670,9 @@ mod tests { assert_eq!(properties["termination"], "exited"); assert_eq!(properties["duration_ms"], 12); assert_eq!(properties["streams_separated"], true); + assert_eq!(properties["output_bytes_observed"], 120); + assert_eq!(properties["output_bytes_retained"], 100); + assert_eq!(properties["output_bytes_omitted"], 20); assert_eq!(properties["exec_output_tail"]["stdout"], "last stdout line"); assert_eq!(properties["visit"], 2); } diff --git a/lib/components/fabro-workflow/src/event/redaction.rs b/lib/components/fabro-workflow/src/event/redaction.rs index 49e86a2ab..339f1edd3 100644 --- a/lib/components/fabro-workflow/src/event/redaction.rs +++ b/lib/components/fabro-workflow/src/event/redaction.rs @@ -83,11 +83,14 @@ mod tests { 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 { + exit_code: Some(7), + termination: ::fabro_types::CommandTermination::Exited, + duration_ms: 12, + streams_separated: true, + output_bytes_observed: 100, + output_bytes_retained: 100, + output_bytes_omitted: 0, + exec_output_tail: Some(fabro_types::ExecOutputTail { stdout: Some(format!("stdout {secret}")), stderr: Some("plain stderr".to_string()), stdout_truncated: false, diff --git a/lib/components/fabro-workflow/src/handler/llm/api.rs b/lib/components/fabro-workflow/src/handler/llm/api.rs index d7e44bf96..30b765f4e 100644 --- a/lib/components/fabro-workflow/src/handler/llm/api.rs +++ b/lib/components/fabro-workflow/src/handler/llm/api.rs @@ -2842,10 +2842,13 @@ reasoning = false track_file_event( &AgentEvent::ToolCallCompleted { - tool_call_id: "tc1".to_string(), - tool_name: "write_file".to_string(), - is_error: false, - output: serde_json::Value::String("ok".to_string()), + tool_call_id: "tc1".to_string(), + tool_name: "write_file".to_string(), + is_error: false, + output: serde_json::Value::String("ok".to_string()), + output_bytes_observed: 2, + output_bytes_retained: 2, + output_bytes_omitted: 0, }, &mut state, ); @@ -2875,10 +2878,13 @@ reasoning = false track_file_event( &AgentEvent::ToolCallCompleted { - tool_call_id: "tc-sub".to_string(), - tool_name: "edit_file".to_string(), - is_error: false, - output: serde_json::Value::String("ok".to_string()), + tool_call_id: "tc-sub".to_string(), + tool_name: "edit_file".to_string(), + is_error: false, + output: serde_json::Value::String("ok".to_string()), + output_bytes_observed: 2, + output_bytes_retained: 2, + output_bytes_omitted: 0, }, &mut state, ); @@ -2902,10 +2908,13 @@ reasoning = false ); track_file_event( &AgentEvent::ToolCallCompleted { - tool_call_id: "tc-kimi".to_string(), - tool_name: "Write".to_string(), - is_error: false, - output: serde_json::Value::String("ok".to_string()), + tool_call_id: "tc-kimi".to_string(), + tool_name: "Write".to_string(), + is_error: false, + output: serde_json::Value::String("ok".to_string()), + output_bytes_observed: 2, + output_bytes_retained: 2, + output_bytes_omitted: 0, }, &mut state, ); @@ -2935,10 +2944,13 @@ reasoning = false track_file_event( &AgentEvent::ToolCallCompleted { - tool_call_id: "tc-err".to_string(), - tool_name: "edit_file".to_string(), - is_error: true, - output: serde_json::Value::String("failed".to_string()), + tool_call_id: "tc-err".to_string(), + tool_name: "edit_file".to_string(), + is_error: true, + output: serde_json::Value::String("failed".to_string()), + output_bytes_observed: 6, + output_bytes_retained: 6, + output_bytes_omitted: 0, }, &mut state, ); diff --git a/lib/foundation/fabro-types/src/run_event/agent.rs b/lib/foundation/fabro-types/src/run_event/agent.rs index a20ee4aaf..171e69aa3 100644 --- a/lib/foundation/fabro-types/src/run_event/agent.rs +++ b/lib/foundation/fabro-types/src/run_event/agent.rs @@ -167,18 +167,27 @@ pub struct AgentToolStartedProps { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct AgentToolCompletedProps { // Narrow legacy fields retained for consumer compatibility. - pub tool_name: String, - pub tool_call_id: String, - pub output: Value, - pub is_error: bool, - pub visit: u32, + pub tool_name: String, + pub tool_call_id: String, + pub output: Value, + pub is_error: bool, + pub visit: u32, + /// UTF-8 bytes in the rendered tool output before hard retention. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_observed: Option, + /// Tool-output bytes kept in `output`, excluding truncation notices. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_retained: Option, + /// Tool-output bytes discarded before this event was emitted. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_omitted: Option, /// Canonical tool result payload. Carries the structured output, error /// state, and supported media/artifact fields. #[serde(default, skip_serializing_if = "Option::is_none")] - pub tool_result: Option, + pub tool_result: Option, /// Turn that owned this tool call. #[serde(default, skip_serializing_if = "Option::is_none")] - pub turn_id: Option, + pub turn_id: Option, } /// Subordinate diagnostic for a tool call that ran a process: the real @@ -189,15 +198,24 @@ pub struct AgentToolCompletedProps { #[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, + 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, + pub streams_separated: bool, #[serde(default, skip_serializing_if = "Option::is_none")] - pub exec_output_tail: Option, - pub visit: u32, + pub exec_output_tail: Option, + /// Raw stdout and stderr bytes drained from the process. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_observed: Option, + /// Raw process-output bytes kept by the streaming capture buffers. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_retained: Option, + /// Raw process-output bytes discarded by the streaming capture buffers. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_omitted: Option, + pub visit: u32, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -642,6 +660,9 @@ mod tests { let props: AgentToolCompletedProps = serde_json::from_value(v).unwrap(); assert!(props.tool_result.is_none()); assert!(props.turn_id.is_none()); + assert!(props.output_bytes_observed.is_none()); + assert!(props.output_bytes_retained.is_none()); + assert!(props.output_bytes_omitted.is_none()); } #[test] @@ -649,16 +670,22 @@ mod tests { let tr = ToolResult::success("call_1", json!({"stdout": "ok"})); let turn = TurnId::new(); let props = AgentToolCompletedProps { - tool_name: "Bash".to_string(), - tool_call_id: "call_1".to_string(), - output: json!({"stdout": "ok"}), - is_error: false, - visit: 1, - tool_result: Some(tr.clone()), - turn_id: Some(turn), + tool_name: "Bash".to_string(), + tool_call_id: "call_1".to_string(), + output: json!({"stdout": "ok"}), + is_error: false, + visit: 1, + output_bytes_observed: Some(120), + output_bytes_retained: Some(100), + output_bytes_omitted: Some(20), + tool_result: Some(tr.clone()), + turn_id: Some(turn), }; let v = serde_json::to_value(&props).unwrap(); assert_eq!(v["tool_result"]["content"]["stdout"], "ok"); + assert_eq!(v["output_bytes_observed"], 120); + assert_eq!(v["output_bytes_retained"], 100); + assert_eq!(v["output_bytes_omitted"], 20); let back: AgentToolCompletedProps = serde_json::from_value(v).unwrap(); assert_eq!(back, props); } diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index 6323e5071..30d56a110 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -2458,6 +2458,9 @@ mod tests { "duration_ms": 12, "streams_separated": true, "exec_output_tail": {"stdout": "out", "stderr": "err"}, + "output_bytes_observed": 150, + "output_bytes_retained": 100, + "output_bytes_omitted": 50, "visit": 1 } }); @@ -2472,6 +2475,9 @@ mod tests { assert_eq!(props.termination, CommandTermination::Exited); assert_eq!(props.duration_ms, 12); assert!(props.streams_separated); + assert_eq!(props.output_bytes_observed, Some(150)); + assert_eq!(props.output_bytes_retained, Some(100)); + assert_eq!(props.output_bytes_omitted, Some(50)); assert_eq!( props.exec_output_tail.as_ref().unwrap().stdout.as_deref(), Some("out") @@ -2483,12 +2489,15 @@ mod tests { #[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, + exit_code: None, + termination: CommandTermination::TimedOut, + duration_ms: 10_000, + streams_separated: false, + exec_output_tail: None, + output_bytes_observed: None, + output_bytes_retained: None, + output_bytes_omitted: None, + visit: 1, }); let value = serde_json::to_value(&body).unwrap(); @@ -2499,6 +2508,9 @@ mod tests { let properties = value["properties"].as_object().unwrap(); assert!(!properties.contains_key("exit_code")); assert!(!properties.contains_key("exec_output_tail")); + assert!(!properties.contains_key("output_bytes_observed")); + assert!(!properties.contains_key("output_bytes_retained")); + assert!(!properties.contains_key("output_bytes_omitted")); let parsed: EventBody = serde_json::from_value(value).unwrap(); assert_eq!(parsed, body); diff --git a/lib/foundation/fabro-types/src/run_event/session.rs b/lib/foundation/fabro-types/src/run_event/session.rs index 78321afa4..02dd98ea1 100644 --- a/lib/foundation/fabro-types/src/run_event/session.rs +++ b/lib/foundation/fabro-types/src/run_event/session.rs @@ -52,11 +52,17 @@ pub struct RunSessionToolCallStartedProps { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct RunSessionToolCallCompletedProps { - pub turn_id: TurnId, - pub tool_name: String, - pub tool_call_id: String, - pub output: Value, - pub is_error: bool, + pub turn_id: TurnId, + pub tool_name: String, + pub tool_call_id: String, + pub output: Value, + pub is_error: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_observed: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_retained: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub output_bytes_omitted: Option, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -115,7 +121,8 @@ pub struct RunSessionTurnInterruptedProps { mod tests { use serde_json::json; - use super::RunSessionCreatedProps; + use super::{RunSessionCreatedProps, RunSessionToolCallCompletedProps}; + use crate::TurnId; #[test] fn session_created_deserializes_legacy_payload_without_provider() { @@ -128,4 +135,20 @@ mod tests { assert_eq!(props.model.as_deref(), Some("gpt-5.4")); assert_eq!(props.provider, None); } + + #[test] + fn tool_completion_deserializes_without_output_byte_counts() { + let props: RunSessionToolCallCompletedProps = serde_json::from_value(json!({ + "turn_id": TurnId::new(), + "tool_name": "shell", + "tool_call_id": "call_1", + "output": "ok", + "is_error": false + })) + .unwrap(); + + assert!(props.output_bytes_observed.is_none()); + assert!(props.output_bytes_retained.is_none()); + assert!(props.output_bytes_omitted.is_none()); + } } From a28a0378dc3945078fd5c46c63ac0f5a99a65e58 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 12:53:53 -0400 Subject: [PATCH 3/6] Report truncated tool output to agents --- .../fabro-agent/src/tool_execution.rs | 39 ++-- lib/components/fabro-agent/src/tools.rs | 3 +- lib/components/fabro-agent/src/truncation.rs | 210 +++++++++++++++--- 3 files changed, 208 insertions(+), 44 deletions(-) diff --git a/lib/components/fabro-agent/src/tool_execution.rs b/lib/components/fabro-agent/src/tool_execution.rs index a2f14d1d9..3cb8b674a 100644 --- a/lib/components/fabro-agent/src/tool_execution.rs +++ b/lib/components/fabro-agent/src/tool_execution.rs @@ -11,7 +11,9 @@ use crate::question_tools::{self, AgentToolRuntime, is_question_tool}; use crate::sandbox::{OutputCaptureStats, Sandbox}; use crate::session::ToolEnvProvider; use crate::tool_registry::{AgentEventEmitter, RegisteredTool, ToolContext, ToolRegistry}; -use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output, truncate_tool_output}; +use crate::truncation::{ + MAX_RETAINED_TOOL_OUTPUT_BYTES, preview_tool_output, truncate_tool_output, +}; use crate::types::AgentEvent; /// Execute tool calls, choosing parallel or sequential based on `parallel` @@ -489,7 +491,7 @@ fn retain_tool_result( serde_json::Value::String(output) => { let previously_omitted = previous_stats.map_or(0, |stats| stats.omitted_bytes); let retained = - retain_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, previously_omitted); + preview_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, previously_omitted); (serde_json::Value::String(retained.output), retained.stats) } other => { @@ -986,7 +988,11 @@ mod tests { .await; let result_output = result.content.as_str().expect("string tool output"); - assert_eq!(result_output.len(), MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert!(result_output.len() <= MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert!(result_output.starts_with("Warning: truncated output")); + assert!(result_output.contains("bytes omitted")); + assert!(result_output.contains("tokens truncated")); + assert!(!result_output.contains("re-run")); let completed = loop { let event = receiver.try_recv().expect("tool completion event"); @@ -1008,17 +1014,16 @@ mod tests { ); } }; - assert_eq!( - completed.0.as_str().expect("string event output").len(), - MAX_RETAINED_TOOL_OUTPUT_BYTES - ); + let event_output = completed.0.as_str().expect("string event output"); + assert_eq!(event_output, result_output); assert!(!completed.1, "truncation must not make the tool an error"); assert_eq!( completed.2, MAX_RETAINED_TOOL_OUTPUT_BYTES + 100 + "echo: ".len() ); - assert_eq!(completed.3, MAX_RETAINED_TOOL_OUTPUT_BYTES); - assert_eq!(completed.4, 100 + "echo: ".len()); + assert!(completed.3 < MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert_eq!(completed.4, completed.2 - completed.3); + assert!(event_output.contains(&format!("... {} bytes omitted ...", completed.4))); } #[tokio::test] @@ -1357,12 +1362,14 @@ mod tests { .iter() .find_map(|event| match &event.event { AgentEvent::ToolCallCompleted { + output, is_error, output_bytes_observed, output_bytes_retained, output_bytes_omitted, .. } => Some(( + output.as_str().expect("string event output"), *is_error, *output_bytes_observed, *output_bytes_retained, @@ -1371,10 +1378,16 @@ mod tests { _ => None, }) .expect("tool completion event"); - assert!(!completed.0, "truncation must not make the tool an error"); - assert!(completed.1 > output_len); - assert_eq!(completed.2, MAX_RETAINED_TOOL_OUTPUT_BYTES); - assert_eq!(completed.3, completed.1 - completed.2); + assert!(completed.0.starts_with("Warning: truncated output")); + assert!(!completed.1, "truncation must not make the tool an error"); + assert!(completed.2 > output_len); + assert!(completed.3 < MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert_eq!(completed.4, completed.2 - completed.3); + assert!( + completed + .0 + .contains(&format!("... {} bytes omitted ...", completed.4)) + ); } #[tokio::test] diff --git a/lib/components/fabro-agent/src/tools.rs b/lib/components/fabro-agent/src/tools.rs index a6060bb9f..24ab550d2 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -1408,7 +1408,8 @@ mod tests { 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.starts_with("Warning: truncated output")); + assert!(truncated.contains("Termination: exited\nExit code: 2\n")); assert!( truncated.contains("stderr:\nthe build failed"), "stderr tail did not survive truncation" diff --git a/lib/components/fabro-agent/src/truncation.rs b/lib/components/fabro-agent/src/truncation.rs index edc3122b3..af85eb2c6 100644 --- a/lib/components/fabro-agent/src/truncation.rs +++ b/lib/components/fabro-agent/src/truncation.rs @@ -8,6 +8,7 @@ pub(crate) const MAX_RETAINED_TOOL_OUTPUT_BYTES: usize = 1024 * 1024; pub(crate) struct RetainedToolOutput { pub output: String, pub stats: OutputCaptureStats, + head_bytes: usize, } /// Keep an equal-sized UTF-8 prefix and suffix within a byte budget. @@ -23,12 +24,17 @@ pub(crate) fn retain_tool_output( let observed_bytes = output.len().saturating_add(previously_omitted_bytes); if output.len() <= max_bytes { return RetainedToolOutput { - output: output.to_string(), - stats: OutputCaptureStats { + output: output.to_string(), + stats: OutputCaptureStats { observed_bytes, retained_bytes: output.len(), omitted_bytes: previously_omitted_bytes, }, + head_bytes: if previously_omitted_bytes == 0 { + output.len() + } else { + output.floor_char_boundary(output.len() / 2) + }, }; } @@ -42,15 +48,85 @@ pub(crate) fn retain_tool_output( retained.push_str(&output[tail_start..]); RetainedToolOutput { - output: retained, - stats: OutputCaptureStats { + output: retained, + stats: OutputCaptureStats { observed_bytes, retained_bytes, omitted_bytes: observed_bytes.saturating_sub(retained_bytes), }, + head_bytes: head_end, } } +/// Build the final model-facing preview, including truncation notices inside +/// the total byte budget. +#[must_use] +pub(crate) fn preview_tool_output( + output: &str, + max_bytes: usize, + previously_omitted_bytes: usize, +) -> RetainedToolOutput { + let mut content_budget = max_bytes; + loop { + let retained = retain_tool_output(output, content_budget, previously_omitted_bytes); + if retained.stats.omitted_bytes == 0 { + return retained; + } + + let rendered = render_retained_output(&retained); + if rendered.len() <= max_bytes { + return RetainedToolOutput { + output: rendered, + stats: retained.stats, + head_bytes: 0, + }; + } + + let excess = rendered.len().saturating_sub(max_bytes).max(1); + let next_budget = content_budget.saturating_sub(excess); + if next_budget == content_budget { + return RetainedToolOutput { + output: truncate_plain_output(&rendered, max_bytes, TruncationMode::HeadTail), + stats: retained.stats, + head_bytes: 0, + }; + } + content_budget = next_budget; + } +} + +fn render_retained_output(retained: &RetainedToolOutput) -> String { + let head = &retained.output[..retained.head_bytes]; + let tail = &retained.output[retained.head_bytes..]; + render_truncated_segments(head, tail, retained.stats, None) +} + +fn render_truncated_segments( + head: &str, + tail: &str, + stats: OutputCaptureStats, + line_count_omitted: Option, +) -> String { + let original_tokens = approximate_tokens(stats.observed_bytes); + let omitted_tokens = approximate_tokens(stats.omitted_bytes); + let middle_marker = line_count_omitted.map_or_else( + || format!("... approximately {omitted_tokens} tokens truncated ..."), + |lines| { + format!( + "... {lines} lines omitted (approximately {omitted_tokens} tokens truncated) ..." + ) + }, + ); + format!( + "Warning: truncated output (original token count: {original_tokens})\n... {} bytes omitted ...\n\n{head}\n\n{middle_marker}\n\n{tail}", + stats.omitted_bytes + ) +} + +fn approximate_tokens(bytes: usize) -> usize { + bytes.div_ceil(4) +} + fn ceil_char_boundary(output: &str, index: usize) -> usize { let mut index = index.min(output.len()); while index < output.len() && !output.is_char_boundary(index) { @@ -98,32 +174,61 @@ pub fn truncate_output(output: &str, max_chars: usize, mode: TruncationMode) -> return output.to_string(); } - let removed = output.len() - max_chars; - match mode { TruncationMode::HeadTail => { let half = max_chars / 2; let head_end = output.floor_char_boundary(half); - let tail_start = output.floor_char_boundary(output.len() - half); + let tail_start = ceil_char_boundary(output, output.len().saturating_sub(half)); let head = &output[..head_end]; let tail = &output[tail_start..]; - format!( - "{head}\n\n[WARNING: Tool output was truncated. {removed} characters were removed from the middle. \ - The full output is available in the event stream. \ - If you need to see specific parts, re-run the tool with more targeted parameters.]\n\n{tail}" + let retained_bytes = head.len().saturating_add(tail.len()); + render_truncated_segments( + head, + tail, + OutputCaptureStats { + observed_bytes: output.len(), + retained_bytes, + omitted_bytes: output.len().saturating_sub(retained_bytes), + }, + None, ) } TruncationMode::Tail => { - let tail_start = output.floor_char_boundary(output.len() - max_chars); + let tail_start = ceil_char_boundary(output, output.len().saturating_sub(max_chars)); let tail = &output[tail_start..]; - format!( - "[WARNING: Tool output was truncated. First {removed} characters were removed. \ - The full output is available in the event stream.]\n\n{tail}" + render_truncated_segments( + "", + tail, + OutputCaptureStats { + observed_bytes: output.len(), + retained_bytes: tail.len(), + omitted_bytes: output.len().saturating_sub(tail.len()), + }, + None, ) } } } +fn truncate_plain_output(output: &str, max_bytes: usize, mode: TruncationMode) -> String { + if output.len() <= max_bytes { + return output.to_string(); + } + + match mode { + TruncationMode::HeadTail => { + let half = max_bytes / 2; + let head_end = output.floor_char_boundary(half); + let tail_start = ceil_char_boundary(output, output.len().saturating_sub(half)); + format!("{}{}", &output[..head_end], &output[tail_start..]) + } + TruncationMode::Tail => { + let tail_start = ceil_char_boundary(output, output.len().saturating_sub(max_bytes)); + output[tail_start..].to_string() + } + } +} + #[must_use] pub fn truncate_lines(output: &str, max_lines: usize) -> String { let lines: Vec<&str> = output.lines().collect(); @@ -131,15 +236,22 @@ pub fn truncate_lines(output: &str, max_lines: usize) -> String { return output.to_string(); } - let half = max_lines / 2; - let head: Vec<&str> = lines[..half].to_vec(); - let tail: Vec<&str> = lines[lines.len() - half..].to_vec(); + let head_count = max_lines / 2; + let tail_count = max_lines.saturating_sub(head_count); + let head = lines[..head_count].join("\n"); + let tail = lines[lines.len() - tail_count..].join("\n"); let omitted = lines.len() - max_lines; + let retained_bytes = head.len().saturating_add(tail.len()); - format!( - "{}\n\n[... {omitted} lines omitted ...]\n\n{}", - head.join("\n"), - tail.join("\n") + render_truncated_segments( + &head, + &tail, + OutputCaptureStats { + observed_bytes: output.len(), + retained_bytes, + omitted_bytes: output.len().saturating_sub(retained_bytes), + }, + Some(omitted), ) } @@ -203,6 +315,41 @@ mod tests { ); } + #[test] + fn model_preview_includes_codex_style_notice_inside_budget() { + let output = format!("HEAD{}TAIL", "x".repeat(1_000)); + let preview = preview_tool_output(&output, 512, 0); + + assert!(preview.output.len() <= 512, "{}", preview.output.len()); + assert!( + preview + .output + .starts_with("Warning: truncated output (original token count: 252)") + ); + assert!(preview.output.contains(&format!( + "... {} bytes omitted ...", + preview.stats.omitted_bytes + ))); + assert!(preview.output.contains("approximately")); + assert!(preview.output.contains("tokens truncated")); + assert!(preview.output.contains("HEAD")); + assert!(preview.output.ends_with("TAIL")); + assert!(!preview.output.contains("re-run")); + assert!(!preview.output.contains("targeted parameters")); + } + + #[test] + fn model_preview_reports_bytes_omitted_before_rendering() { + let preview = preview_tool_output("abcdefgh", 512, 100); + + assert_eq!(preview.stats.observed_bytes, 108); + assert_eq!(preview.stats.retained_bytes, 8); + assert_eq!(preview.stats.omitted_bytes, 100); + assert!(preview.output.contains("... 100 bytes omitted ...")); + assert!(preview.output.contains("abcd")); + assert!(preview.output.ends_with("efgh")); + } + #[test] fn under_limit_passthrough_chars() { let output = "short output"; @@ -222,16 +369,18 @@ mod tests { let output = "a".repeat(100); let result = truncate_output(&output, 40, TruncationMode::HeadTail); assert!(result.contains(&"a".repeat(20))); - assert!(result.contains("Tool output was truncated")); - assert!(result.contains("60 characters were removed from the middle")); + assert!(result.starts_with("Warning: truncated output (original token count: 25)")); + assert!(result.contains("... 60 bytes omitted ...")); + assert!(result.contains("approximately 15 tokens truncated")); } #[test] fn tail_mode() { let output = format!("{}BBB", "A".repeat(100)); let result = truncate_output(&output, 10, TruncationMode::Tail); - assert!(result.contains("Tool output was truncated")); - assert!(result.contains("First 93 characters were removed")); + assert!(result.starts_with("Warning: truncated output")); + assert!(result.contains("... 93 bytes omitted ...")); + assert!(result.contains("approximately 24 tokens truncated")); assert!(result.ends_with("AAAAAAABBB")); } @@ -244,7 +393,8 @@ mod tests { assert!(result.contains("line 3")); assert!(result.contains("line 18")); assert!(result.contains("line 20")); - assert!(result.contains("... 14 lines omitted ...")); + assert!(result.contains("14 lines omitted")); + assert!(result.contains("tokens truncated")); } #[test] @@ -273,7 +423,7 @@ mod tests { let mut config = SessionOptions::default(); config.tool_output_limits.insert("shell".into(), 100); let result = truncate_tool_output(&"x".repeat(1_000), "Bash", &config); - assert!(result.contains("Tool output was truncated")); + assert!(result.contains("Warning: truncated output")); } #[test] @@ -283,7 +433,7 @@ mod tests { config.tool_output_limits.insert("my_tool".into(), 100); let result = truncate_tool_output(&output, "my_tool", &config); assert!(result.len() < output.len()); - assert!(result.contains("Tool output was truncated")); + assert!(result.contains("Warning: truncated output")); } #[test] @@ -344,6 +494,6 @@ mod tests { fn truncate_output_multibyte_no_panic() { let output = "✅".repeat(100); // 300 bytes let result = truncate_output(&output, 10, TruncationMode::HeadTail); - assert!(result.contains("Tool output was truncated")); + assert!(result.contains("Warning: truncated output")); } } From 6e19fb2eec94c6799e34890af7a949245edbc862 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 13:00:58 -0400 Subject: [PATCH 4/6] Reserve event space for serialized tool output --- .../fabro-server/src/server/handler/events.rs | 7 +- lib/apps/fabro-server/src/server/tests.rs | 73 ++++++++++++++ .../fabro-agent/src/tool_execution.rs | 94 +++++++++++++++++++ lib/components/fabro-agent/src/truncation.rs | 69 ++++++++++++-- 4 files changed, 233 insertions(+), 10 deletions(-) diff --git a/lib/apps/fabro-server/src/server/handler/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index 926fefdca..247823cb7 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -1,5 +1,6 @@ use std::sync::Arc; +use axum::extract::DefaultBodyLimit; use fabro_types::{ RunEventDetailContent, RunEventDetailContentKind, RunEventDetailEnvelope, RunEventDetailResponse, @@ -15,12 +16,16 @@ use super::super::{ reject_if_archived, update_live_run_from_event, }; +const MAX_RUN_EVENT_BODY_BYTES: usize = 3 * 1024 * 1024; + pub(super) fn routes() -> Router> { Router::new() .route("/attach", get(attach_events)) .route( "/runs/{id}/events", - get(list_run_events).post(append_run_event), + get(list_run_events) + .post(append_run_event) + .layer(DefaultBodyLimit::max(MAX_RUN_EVENT_BODY_BYTES)), ) .route("/runs/{id}/events/{seq}", get(get_run_event_detail)) .route( diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index dc21a12e6..61429e6e1 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -10985,6 +10985,79 @@ async fn append_run_event_rejects_run_id_mismatch() { assert_status!(response, StatusCode::BAD_REQUEST).await; } +#[tokio::test] +async fn append_run_event_accepts_a_body_larger_than_two_mib() { + let state = test_app_state(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = create_run(&app, MINIMAL_DOT).await; + let payload = json!({ + "id": "evt-large-agent-output", + "ts": "2026-08-24T12:00:00Z", + "run_id": run_id, + "event": "agent.tool.completed", + "properties": { + "tool_name": "shell", + "tool_call_id": "call-large", + "output": "x".repeat(2 * 1024 * 1024), + "is_error": false, + "visit": 1 + } + }) + .to_string(); + assert!(payload.len() > 2 * 1024 * 1024); + assert!(payload.len() < 3 * 1024 * 1024); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(api(&format!("/runs/{run_id}/events"))) + .header("content-type", "application/json") + .body(Body::from(payload)) + .unwrap(), + ) + .await + .unwrap(); + + assert_status!(response, StatusCode::OK).await; +} + +#[tokio::test] +async fn append_run_event_rejects_a_body_larger_than_three_mib() { + let state = test_app_state(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = create_run(&app, MINIMAL_DOT).await; + let payload = json!({ + "id": "evt-oversized-agent-output", + "ts": "2026-08-24T12:00:00Z", + "run_id": run_id, + "event": "agent.tool.completed", + "properties": { + "tool_name": "shell", + "tool_call_id": "call-oversized", + "output": "x".repeat(3 * 1024 * 1024), + "is_error": false, + "visit": 1 + } + }) + .to_string(); + assert!(payload.len() > 3 * 1024 * 1024); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(api(&format!("/runs/{run_id}/events"))) + .header("content-type", "application/json") + .body(Body::from(payload)) + .unwrap(), + ) + .await + .unwrap(); + + assert_status!(response, StatusCode::PAYLOAD_TOO_LARGE).await; +} + #[tokio::test] async fn append_run_event_rejects_reserved_archive_event() { let state = test_app_state(); diff --git a/lib/components/fabro-agent/src/tool_execution.rs b/lib/components/fabro-agent/src/tool_execution.rs index 3cb8b674a..7b377e586 100644 --- a/lib/components/fabro-agent/src/tool_execution.rs +++ b/lib/components/fabro-agent/src/tool_execution.rs @@ -645,6 +645,7 @@ mod tests { use async_trait::async_trait; use fabro_llm::types::{ToolCall, ToolDefinition}; use fabro_model::AgentProfileKind; + use fabro_types::run_event::AgentToolCompletedProps; use tokio::sync::broadcast; use super::*; @@ -660,6 +661,7 @@ mod tests { use crate::test_support::MockSandbox; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource}; use crate::tools::make_shell_tool; + use crate::truncation::MAX_SERIALIZED_TOOL_OUTPUT_BYTES; use crate::types::SessionEvent; struct NamedPolicy { @@ -1026,6 +1028,98 @@ mod tests { assert!(event_output.contains(&format!("... {} bytes omitted ...", completed.4))); } + #[tokio::test] + async fn serialized_tool_output_and_full_event_stay_within_reserved_budgets() { + let mut registry = ToolRegistry::new(); + registry.register(make_echo_tool()); + let text = format!( + "HEAD{}TAIL", + "\0".repeat(MAX_RETAINED_TOOL_OUTPUT_BYTES - "echo: HEADTAIL".len()) + ); + let tc = make_tool_call("echo", "call_escaped", serde_json::json!({"text": text})); + let emitter = Emitter::new(); + let mut receiver = emitter.subscribe(); + + let result = execute_and_emit_one_tool( + &tc, + ®istry, + make_sandbox(), + None, + CancellationToken::new(), + &SessionOptions::default(), + &emitter, + "test-session", + "test-session", + None, + ) + .await; + + assert!(!result.is_error); + let completed = loop { + let event = receiver.try_recv().expect("tool completion event"); + if let AgentEvent::ToolCallCompleted { + tool_name, + tool_call_id, + output, + is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, + } = event.event + { + break ( + tool_name, + tool_call_id, + output, + is_error, + output_bytes_observed, + output_bytes_retained, + output_bytes_omitted, + ); + } + }; + + let serialized_output_bytes = serde_json::to_vec(&completed.2) + .expect("tool output serializes") + .len(); + assert!(serialized_output_bytes <= MAX_SERIALIZED_TOOL_OUTPUT_BYTES); + + let run_id = fabro_types::RunId::new(); + let run_event = fabro_types::RunEvent { + id: "evt-escaped-output".to_string(), + ts: chrono::Utc::now(), + run_id, + node_id: None, + node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, + session_id: Some("test-session".to_string()), + parent_session_id: None, + tool_call_id: Some(completed.1.clone()), + actor: None, + body: fabro_types::EventBody::AgentToolCompleted(AgentToolCompletedProps { + tool_name: completed.0, + tool_call_id: completed.1, + output: completed.2, + is_error: completed.3, + visit: 1, + output_bytes_observed: Some(completed.4 as u64), + output_bytes_retained: Some(completed.5 as u64), + output_bytes_omitted: Some(completed.6 as u64), + tool_result: None, + turn_id: None, + }), + }; + let serialized_event_bytes = serde_json::to_vec(&run_event) + .expect("run event serializes") + .len(); + assert!( + serialized_event_bytes < 2 * 1024 * 1024, + "serialized event was {serialized_event_bytes} bytes" + ); + } + #[tokio::test] async fn post_tool_use_hook_fires_on_success() { let mut registry = ToolRegistry::new(); diff --git a/lib/components/fabro-agent/src/truncation.rs b/lib/components/fabro-agent/src/truncation.rs index af85eb2c6..6e155da86 100644 --- a/lib/components/fabro-agent/src/truncation.rs +++ b/lib/components/fabro-agent/src/truncation.rs @@ -3,6 +3,7 @@ use crate::sandbox::OutputCaptureStats; use crate::tool_permissions::canonical_tool_name; pub(crate) const MAX_RETAINED_TOOL_OUTPUT_BYTES: usize = 1024 * 1024; +pub(crate) const MAX_SERIALIZED_TOOL_OUTPUT_BYTES: usize = 3 * 1024 * 1024 / 2; #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) struct RetainedToolOutput { @@ -59,7 +60,7 @@ pub(crate) fn retain_tool_output( } /// Build the final model-facing preview, including truncation notices inside -/// the total byte budget. +/// the total byte budget and JSON serialization limit. #[must_use] pub(crate) fn preview_tool_output( output: &str, @@ -69,12 +70,16 @@ pub(crate) fn preview_tool_output( let mut content_budget = max_bytes; loop { let retained = retain_tool_output(output, content_budget, previously_omitted_bytes); - if retained.stats.omitted_bytes == 0 { - return retained; - } - - let rendered = render_retained_output(&retained); - if rendered.len() <= max_bytes { + let rendered = if retained.stats.omitted_bytes == 0 { + retained.output.clone() + } else { + render_retained_output(&retained) + }; + let serialized_bytes = serialized_json_string_bytes(&rendered); + if rendered.len() <= max_bytes && serialized_bytes <= MAX_SERIALIZED_TOOL_OUTPUT_BYTES { + if retained.stats.omitted_bytes == 0 { + return retained; + } return RetainedToolOutput { output: rendered, stats: retained.stats, @@ -82,8 +87,20 @@ pub(crate) fn preview_tool_output( }; } - let excess = rendered.len().saturating_sub(max_bytes).max(1); - let next_budget = content_budget.saturating_sub(excess); + let mut next_budget = content_budget; + if rendered.len() > max_bytes { + let excess = rendered.len().saturating_sub(max_bytes); + next_budget = next_budget.min(content_budget.saturating_sub(excess)); + } + if serialized_bytes > MAX_SERIALIZED_TOOL_OUTPUT_BYTES { + let scaled_budget = (content_budget as u128) + .saturating_mul(MAX_SERIALIZED_TOOL_OUTPUT_BYTES as u128) + .checked_div(serialized_bytes as u128) + .and_then(|budget| usize::try_from(budget).ok()) + .unwrap_or(0); + next_budget = next_budget.min(scaled_budget); + } + next_budget = next_budget.min(content_budget.saturating_sub(1)); if next_budget == content_budget { return RetainedToolOutput { output: truncate_plain_output(&rendered, max_bytes, TruncationMode::HeadTail), @@ -95,6 +112,12 @@ pub(crate) fn preview_tool_output( } } +fn serialized_json_string_bytes(output: &str) -> usize { + serde_json::to_vec(output) + .expect("strings always serialize as JSON") + .len() +} + fn render_retained_output(retained: &RetainedToolOutput) -> String { let head = &retained.output[..retained.head_bytes]; let tail = &retained.output[retained.head_bytes..]; @@ -350,6 +373,34 @@ mod tests { assert!(preview.output.ends_with("efgh")); } + #[test] + fn model_preview_bounds_pathological_json_serialization() { + let output = format!( + "HEAD{}TAIL", + "\0".repeat(MAX_RETAINED_TOOL_OUTPUT_BYTES - "HEADTAIL".len()) + ); + assert_eq!(output.len(), MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert!(serialized_json_string_bytes(&output) > MAX_SERIALIZED_TOOL_OUTPUT_BYTES); + + let preview = preview_tool_output(&output, MAX_RETAINED_TOOL_OUTPUT_BYTES, 0); + let serialized_bytes = serialized_json_string_bytes(&preview.output); + + assert!(preview.output.len() <= MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert!( + serialized_bytes <= MAX_SERIALIZED_TOOL_OUTPUT_BYTES, + "serialized preview was {serialized_bytes} bytes" + ); + assert!(preview.output.starts_with("Warning: truncated output")); + assert!(preview.output.contains("HEAD")); + assert!(preview.output.ends_with("TAIL")); + assert_eq!(preview.stats.observed_bytes, MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert!(preview.stats.retained_bytes < MAX_RETAINED_TOOL_OUTPUT_BYTES); + assert_eq!( + preview.stats.omitted_bytes, + preview.stats.observed_bytes - preview.stats.retained_bytes + ); + } + #[test] fn under_limit_passthrough_chars() { let output = "short output"; From 2626e5ab4a4fe966e38cfa124625459d422c67db Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 13:55:10 -0400 Subject: [PATCH 5/6] Fix the daytona-only build of fabro-sandbox duration_to_minutes_i32 carried stacked docker and daytona cfg attributes, which combine as AND, so building with the daytona feature alone failed to find the function. fabro-workflow and fabro-cli enable daytona without docker in their production dependencies, so that combination is real. Removing the stray docker gate surfaced items that only docker-gated code uses: the ResolveError import in from_environment and four exact-checkout command builders in clone_source. Gate those on the docker feature, keeping the command builders available to clone_source's own tests under cfg(test). Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01TK3QTWQHiXhRbFwTr57LzX --- lib/components/fabro-sandbox/src/clone_source.rs | 4 ++++ lib/components/fabro-sandbox/src/from_environment.rs | 2 +- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/lib/components/fabro-sandbox/src/clone_source.rs b/lib/components/fabro-sandbox/src/clone_source.rs index 9f03a93e3..d8268ab04 100644 --- a/lib/components/fabro-sandbox/src/clone_source.rs +++ b/lib/components/fabro-sandbox/src/clone_source.rs @@ -72,6 +72,7 @@ pub(crate) fn repo_symlink_command(layout: &GitHubRepoLayout) -> String { ) } +#[cfg(any(feature = "docker", test))] pub(crate) fn exact_repository_init_command(clone_url: &str, checkout_path: &str) -> String { format!( "{git} init -- {path} && git -C {path} remote add origin {origin}", @@ -87,6 +88,7 @@ pub(crate) fn exact_repository_init_command(clone_url: &str, checkout_path: &str /// The fetch names the commit directly rather than the branch: reachability of /// the commit from the admitted branch is an admission-time invariant, not /// something this layer re-verifies. +#[cfg(any(feature = "docker", test))] pub(crate) fn exact_fetch_command( checkout_path: &str, fetch_source: &str, @@ -105,6 +107,7 @@ pub(crate) fn exact_fetch_command( /// Leading-space ` --depth N` fragment for a Git command, or empty when /// `depth` is `None` to fetch full history. +#[cfg(any(feature = "docker", test))] pub(crate) fn depth_argument(depth: Option) -> String { depth.map_or_else(String::new, |depth| format!(" --depth {depth}")) } @@ -139,6 +142,7 @@ pub(crate) fn exact_head_revision_command(checkout_path: &str) -> String { /// Check out the admitted branch and print the resulting HEAD in one shell /// command; stdout is the `rev-parse HEAD` output for [`verify_exact_head`]. +#[cfg(any(feature = "docker", test))] pub(crate) fn exact_checkout_verify_command( checkout_path: &str, branch: &str, diff --git a/lib/components/fabro-sandbox/src/from_environment.rs b/lib/components/fabro-sandbox/src/from_environment.rs index 657b4c179..825fca580 100644 --- a/lib/components/fabro-sandbox/src/from_environment.rs +++ b/lib/components/fabro-sandbox/src/from_environment.rs @@ -5,6 +5,7 @@ use std::path::{Path, PathBuf}; +#[cfg(feature = "docker")] use fabro_types::settings::ResolveError; #[cfg(feature = "daytona")] use fabro_types::settings::run::DockerfileSource as ResolvedDockerfileSource; @@ -164,7 +165,6 @@ pub fn local_working_directory_from_environment( ))) } -#[cfg(feature = "docker")] #[cfg(feature = "daytona")] fn duration_to_minutes_i32(duration: std::time::Duration) -> i32 { let minutes = duration.as_secs() / 60; From 1e284c625e67be96287b8caedf64819479210acd Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 13:56:16 -0400 Subject: [PATCH 6/6] Simplify bounded tool output capture Apply cleanups from a reuse/simplification/efficiency review of the bounded-tool-output changes: - Share one MAX_RUN_EVENT_BODY_BYTES constant in fabro-types; the server body limit, the agent's serialized-output reservation, and the event headroom test all derive from it. - Rework truncation.rs around one split_head_tail helper: drop the hand-rolled ceil_char_boundary (std's is stable), the duplicate truncate_plain_output splitter and its dead Tail arm, and the head_bytes field with its sentinel values. - Return Cow from preview_tool_output and take retain_tool_output's input by value, so untruncated output crosses the pipeline without full copies. Measure serialized JSON size with a counting writer instead of materializing the payload. - Reuse fabro-llm's byte-token estimate (now public) instead of a third copy of the 4-bytes-per-token heuristic. - Take retain_tool_result's ToolResult by value and mutate content in place; extract the triplicated error retain-emit-truncate block into finish_error_result. - Share the shell retain-and-record sequence between the native and kimi shell tools as retain_shell_output. - Move OutputCaptureBuffer::into_parts to reuse the head allocation, skip the buffer round-trip in replay_exec_result when output fits, and replace daytona's byte-iterator suffix matching with contiguous slice comparisons behind one retained_slices accessor. - Make SessionBoundEmitter's fields private. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01TK3QTWQHiXhRbFwTr57LzX --- .../fabro-server/src/server/handler/events.rs | 3 +- lib/components/fabro-agent/src/event.rs | 6 +- .../fabro-agent/src/profiles/kimi_tools.rs | 11 +- .../fabro-agent/src/tool_execution.rs | 74 +++-- lib/components/fabro-agent/src/tools.rs | 30 ++- lib/components/fabro-agent/src/truncation.rs | 252 +++++++++--------- lib/components/fabro-llm/src/token_count.rs | 12 +- .../fabro-sandbox/src/daytona/mod.rs | 64 +++-- lib/components/fabro-sandbox/src/sandbox.rs | 61 +++-- .../fabro-types/src/run_event/mod.rs | 8 + 10 files changed, 266 insertions(+), 255 deletions(-) diff --git a/lib/apps/fabro-server/src/server/handler/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index 247823cb7..01e8085eb 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -1,6 +1,7 @@ use std::sync::Arc; use axum::extract::DefaultBodyLimit; +use fabro_types::run_event::MAX_RUN_EVENT_BODY_BYTES; use fabro_types::{ RunEventDetailContent, RunEventDetailContentKind, RunEventDetailEnvelope, RunEventDetailResponse, @@ -16,8 +17,6 @@ use super::super::{ reject_if_archived, update_live_run_from_event, }; -const MAX_RUN_EVENT_BODY_BYTES: usize = 3 * 1024 * 1024; - pub(super) fn routes() -> Router> { Router::new() .route("/attach", get(attach_events)) diff --git a/lib/components/fabro-agent/src/event.rs b/lib/components/fabro-agent/src/event.rs index eaea439bd..e16f71e50 100644 --- a/lib/components/fabro-agent/src/event.rs +++ b/lib/components/fabro-agent/src/event.rs @@ -64,9 +64,9 @@ impl Default for Emitter { /// when a subagent's events are forwarded through its parent. #[derive(Clone)] pub struct SessionBoundEmitter { - pub emitter: Emitter, - pub session_id: String, - pub tool_call_id: Option, + emitter: Emitter, + session_id: String, + tool_call_id: Option, tool_output_stats: Arc>>, } diff --git a/lib/components/fabro-agent/src/profiles/kimi_tools.rs b/lib/components/fabro-agent/src/profiles/kimi_tools.rs index 6f13f2b57..07b6ec05e 100644 --- a/lib/components/fabro-agent/src/profiles/kimi_tools.rs +++ b/lib/components/fabro-agent/src/profiles/kimi_tools.rs @@ -31,9 +31,8 @@ use crate::sandbox::{GrepOptions, format_lines_numbered}; use crate::tool_registry::{RegisteredTool, ToolSource}; use crate::tools::{ DEFAULT_READ_LINES, emit_shell_process_completed, execute_grep, execute_shell_command, - grep_result_path, make_edit_file_tool, optional_usize_arg, required_str, + grep_result_path, make_edit_file_tool, optional_usize_arg, required_str, retain_shell_output, }; -use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output}; const DEFAULT_GREP_RESULTS: usize = 250; const MAX_GREP_RESULTS: usize = 2000; @@ -144,13 +143,7 @@ explicitly asked. Never run commands requiring superuser privileges unless expli let _ = write!(out, "Command failed with exit code: {code}"); } let is_success = result.is_success(); - let retained = retain_tool_output( - &out, - MAX_RETAINED_TOOL_OUTPUT_BYTES, - streaming.output_capture().omitted_bytes, - ); - ctx.record_tool_output_stats(retained.stats); - let out = retained.output; + let out = retain_shell_output(&ctx, &streaming, out); emit_shell_process_completed(&ctx, streaming).await; if is_success { Ok(out) } else { Err(out) } }) diff --git a/lib/components/fabro-agent/src/tool_execution.rs b/lib/components/fabro-agent/src/tool_execution.rs index 7b377e586..cc2943aaf 100644 --- a/lib/components/fabro-agent/src/tool_execution.rs +++ b/lib/components/fabro-agent/src/tool_execution.rs @@ -1,3 +1,4 @@ +use std::borrow::Cow; use std::sync::Arc; use fabro_llm::types::{ToolCall, ToolResult}; @@ -12,7 +13,8 @@ use crate::sandbox::{OutputCaptureStats, Sandbox}; use crate::session::ToolEnvProvider; use crate::tool_registry::{AgentEventEmitter, RegisteredTool, ToolContext, ToolRegistry}; use crate::truncation::{ - MAX_RETAINED_TOOL_OUTPUT_BYTES, preview_tool_output, truncate_tool_output, + MAX_RETAINED_TOOL_OUTPUT_BYTES, preview_tool_output, serialized_json_bytes, + truncate_tool_output, }; use crate::types::AgentEvent; @@ -267,7 +269,19 @@ fn error_tool_result_with_events( message: &str, ) -> ToolResult { emit_tool_call_started(emitter, session_id, tc); - let retained = retain_tool_result(&ToolResult::error(&tc.id, message), None); + finish_error_result(tc, emitter, session_id, config, message) +} + +/// Bound, emit, and truncate an error result for a tool call whose +/// started event was already emitted. +fn finish_error_result( + tc: &ToolCall, + emitter: &Emitter, + session_id: &str, + config: &SessionOptions, + message: &str, +) -> ToolResult { + let retained = retain_tool_result(ToolResult::error(&tc.id, message), None); emit_tool_call_result( emitter, session_id, @@ -403,15 +417,7 @@ async fn execute_and_emit_one_tool_with_lookup( emit_tool_call_started(emitter, session_id, tc); if let Some(reason) = access_denial { - let retained = retain_tool_result(&ToolResult::error(&tc.id, &reason), None); - emit_tool_call_result( - emitter, - session_id, - tc, - &retained.result, - retained.output_stats, - ); - return truncate_tool_result(&retained.result, &tc.name, config); + return finish_error_result(tc, emitter, session_id, config, &reason); } // Pre-tool-use hook @@ -423,15 +429,7 @@ async fn execute_and_emit_one_tool_with_lookup( debug!(tool = %tc.name, hook_event = "pre_tool_use", ?decision, duration_ms = elapsed, "Tool hook complete"); if let ToolHookDecision::Block { reason } = decision { - let retained = retain_tool_result(&ToolResult::error(&tc.id, &reason), None); - emit_tool_call_result( - emitter, - session_id, - tc, - &retained.result, - retained.output_stats, - ); - return truncate_tool_result(&retained.result, &tc.name, config); + return finish_error_result(tc, emitter, session_id, config, &reason); } } @@ -447,7 +445,7 @@ async fn execute_and_emit_one_tool_with_lookup( agent_tool_runtime, ) .await; - let retained = retain_tool_result(&executed.result, executed.output_stats); + let retained = retain_tool_result(executed.result, executed.output_stats); let result = retained.result; emit_tool_call_result(emitter, session_id, tc, &result, retained.output_stats); @@ -484,32 +482,25 @@ struct RetainedToolResult { /// Bound model-native tool output before it reaches hooks, events, or history. fn retain_tool_result( - result: &ToolResult, + mut result: ToolResult, previous_stats: Option, ) -> RetainedToolResult { - let (retained_content, output_stats) = match &result.content { + let output_stats = match &mut result.content { serde_json::Value::String(output) => { let previously_omitted = previous_stats.map_or(0, |stats| stats.omitted_bytes); - let retained = + let previewed = preview_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, previously_omitted); - (serde_json::Value::String(retained.output), retained.stats) - } - other => { - let byte_count = serde_json::to_vec(other) - .expect("serde_json::Value always serializes") - .len(); - (other.clone(), OutputCaptureStats::complete(byte_count)) + let stats = previewed.stats; + if let Cow::Owned(previewed_output) = previewed.output { + *output = previewed_output; + } + stats } + other => OutputCaptureStats::complete(serialized_json_bytes(other)), }; RetainedToolResult { - result: ToolResult { - tool_call_id: result.tool_call_id.clone(), - content: retained_content, - is_error: result.is_error, - image_data: result.image_data.clone(), - image_media_type: result.image_media_type.clone(), - }, + result, output_stats, } } @@ -645,7 +636,7 @@ mod tests { use async_trait::async_trait; use fabro_llm::types::{ToolCall, ToolDefinition}; use fabro_model::AgentProfileKind; - use fabro_types::run_event::AgentToolCompletedProps; + use fabro_types::run_event::{AgentToolCompletedProps, MAX_RUN_EVENT_BODY_BYTES}; use tokio::sync::broadcast; use super::*; @@ -1114,8 +1105,11 @@ mod tests { let serialized_event_bytes = serde_json::to_vec(&run_event) .expect("run event serializes") .len(); + // Leave at least 1 MiB of envelope headroom under the server's + // run-event body limit. + let event_body_budget = MAX_RUN_EVENT_BODY_BYTES - 1024 * 1024; assert!( - serialized_event_bytes < 2 * 1024 * 1024, + serialized_event_bytes < event_body_budget, "serialized event was {serialized_event_bytes} bytes" ); } diff --git a/lib/components/fabro-agent/src/tools.rs b/lib/components/fabro-agent/src/tools.rs index 24ab550d2..24ad2c80f 100644 --- a/lib/components/fabro-agent/src/tools.rs +++ b/lib/components/fabro-agent/src/tools.rs @@ -13,7 +13,7 @@ use tokio::task; use crate::config::NativeToolOptions; use crate::sandbox::{ExecStreamingResult, GrepOptions}; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource}; -use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, RetainedToolOutput, retain_tool_output}; +use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output}; use crate::types::AgentEvent; use crate::web_search::{SearchBackend, make_web_search_tool}; @@ -320,15 +320,29 @@ pub(crate) async fn run_shell_command( cwd: Option<&str>, ) -> Result { let streaming = execute_shell_command(ctx, command, timeout_ms, cwd).await?; - let retained = render_shell_result(&streaming); - ctx.record_tool_output_stats(retained.stats); - let text = retained.output; + let text = retain_shell_output(ctx, &streaming, 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) } } +/// Bound rendered shell output to the retention budget and record the capture +/// stats for the executing tool call. +pub(crate) fn retain_shell_output( + ctx: &ToolContext, + streaming: &ExecStreamingResult, + output: String, +) -> String { + let retained = retain_tool_output( + output, + MAX_RETAINED_TOOL_OUTPUT_BYTES, + streaming.output_capture().omitted_bytes, + ); + ctx.record_tool_output_stats(retained.stats); + retained.output +} + /// 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. @@ -372,7 +386,7 @@ pub(crate) async fn emit_shell_process_completed( /// 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) -> RetainedToolOutput { +fn render_shell_result(streaming: &ExecStreamingResult) -> String { let result = &streaming.result; let mut output = format!( "Termination: {}\nExit code: {}\nDuration: {}ms\n", @@ -392,11 +406,7 @@ fn render_shell_result(streaming: &ExecStreamingResult) -> RetainedToolOutput { } else if !result.stdout.is_empty() { let _ = write!(output, "output (combined):\n{}\n", result.stdout); } - retain_tool_output( - &output, - MAX_RETAINED_TOOL_OUTPUT_BYTES, - streaming.output_capture().omitted_bytes, - ) + output } #[must_use] diff --git a/lib/components/fabro-agent/src/truncation.rs b/lib/components/fabro-agent/src/truncation.rs index 6e155da86..719cda287 100644 --- a/lib/components/fabro-agent/src/truncation.rs +++ b/lib/components/fabro-agent/src/truncation.rs @@ -1,15 +1,43 @@ +use std::borrow::Cow; + +use fabro_llm::token_count; +use fabro_types::run_event::MAX_RUN_EVENT_BODY_BYTES; +use serde::Serialize; + use crate::config::SessionOptions; use crate::sandbox::OutputCaptureStats; use crate::tool_permissions::canonical_tool_name; pub(crate) const MAX_RETAINED_TOOL_OUTPUT_BYTES: usize = 1024 * 1024; -pub(crate) const MAX_SERIALIZED_TOOL_OUTPUT_BYTES: usize = 3 * 1024 * 1024 / 2; +/// Reserve half the run-event body limit for serialized tool output; the +/// other half is headroom for the rest of the event envelope. +pub(crate) const MAX_SERIALIZED_TOOL_OUTPUT_BYTES: usize = MAX_RUN_EVENT_BODY_BYTES / 2; -#[derive(Debug, Clone, PartialEq, Eq)] +#[derive(Debug)] pub(crate) struct RetainedToolOutput { pub output: String, pub stats: OutputCaptureStats, - head_bytes: usize, +} + +/// Model-facing preview of a tool output. Borrows the input when no +/// truncation notice was needed. +#[derive(Debug)] +pub(crate) struct PreviewedToolOutput<'a> { + pub output: Cow<'a, str>, + pub stats: OutputCaptureStats, +} + +/// Boundaries of an equal-sized UTF-8 head and tail fitting `max_bytes`, or +/// `None` when `output` already fits. +fn split_head_tail(output: &str, max_bytes: usize) -> Option<(usize, usize)> { + if output.len() <= max_bytes { + return None; + } + let head_budget = max_bytes / 2; + let tail_budget = max_bytes - head_budget; + let head_end = output.floor_char_boundary(head_budget); + let tail_start = output.ceil_char_boundary(output.len() - tail_budget); + Some((head_end, tail_start)) } /// Keep an equal-sized UTF-8 prefix and suffix within a byte budget. @@ -18,44 +46,34 @@ pub(crate) struct RetainedToolOutput { /// discarded before the rendered result was assembled. #[must_use] pub(crate) fn retain_tool_output( - output: &str, + output: String, max_bytes: usize, previously_omitted_bytes: usize, ) -> RetainedToolOutput { let observed_bytes = output.len().saturating_add(previously_omitted_bytes); - if output.len() <= max_bytes { + let Some((head_end, tail_start)) = split_head_tail(&output, max_bytes) else { return RetainedToolOutput { - output: output.to_string(), - stats: OutputCaptureStats { + stats: OutputCaptureStats { observed_bytes, retained_bytes: output.len(), omitted_bytes: previously_omitted_bytes, }, - head_bytes: if previously_omitted_bytes == 0 { - output.len() - } else { - output.floor_char_boundary(output.len() / 2) - }, + output, }; - } + }; - let head_budget = max_bytes / 2; - let tail_budget = max_bytes.saturating_sub(head_budget); - let head_end = output.floor_char_boundary(head_budget); - let tail_start = ceil_char_boundary(output, output.len().saturating_sub(tail_budget)); - let retained_bytes = head_end.saturating_add(output.len().saturating_sub(tail_start)); + let retained_bytes = head_end + (output.len() - tail_start); let mut retained = String::with_capacity(retained_bytes); retained.push_str(&output[..head_end]); retained.push_str(&output[tail_start..]); RetainedToolOutput { - output: retained, - stats: OutputCaptureStats { + output: retained, + stats: OutputCaptureStats { observed_bytes, retained_bytes, omitted_bytes: observed_bytes.saturating_sub(retained_bytes), }, - head_bytes: head_end, } } @@ -66,30 +84,64 @@ pub(crate) fn preview_tool_output( output: &str, max_bytes: usize, previously_omitted_bytes: usize, -) -> RetainedToolOutput { +) -> PreviewedToolOutput<'_> { + let observed_bytes = output.len().saturating_add(previously_omitted_bytes); let mut content_budget = max_bytes; loop { - let retained = retain_tool_output(output, content_budget, previously_omitted_bytes); - let rendered = if retained.stats.omitted_bytes == 0 { - retained.output.clone() + let (head_end, tail_start, stats) = + if let Some((head_end, tail_start)) = split_head_tail(output, content_budget) { + let retained_bytes = head_end + (output.len() - tail_start); + (head_end, tail_start, OutputCaptureStats { + observed_bytes, + retained_bytes, + omitted_bytes: observed_bytes.saturating_sub(retained_bytes), + }) + } else { + // The whole output fits. A notice is still rendered when the + // stream itself omitted bytes; equal-sized retention keeps + // that omission gap at the midpoint. + let mid = output.floor_char_boundary(output.len() / 2); + (mid, mid, OutputCaptureStats { + observed_bytes, + retained_bytes: output.len(), + omitted_bytes: previously_omitted_bytes, + }) + }; + let rendered: Cow<'_, str> = if stats.omitted_bytes == 0 { + Cow::Borrowed(output) } else { - render_retained_output(&retained) + Cow::Owned(render_truncated_segments( + &output[..head_end], + &output[tail_start..], + stats, + None, + )) }; - let serialized_bytes = serialized_json_string_bytes(&rendered); + let serialized_bytes = serialized_json_bytes(rendered.as_ref()); if rendered.len() <= max_bytes && serialized_bytes <= MAX_SERIALIZED_TOOL_OUTPUT_BYTES { - if retained.stats.omitted_bytes == 0 { - return retained; - } - return RetainedToolOutput { - output: rendered, - stats: retained.stats, - head_bytes: 0, + return PreviewedToolOutput { + output: rendered, + stats, }; } - let mut next_budget = content_budget; + let Some(reduced_budget) = content_budget.checked_sub(1) else { + // The content budget is exhausted and the notice text alone still + // overflows. Hard-cut the rendered notice to fit. + let output = match split_head_tail(&rendered, max_bytes) { + Some((head_end, tail_start)) => { + format!("{}{}", &rendered[..head_end], &rendered[tail_start..]) + } + None => rendered.into_owned(), + }; + return PreviewedToolOutput { + output: Cow::Owned(output), + stats, + }; + }; + let mut next_budget = reduced_budget; if rendered.len() > max_bytes { - let excess = rendered.len().saturating_sub(max_bytes); + let excess = rendered.len() - max_bytes; next_budget = next_budget.min(content_budget.saturating_sub(excess)); } if serialized_bytes > MAX_SERIALIZED_TOOL_OUTPUT_BYTES { @@ -100,28 +152,27 @@ pub(crate) fn preview_tool_output( .unwrap_or(0); next_budget = next_budget.min(scaled_budget); } - next_budget = next_budget.min(content_budget.saturating_sub(1)); - if next_budget == content_budget { - return RetainedToolOutput { - output: truncate_plain_output(&rendered, max_bytes, TruncationMode::HeadTail), - stats: retained.stats, - head_bytes: 0, - }; - } content_budget = next_budget; } } -fn serialized_json_string_bytes(output: &str) -> usize { - serde_json::to_vec(output) - .expect("strings always serialize as JSON") - .len() -} +/// Serialized JSON size in bytes, counted without materializing the payload. +pub(crate) fn serialized_json_bytes(value: &T) -> usize { + struct CountingWriter(usize); + impl std::io::Write for CountingWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0 += buf.len(); + Ok(buf.len()) + } -fn render_retained_output(retained: &RetainedToolOutput) -> String { - let head = &retained.output[..retained.head_bytes]; - let tail = &retained.output[retained.head_bytes..]; - render_truncated_segments(head, tail, retained.stats, None) + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + let mut writer = CountingWriter(0); + serde_json::to_writer(&mut writer, value).expect("JSON tool output always serializes"); + writer.0 } fn render_truncated_segments( @@ -130,8 +181,8 @@ fn render_truncated_segments( stats: OutputCaptureStats, line_count_omitted: Option, ) -> String { - let original_tokens = approximate_tokens(stats.observed_bytes); - let omitted_tokens = approximate_tokens(stats.omitted_bytes); + let original_tokens = token_count::estimate_byte_tokens(stats.observed_bytes); + let omitted_tokens = token_count::estimate_byte_tokens(stats.omitted_bytes); let middle_marker = line_count_omitted.map_or_else( || format!("... approximately {omitted_tokens} tokens truncated ..."), |lines| { @@ -146,18 +197,6 @@ fn render_truncated_segments( ) } -fn approximate_tokens(bytes: usize) -> usize { - bytes.div_ceil(4) -} - -fn ceil_char_boundary(output: &str, index: usize) -> usize { - let mut index = index.min(output.len()); - while index < output.len() && !output.is_char_boundary(index) { - index += 1; - } - index -} - #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum TruncationMode { HeadTail, @@ -193,63 +232,28 @@ fn default_truncation_mode(tool_name: &str) -> TruncationMode { #[must_use] pub fn truncate_output(output: &str, max_chars: usize, mode: TruncationMode) -> String { - if output.len() <= max_chars { + let Some((head_end, tail_start)) = split_head_tail(output, max_chars) else { return output.to_string(); - } + }; - match mode { - TruncationMode::HeadTail => { - let half = max_chars / 2; - let head_end = output.floor_char_boundary(half); - let tail_start = ceil_char_boundary(output, output.len().saturating_sub(half)); - let head = &output[..head_end]; - let tail = &output[tail_start..]; - let retained_bytes = head.len().saturating_add(tail.len()); - render_truncated_segments( - head, - tail, - OutputCaptureStats { - observed_bytes: output.len(), - retained_bytes, - omitted_bytes: output.len().saturating_sub(retained_bytes), - }, - None, - ) - } + let (head, tail) = match mode { + TruncationMode::HeadTail => (&output[..head_end], &output[tail_start..]), TruncationMode::Tail => { - let tail_start = ceil_char_boundary(output, output.len().saturating_sub(max_chars)); - let tail = &output[tail_start..]; - render_truncated_segments( - "", - tail, - OutputCaptureStats { - observed_bytes: output.len(), - retained_bytes: tail.len(), - omitted_bytes: output.len().saturating_sub(tail.len()), - }, - None, - ) + let tail_start = output.ceil_char_boundary(output.len() - max_chars); + ("", &output[tail_start..]) } - } -} - -fn truncate_plain_output(output: &str, max_bytes: usize, mode: TruncationMode) -> String { - if output.len() <= max_bytes { - return output.to_string(); - } - - match mode { - TruncationMode::HeadTail => { - let half = max_bytes / 2; - let head_end = output.floor_char_boundary(half); - let tail_start = ceil_char_boundary(output, output.len().saturating_sub(half)); - format!("{}{}", &output[..head_end], &output[tail_start..]) - } - TruncationMode::Tail => { - let tail_start = ceil_char_boundary(output, output.len().saturating_sub(max_bytes)); - output[tail_start..].to_string() - } - } + }; + let retained_bytes = head.len().saturating_add(tail.len()); + render_truncated_segments( + head, + tail, + OutputCaptureStats { + observed_bytes: output.len(), + retained_bytes, + omitted_bytes: output.len().saturating_sub(retained_bytes), + }, + None, + ) } #[must_use] @@ -316,7 +320,7 @@ mod tests { #[test] fn retained_tool_output_keeps_equal_head_and_tail() { - let retained = retain_tool_output("abcdefghijkl", 8, 0); + let retained = retain_tool_output("abcdefghijkl".to_string(), 8, 0); assert_eq!(retained.output, "abcdijkl"); assert_eq!(retained.stats.observed_bytes, 12); @@ -326,7 +330,7 @@ mod tests { #[test] fn retained_tool_output_stays_within_budget_at_utf8_boundaries() { - let retained = retain_tool_output("aa😀😀zz", 7, 3); + let retained = retain_tool_output("aa😀😀zz".to_string(), 7, 3); assert!(retained.output.len() <= 7, "{}", retained.output.len()); assert!(retained.output.starts_with("aa")); @@ -380,10 +384,10 @@ mod tests { "\0".repeat(MAX_RETAINED_TOOL_OUTPUT_BYTES - "HEADTAIL".len()) ); assert_eq!(output.len(), MAX_RETAINED_TOOL_OUTPUT_BYTES); - assert!(serialized_json_string_bytes(&output) > MAX_SERIALIZED_TOOL_OUTPUT_BYTES); + assert!(serialized_json_bytes(output.as_str()) > MAX_SERIALIZED_TOOL_OUTPUT_BYTES); let preview = preview_tool_output(&output, MAX_RETAINED_TOOL_OUTPUT_BYTES, 0); - let serialized_bytes = serialized_json_string_bytes(&preview.output); + let serialized_bytes = serialized_json_bytes(preview.output.as_ref()); assert!(preview.output.len() <= MAX_RETAINED_TOOL_OUTPUT_BYTES); assert!( diff --git a/lib/components/fabro-llm/src/token_count.rs b/lib/components/fabro-llm/src/token_count.rs index 58d9ec74d..fa64abda7 100644 --- a/lib/components/fabro-llm/src/token_count.rs +++ b/lib/components/fabro-llm/src/token_count.rs @@ -221,7 +221,7 @@ impl Estimator { let mut tokens = estimate_text_tokens(&result.tool_call_id) + estimate_json_tokens(&result.content); if let Some(image_data) = &result.image_data { - tokens += estimate_embedded_bytes(image_data.len()); + tokens += estimate_byte_tokens(image_data.len()); self.warn( MEDIA_ESTIMATE_WARNING, "Media content couldn't be precisely tokenized; total is approximate.", @@ -243,7 +243,7 @@ impl Estimator { + image .data .as_ref() - .map_or(2000, |data| estimate_embedded_bytes(data.len()).max(2000)) + .map_or(2000, |data| estimate_byte_tokens(data.len()).max(2000)) } fn estimate_audio(&mut self, audio: &AudioData) -> usize { @@ -252,7 +252,7 @@ impl Estimator { + audio .data .as_ref() - .map_or(2000, |data| estimate_embedded_bytes(data.len())) + .map_or(2000, |data| estimate_byte_tokens(data.len())) } fn estimate_document(&mut self, document: &DocumentData) -> usize { @@ -265,7 +265,7 @@ impl Estimator { + document .data .as_ref() - .map_or(2000, |data| estimate_embedded_bytes(data.len())) + .map_or(2000, |data| estimate_byte_tokens(data.len())) } fn estimate_media_common(&mut self, url: Option<&str>, media_type: Option<&str>) -> usize { @@ -292,7 +292,9 @@ fn estimate_tool(tool: &ToolDefinition) -> usize { + estimate_json_tokens(&tool.parameters) } -fn estimate_embedded_bytes(byte_len: usize) -> usize { +/// Approximate the token cost of a raw byte payload (4 bytes per token). +#[must_use] +pub fn estimate_byte_tokens(byte_len: usize) -> usize { byte_len.div_ceil(4) } diff --git a/lib/components/fabro-sandbox/src/daytona/mod.rs b/lib/components/fabro-sandbox/src/daytona/mod.rs index 6da9c6e72..a7a0a8d12 100644 --- a/lib/components/fabro-sandbox/src/daytona/mod.rs +++ b/lib/components/fabro-sandbox/src/daytona/mod.rs @@ -33,8 +33,8 @@ use crate::push_credentials::{self, PushCredentialState}; use crate::redact::redact_auth_url; use crate::sandbox::{ self, BASH_ENV_VAR, BASH_PROBE_MARKER, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, - OutputCaptureBuffer, REMOTE_BASH, REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, optional_timeout, - resolve_path, validate_bash_probe, + OutputCaptureBuffer, OutputCaptureStats, REMOTE_BASH, REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, + optional_timeout, resolve_path, validate_bash_probe, }; use crate::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingRequest, ExecStreamingResult, @@ -2375,20 +2375,8 @@ impl Sandbox for DaytonaSandbox { .await?; } - let (stdout, stdout_capture) = { - let stdout_seen = stdout_seen.lock().await; - ( - String::from_utf8_lossy(&stdout_seen.to_bytes()).into_owned(), - stdout_seen.stats(), - ) - }; - let (stderr, stderr_capture) = { - let stderr_seen = stderr_seen.lock().await; - ( - String::from_utf8_lossy(&stderr_seen.to_bytes()).into_owned(), - stderr_seen.stats(), - ) - }; + let (stdout, stdout_capture) = drain_captured_stream(&stdout_seen).await; + let (stderr, stderr_capture) = drain_captured_stream(&stderr_seen).await; let result = ExecStreamingResult { result: ExecResult { @@ -2914,7 +2902,7 @@ async fn append_missing_log_suffix( } let mut seen = seen.lock().await; - let offset = captured_log_suffix_offset(&seen, final_bytes); + let offset = captured_log_suffix_offset(&mut seen, final_bytes); if offset >= final_bytes.len() { return Ok(()); } @@ -2928,21 +2916,34 @@ async fn append_missing_log_suffix( } } -fn captured_log_suffix_offset(seen: &OutputCaptureBuffer, final_bytes: &[u8]) -> usize { +/// Take the captured stream bytes out of their shared buffer as a lossy +/// string, avoiding a copy when the bytes are valid UTF-8. +async fn drain_captured_stream( + seen: &Arc>, +) -> (String, OutputCaptureStats) { + let buffer = { + let mut seen = seen.lock().await; + std::mem::replace(&mut *seen, OutputCaptureBuffer::new(None)) + }; + let (bytes, stats) = buffer.into_parts(); + let text = match String::from_utf8(bytes) { + Ok(text) => text, + Err(err) => String::from_utf8_lossy(err.as_bytes()).into_owned(), + }; + (text, stats) +} + +fn captured_log_suffix_offset(seen: &mut OutputCaptureBuffer, final_bytes: &[u8]) -> usize { let stats = seen.stats(); if stats.omitted_bytes == 0 { return missing_log_suffix_offset(&seen.to_bytes(), final_bytes); } - let observed_bytes = seen.observed_bytes(); - let head = seen.retained_head(); - let tail = seen.retained_tail(); + let observed_bytes = stats.observed_bytes; + let (head, tail) = seen.retained_slices(); if final_bytes.len() >= observed_bytes && final_bytes.starts_with(head) - && tail.iter().copied().eq(final_bytes - [observed_bytes.saturating_sub(tail.len())..observed_bytes] - .iter() - .copied()) + && tail == &final_bytes[observed_bytes.saturating_sub(tail.len())..observed_bytes] { return observed_bytes; } @@ -2952,12 +2953,7 @@ fn captured_log_suffix_offset(seen: &OutputCaptureBuffer, final_bytes: &[u8]) -> let max_overlap = tail.len().min(final_bytes.len()); for overlap in (1..=max_overlap).rev() { - if tail - .iter() - .skip(tail.len() - overlap) - .copied() - .eq(final_bytes[..overlap].iter().copied()) - { + if tail[tail.len() - overlap..] == final_bytes[..overlap] { return overlap; } } @@ -4791,9 +4787,9 @@ mod tests { let mut seen = OutputCaptureBuffer::new(Some(6)); seen.push(b"abcdefgh"); - assert_eq!(captured_log_suffix_offset(&seen, b"abcdefghij"), 8); - assert_eq!(captured_log_suffix_offset(&seen, b"abcdefgh"), 8); - assert_eq!(captured_log_suffix_offset(&seen, b"abcd"), 4); + assert_eq!(captured_log_suffix_offset(&mut seen, b"abcdefghij"), 8); + assert_eq!(captured_log_suffix_offset(&mut seen, b"abcdefgh"), 8); + assert_eq!(captured_log_suffix_offset(&mut seen, b"abcd"), 4); } #[test] diff --git a/lib/components/fabro-sandbox/src/sandbox.rs b/lib/components/fabro-sandbox/src/sandbox.rs index 473cfa71b..a08094df6 100644 --- a/lib/components/fabro-sandbox/src/sandbox.rs +++ b/lib/components/fabro-sandbox/src/sandbox.rs @@ -870,38 +870,34 @@ impl OutputCaptureBuffer { #[cfg(feature = "daytona")] #[must_use] pub(crate) fn to_bytes(&self) -> Vec { - let stats = self.stats(); - let mut bytes = Vec::with_capacity(stats.retained_bytes); + let mut bytes = Vec::with_capacity(self.head.len().saturating_add(self.tail.len())); bytes.extend_from_slice(&self.head); - bytes.extend(self.tail.iter().copied()); + let (front, back) = self.tail.as_slices(); + bytes.extend_from_slice(front); + bytes.extend_from_slice(back); bytes } #[must_use] pub(crate) fn into_parts(self) -> (Vec, OutputCaptureStats) { let stats = self.stats(); - let mut bytes = Vec::with_capacity(stats.retained_bytes); - bytes.extend(self.head); - bytes.extend(self.tail); + let Self { + head: mut bytes, + tail, + .. + } = self; + let (front, back) = tail.as_slices(); + bytes.extend_from_slice(front); + bytes.extend_from_slice(back); (bytes, stats) } + /// Retained bytes as two contiguous slices: the stable head, then the + /// rolling tail. #[cfg(feature = "daytona")] #[must_use] - pub(crate) fn observed_bytes(&self) -> usize { - self.observed_bytes - } - - #[cfg(feature = "daytona")] - #[must_use] - pub(crate) fn retained_head(&self) -> &[u8] { - &self.head - } - - #[cfg(feature = "daytona")] - #[must_use] - pub(crate) fn retained_tail(&self) -> &VecDeque { - &self.tail + pub(crate) fn retained_slices(&mut self) -> (&[u8], &[u8]) { + (&self.head, self.tail.make_contiguous()) } } @@ -1008,14 +1004,8 @@ pub(crate) async fn replay_exec_result( .await?; } } - let mut stdout_capture = OutputCaptureBuffer::new(stream_output_bytes_cap); - stdout_capture.push(result.stdout.as_bytes()); - let mut stderr_capture = OutputCaptureBuffer::new(stream_output_bytes_cap); - stderr_capture.push(result.stderr.as_bytes()); - let (stdout, stdout_capture) = stdout_capture.into_parts(); - let (stderr, stderr_capture) = stderr_capture.into_parts(); - result.stdout = String::from_utf8_lossy(&stdout).into_owned(); - result.stderr = String::from_utf8_lossy(&stderr).into_owned(); + let stdout_capture = capture_replayed_stream(&mut result.stdout, stream_output_bytes_cap); + let stderr_capture = capture_replayed_stream(&mut result.stderr, stream_output_bytes_cap); Ok(ExecStreamingResult { result, @@ -1026,6 +1016,21 @@ pub(crate) async fn replay_exec_result( }) } +/// Bound one replayed stream in place, leaving it untouched when it already +/// fits the cap. +fn capture_replayed_stream(text: &mut String, cap: Option) -> OutputCaptureStats { + match cap { + Some(cap) if text.len() > cap => { + let mut buffer = OutputCaptureBuffer::new(Some(cap)); + buffer.push(text.as_bytes()); + let (bytes, stats) = buffer.into_parts(); + *text = String::from_utf8_lossy(&bytes).into_owned(); + stats + } + _ => OutputCaptureStats::complete(text.len()), + } +} + pub struct StdioProcess { pub stdin: Pin>, pub stdout: Pin>, diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index 30d56a110..4dc457455 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -22,6 +22,14 @@ pub use todo::*; use crate::{ParallelBranchId, Principal, RunId, StageId}; +/// Maximum accepted body size for `POST /runs/{id}/events`. +/// +/// Producers that embed large payloads in an event (serialized tool output in +/// particular) must budget against this limit, leaving headroom for the rest +/// of the event envelope. The agent layer reserves half of it for serialized +/// tool output. +pub const MAX_RUN_EVENT_BODY_BYTES: usize = 3 * 1024 * 1024; + #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum RunNoticeLevel {