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/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index fb9ff4c1d..c032d57d8 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -1,5 +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, @@ -20,7 +22,9 @@ pub(super) fn routes() -> Router> { .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/handler/sessions.rs b/lib/apps/fabro-server/src/server/handler/sessions.rs index 773b2edf6..7225fca3d 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/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index a704795ef..0c4c63515 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -11408,6 +11408,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/event.rs b/lib/components/fabro-agent/src/event.rs index f2dafe507..e16f71e50 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, + emitter: Emitter, + session_id: String, + 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/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..07b6ec05e 100644 --- a/lib/components/fabro-agent/src/profiles/kimi_tools.rs +++ b/lib/components/fabro-agent/src/profiles/kimi_tools.rs @@ -31,7 +31,7 @@ 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, }; const DEFAULT_GREP_RESULTS: usize = 250; @@ -143,6 +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 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/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..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}; @@ -8,10 +9,13 @@ 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::truncate_tool_output; +use crate::truncation::{ + MAX_RETAINED_TOOL_OUTPUT_BYTES, preview_tool_output, serialized_json_bytes, + truncate_tool_output, +}; use crate::types::AgentEvent; /// Execute tool calls, choosing parallel or sequential based on `parallel` @@ -265,9 +269,27 @@ fn error_tool_result_with_events( message: &str, ) -> ToolResult { emit_tool_call_started(emitter, session_id, tc); - let result = ToolResult::error(&tc.id, message); - emit_tool_call_result(emitter, session_id, tc, &result); - truncate_tool_result(&result, &tc.name, config) + 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, + 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 +300,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 +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 result = ToolResult::error(&tc.id, &reason); - emit_tool_call_result(emitter, session_id, tc, &result); - return truncate_tool_result(&result, &tc.name, config); + return finish_error_result(tc, emitter, session_id, config, &reason); } // Pre-tool-use hook @@ -400,13 +429,11 @@ 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); - emit_tool_call_result(emitter, session_id, tc, &result); - return truncate_tool_result(&result, &tc.name, config); + return finish_error_result(tc, emitter, session_id, config, &reason); } } - let result = execute_one_tool( + let executed = execute_one_tool( tc, registered_tool, env, @@ -418,8 +445,10 @@ async fn execute_and_emit_one_tool_with_lookup( 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 { @@ -446,6 +475,41 @@ 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( + mut result: ToolResult, + previous_stats: Option, +) -> RetainedToolResult { + 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 previewed = + preview_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, previously_omitted); + 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, + output_stats, + } +} + +struct ExecutedToolResult { + result: ToolResult, + output_stats: Option, +} + /// Execute a single tool call: argument validation and execution. #[allow( clippy::too_many_arguments, @@ -461,23 +525,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, @@ -488,14 +556,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, + }, } } @@ -558,6 +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, MAX_RUN_EVENT_BODY_BYTES}; use tokio::sync::broadcast; use super::*; @@ -573,6 +652,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 { @@ -877,6 +957,163 @@ 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!(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"); + 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, + ); + } + }; + 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!(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] + 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(); + // 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 < event_body_budget, + "serialized event was {serialized_event_bytes} bytes" + ); + } + #[tokio::test] async fn post_tool_use_hook_fires_on_success() { let mut registry = ToolRegistry::new(); @@ -1171,6 +1408,76 @@ 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 { + 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, + *output_bytes_omitted, + )), + _ => None, + }) + .expect("tool completion event"); + 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] 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 6d9f673e1..24ad2c80f 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, 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,13 +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 text = render_shell_result(&streaming); + 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. @@ -340,6 +358,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 { @@ -358,6 +377,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, }); } @@ -1093,11 +1115,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) } } @@ -1313,11 +1335,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"); @@ -1391,7 +1418,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 d7fbe3ffc..719cda287 100644 --- a/lib/components/fabro-agent/src/truncation.rs +++ b/lib/components/fabro-agent/src/truncation.rs @@ -1,6 +1,202 @@ +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; +/// 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)] +pub(crate) struct RetainedToolOutput { + pub output: String, + pub stats: OutputCaptureStats, +} + +/// 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. +/// +/// `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: String, + max_bytes: usize, + previously_omitted_bytes: usize, +) -> RetainedToolOutput { + let observed_bytes = output.len().saturating_add(previously_omitted_bytes); + let Some((head_end, tail_start)) = split_head_tail(&output, max_bytes) else { + return RetainedToolOutput { + stats: OutputCaptureStats { + observed_bytes, + retained_bytes: output.len(), + omitted_bytes: previously_omitted_bytes, + }, + output, + }; + }; + + 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 { + observed_bytes, + retained_bytes, + omitted_bytes: observed_bytes.saturating_sub(retained_bytes), + }, + } +} + +/// Build the final model-facing preview, including truncation notices inside +/// the total byte budget and JSON serialization limit. +#[must_use] +pub(crate) fn preview_tool_output( + output: &str, + max_bytes: usize, + previously_omitted_bytes: usize, +) -> PreviewedToolOutput<'_> { + let observed_bytes = output.len().saturating_add(previously_omitted_bytes); + let mut content_budget = max_bytes; + loop { + 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 { + Cow::Owned(render_truncated_segments( + &output[..head_end], + &output[tail_start..], + stats, + None, + )) + }; + let serialized_bytes = serialized_json_bytes(rendered.as_ref()); + if rendered.len() <= max_bytes && serialized_bytes <= MAX_SERIALIZED_TOOL_OUTPUT_BYTES { + return PreviewedToolOutput { + output: rendered, + stats, + }; + } + + 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() - 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); + } + content_budget = next_budget; + } +} + +/// 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 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( + head: &str, + tail: &str, + stats: OutputCaptureStats, + line_count_omitted: Option, +) -> String { + 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| { + 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 + ) +} + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum TruncationMode { HeadTail, @@ -36,34 +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(); - } + }; - 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 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 (head, tail) = match mode { + TruncationMode::HeadTail => (&output[..head_end], &output[tail_start..]), TruncationMode::Tail => { - let tail_start = output.floor_char_boundary(output.len() - 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}" - ) + let tail_start = output.ceil_char_boundary(output.len() - max_chars); + ("", &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, + ) } #[must_use] @@ -73,15 +263,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), ) } @@ -121,6 +318,93 @@ 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".to_string(), 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".to_string(), 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 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 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_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_bytes(preview.output.as_ref()); + + 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"; @@ -140,16 +424,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")); } @@ -162,7 +448,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] @@ -191,7 +478,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] @@ -201,7 +488,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] @@ -262,6 +549,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")); } } 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-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/clone_source.rs b/lib/components/fabro-sandbox/src/clone_source.rs index b3c5c9dd5..7c787b366 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. No layer proves /// that the submitted commit belongs to the submitted branch: the branch names /// the working branch, while a fetchable exact commit is checked out as-is. +#[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/daytona/mod.rs b/lib/components/fabro-sandbox/src/daytona/mod.rs index 14af37bd7..a7a0a8d12 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, OutputCaptureStats, 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,8 @@ 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) = drain_captured_stream(&stdout_seen).await; + let (stderr, stderr_capture) = drain_captured_stream(&stderr_seen).await; let result = ExecStreamingResult { result: ExecResult { @@ -2384,6 +2390,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 +2894,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 +2902,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(&mut 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 +2916,50 @@ async fn append_missing_log_suffix( } } +/// 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 = stats.observed_bytes; + let (head, tail) = seen.retained_slices(); + if final_bytes.len() >= observed_bytes + && final_bytes.starts_with(head) + && tail == &final_bytes[observed_bytes.saturating_sub(tail.len())..observed_bytes] + { + 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[tail.len() - overlap..] == final_bytes[..overlap] { + 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 +4782,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(&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] 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/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; 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..a08094df6 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,139 @@ 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 mut bytes = Vec::with_capacity(self.head.len().saturating_add(self.tail.len())); + bytes.extend_from_slice(&self.head); + 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 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 retained_slices(&mut self) -> (&[u8], &[u8]) { + (&self.head, self.tail.make_contiguous()) + } } pub type CommandOutputCallback = Arc< @@ -785,13 +918,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 +941,7 @@ impl<'a> ExecStreamingRequest<'a> { cancel_token: None, stdin: None, output_callback: None, + stream_output_bytes_cap: None, } } } @@ -846,9 +983,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,13 +1004,33 @@ pub(crate) async fn replay_exec_result( .await?; } } + 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, streams_separated, live_streaming: false, + stdout_capture, + stderr_capture, }) } +/// 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>, @@ -1179,7 +1337,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 +2681,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( 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 3407ab92f..c8fc9716e 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -1762,13 +1762,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 8d0827835..34db0184b 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. /// @@ -692,12 +696,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, }), @@ -707,6 +717,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, @@ -714,6 +727,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, }, ), @@ -1622,11 +1638,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()), @@ -1653,6 +1672,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..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 { @@ -2458,6 +2466,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 +2483,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 +2497,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 +2516,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()); + } }