From ddcdafa06b49098444a117b18d91a5116d0b897b Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 24 Aug 2026 12:34:43 -0400 Subject: [PATCH] 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(