mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
Merge pull request #793 from fabro-sh/feature/bounded-agent-tool-output
Bound oversized agent tool output
This commit is contained in:
commit
8a24046b94
32 changed files with 1467 additions and 219 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}),
|
||||
|
|
|
|||
|
|
@ -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<Arc<AppState>> {
|
|||
.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(
|
||||
|
|
|
|||
|
|
@ -1271,6 +1271,9 @@ fn agent_event_payload(event_turn_id: TurnId, event: AgentEvent) -> Option<Event
|
|||
tool_call_id,
|
||||
output,
|
||||
is_error,
|
||||
output_bytes_observed,
|
||||
output_bytes_retained,
|
||||
output_bytes_omitted,
|
||||
} => Some(EventBody::RunSessionToolCallCompleted(
|
||||
RunSessionToolCallCompletedProps {
|
||||
turn_id: event_turn_id,
|
||||
|
|
@ -1278,6 +1281,13 @@ fn agent_event_payload(event_turn_id: TurnId, event: AgentEvent) -> Option<Event
|
|||
tool_call_id,
|
||||
output,
|
||||
is_error,
|
||||
output_bytes_observed: Some(
|
||||
u64::try_from(output_bytes_observed).unwrap_or(u64::MAX),
|
||||
),
|
||||
output_bytes_retained: Some(
|
||||
u64::try_from(output_bytes_retained).unwrap_or(u64::MAX),
|
||||
),
|
||||
output_bytes_omitted: Some(u64::try_from(output_bytes_omitted).unwrap_or(u64::MAX)),
|
||||
},
|
||||
)),
|
||||
_ => None,
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
emitter: Emitter,
|
||||
session_id: String,
|
||||
tool_call_id: Option<String>,
|
||||
tool_output_stats: Arc<Mutex<Option<OutputCaptureStats>>>,
|
||||
}
|
||||
|
||||
impl SessionBoundEmitter {
|
||||
#[must_use]
|
||||
pub fn new(emitter: Emitter, session_id: String, tool_call_id: Option<String>) -> Self {
|
||||
Self {
|
||||
emitter,
|
||||
session_id,
|
||||
tool_call_id,
|
||||
tool_output_stats: Arc::new(Mutex::new(None)),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn take_tool_output_stats(&self) -> Option<OutputCaptureStats> {
|
||||
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)]
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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) }
|
||||
})
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<OutputCaptureStats>,
|
||||
) -> 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<OutputCaptureStats>,
|
||||
}
|
||||
|
||||
/// 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<dyn ToolEnvProvider>>,
|
||||
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<Arc<dyn AgentEventEmitter>> =
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<
|
||||
|
|
|
|||
|
|
@ -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<String, String> {
|
||||
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"
|
||||
|
|
|
|||
|
|
@ -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<T: Serialize + ?Sized>(value: &T) -> usize {
|
||||
struct CountingWriter(usize);
|
||||
impl std::io::Write for CountingWriter {
|
||||
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
|
||||
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<usize>,
|
||||
) -> 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"));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<i32>,
|
||||
termination: CommandTermination,
|
||||
duration_ms: u64,
|
||||
streams_separated: bool,
|
||||
exit_code: Option<i32>,
|
||||
termination: CommandTermination,
|
||||
duration_ms: u64,
|
||||
streams_separated: bool,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
exec_output_tail: Option<ExecOutputTail>,
|
||||
exec_output_tail: Option<ExecOutputTail>,
|
||||
#[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 {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<usize>) -> 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,
|
||||
|
|
|
|||
|
|
@ -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<Mutex<Vec<u8>>>,
|
||||
seen: &Arc<Mutex<OutputCaptureBuffer>>,
|
||||
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<Mutex<OutputCaptureBuffer>>,
|
||||
) -> (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();
|
||||
|
|
|
|||
|
|
@ -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<Vec<String>>,
|
||||
stdin: Option<Vec<u8>>,
|
||||
output_callback: Option<CommandOutputCallback>,
|
||||
) -> crate::Result<(Vec<u8>, Vec<u8>, i32)> {
|
||||
stream_output_bytes_cap: Option<usize>,
|
||||
) -> 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,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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<R>(
|
|||
mut reader: Option<R>,
|
||||
stream: CommandOutputStream,
|
||||
output_callback: Option<CommandOutputCallback>,
|
||||
) -> crate::Result<Vec<u8>>
|
||||
stream_output_bytes_cap: Option<usize>,
|
||||
) -> crate::Result<OutputCaptureBuffer>
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<usize>,
|
||||
head: Vec<u8>,
|
||||
tail: VecDeque<u8>,
|
||||
observed_bytes: usize,
|
||||
}
|
||||
|
||||
impl OutputCaptureBuffer {
|
||||
#[must_use]
|
||||
pub(crate) fn new(max_bytes: Option<usize>) -> 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<u8> {
|
||||
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<u8>, 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<u64>,
|
||||
pub working_dir: Option<&'a str>,
|
||||
pub env_vars: Option<&'a HashMap<String, String>>,
|
||||
pub cancel_token: Option<CancellationToken>,
|
||||
pub stdin: Option<Vec<u8>>,
|
||||
pub output_callback: Option<CommandOutputCallback>,
|
||||
pub command: &'a str,
|
||||
pub timeout_ms: Option<u64>,
|
||||
pub working_dir: Option<&'a str>,
|
||||
pub env_vars: Option<&'a HashMap<String, String>>,
|
||||
pub cancel_token: Option<CancellationToken>,
|
||||
pub stdin: Option<Vec<u8>>,
|
||||
pub output_callback: Option<CommandOutputCallback>,
|
||||
/// Maximum bytes retained from each stream. Providers continue draining
|
||||
/// stdout and stderr after the cap is reached.
|
||||
pub stream_output_bytes_cap: Option<usize>,
|
||||
}
|
||||
|
||||
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<usize>,
|
||||
) -> crate::Result<ExecStreamingResult> {
|
||||
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<usize>) -> 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<Box<dyn AsyncWrite + Send>>,
|
||||
pub stdout: Pin<Box<dyn AsyncRead + Send>>,
|
||||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}),
|
||||
),
|
||||
];
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
);
|
||||
|
|
|
|||
|
|
@ -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<u64>,
|
||||
/// Tool-output bytes kept in `output`, excluding truncation notices.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_retained: Option<u64>,
|
||||
/// Tool-output bytes discarded before this event was emitted.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_omitted: Option<u64>,
|
||||
/// 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<ToolResult>,
|
||||
pub tool_result: Option<ToolResult>,
|
||||
/// Turn that owned this tool call.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub turn_id: Option<TurnId>,
|
||||
pub turn_id: Option<TurnId>,
|
||||
}
|
||||
|
||||
/// 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<i32>,
|
||||
pub termination: CommandTermination,
|
||||
pub duration_ms: u64,
|
||||
pub exit_code: Option<i32>,
|
||||
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<ExecOutputTail>,
|
||||
pub visit: u32,
|
||||
pub exec_output_tail: Option<ExecOutputTail>,
|
||||
/// Raw stdout and stderr bytes drained from the process.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_observed: Option<u64>,
|
||||
/// Raw process-output bytes kept by the streaming capture buffers.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_retained: Option<u64>,
|
||||
/// Raw process-output bytes discarded by the streaming capture buffers.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_omitted: Option<u64>,
|
||||
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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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<u64>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_retained: Option<u64>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub output_bytes_omitted: Option<u64>,
|
||||
}
|
||||
|
||||
#[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());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue