mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
Bound agent tool output capture
This commit is contained in:
parent
b3f602f6e9
commit
ddcdafa06b
12 changed files with 548 additions and 65 deletions
|
|
@ -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,
|
||||
};
|
||||
|
|
|
|||
|
|
@ -33,6 +33,7 @@ use crate::tools::{
|
|||
DEFAULT_READ_LINES, emit_shell_process_completed, execute_grep, execute_shell_command,
|
||||
grep_result_path, make_edit_file_tool, optional_usize_arg, required_str,
|
||||
};
|
||||
use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output};
|
||||
|
||||
const DEFAULT_GREP_RESULTS: usize = 250;
|
||||
const MAX_GREP_RESULTS: usize = 2000;
|
||||
|
|
@ -143,6 +144,12 @@ explicitly asked. Never run commands requiring superuser privileges unless expli
|
|||
let _ = write!(out, "Command failed with exit code: {code}");
|
||||
}
|
||||
let is_success = result.is_success();
|
||||
let out = retain_tool_output(
|
||||
&out,
|
||||
MAX_RETAINED_TOOL_OUTPUT_BYTES,
|
||||
streaming.output_capture().omitted_bytes,
|
||||
)
|
||||
.output;
|
||||
emit_shell_process_completed(&ctx, streaming).await;
|
||||
if is_success { Ok(out) } else { Err(out) }
|
||||
})
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ use crate::question_tools::{self, AgentToolRuntime, is_question_tool};
|
|||
use crate::sandbox::Sandbox;
|
||||
use crate::session::ToolEnvProvider;
|
||||
use crate::tool_registry::{AgentEventEmitter, RegisteredTool, ToolContext, ToolRegistry};
|
||||
use crate::truncation::truncate_tool_output;
|
||||
use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, retain_tool_output, truncate_tool_output};
|
||||
use crate::types::AgentEvent;
|
||||
|
||||
/// Execute tool calls, choosing parallel or sequential based on `parallel`
|
||||
|
|
@ -265,7 +265,7 @@ fn error_tool_result_with_events(
|
|||
message: &str,
|
||||
) -> ToolResult {
|
||||
emit_tool_call_started(emitter, session_id, tc);
|
||||
let result = ToolResult::error(&tc.id, message);
|
||||
let result = retain_tool_result(&ToolResult::error(&tc.id, message));
|
||||
emit_tool_call_result(emitter, session_id, tc, &result);
|
||||
truncate_tool_result(&result, &tc.name, config)
|
||||
}
|
||||
|
|
@ -386,7 +386,7 @@ async fn execute_and_emit_one_tool_with_lookup(
|
|||
emit_tool_call_started(emitter, session_id, tc);
|
||||
|
||||
if let Some(reason) = access_denial {
|
||||
let result = ToolResult::error(&tc.id, &reason);
|
||||
let result = retain_tool_result(&ToolResult::error(&tc.id, &reason));
|
||||
emit_tool_call_result(emitter, session_id, tc, &result);
|
||||
return truncate_tool_result(&result, &tc.name, config);
|
||||
}
|
||||
|
|
@ -400,24 +400,26 @@ async fn execute_and_emit_one_tool_with_lookup(
|
|||
debug!(tool = %tc.name, hook_event = "pre_tool_use", ?decision, duration_ms = elapsed, "Tool hook complete");
|
||||
|
||||
if let ToolHookDecision::Block { reason } = decision {
|
||||
let result = ToolResult::error(&tc.id, &reason);
|
||||
let result = retain_tool_result(&ToolResult::error(&tc.id, &reason));
|
||||
emit_tool_call_result(emitter, session_id, tc, &result);
|
||||
return truncate_tool_result(&result, &tc.name, config);
|
||||
}
|
||||
}
|
||||
|
||||
let result = execute_one_tool(
|
||||
tc,
|
||||
registered_tool,
|
||||
env,
|
||||
cancel_token,
|
||||
emitter,
|
||||
session_id,
|
||||
root_session_id,
|
||||
tool_env_provider,
|
||||
agent_tool_runtime,
|
||||
)
|
||||
.await;
|
||||
let result = retain_tool_result(
|
||||
&execute_one_tool(
|
||||
tc,
|
||||
registered_tool,
|
||||
env,
|
||||
cancel_token,
|
||||
emitter,
|
||||
session_id,
|
||||
root_session_id,
|
||||
tool_env_provider,
|
||||
agent_tool_runtime,
|
||||
)
|
||||
.await,
|
||||
);
|
||||
|
||||
emit_tool_call_result(emitter, session_id, tc, &result);
|
||||
|
||||
|
|
@ -446,6 +448,24 @@ async fn execute_and_emit_one_tool_with_lookup(
|
|||
truncate_tool_result(&result, &tc.name, config)
|
||||
}
|
||||
|
||||
/// Bound model-native tool output before it reaches hooks, events, or history.
|
||||
fn retain_tool_result(result: &ToolResult) -> ToolResult {
|
||||
let retained_content = match &result.content {
|
||||
serde_json::Value::String(output) => serde_json::Value::String(
|
||||
retain_tool_output(output, MAX_RETAINED_TOOL_OUTPUT_BYTES, 0).output,
|
||||
),
|
||||
other => other.clone(),
|
||||
};
|
||||
|
||||
ToolResult {
|
||||
tool_call_id: result.tool_call_id.clone(),
|
||||
content: retained_content,
|
||||
is_error: result.is_error,
|
||||
image_data: result.image_data.clone(),
|
||||
image_media_type: result.image_media_type.clone(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Execute a single tool call: argument validation and execution.
|
||||
#[allow(
|
||||
clippy::too_many_arguments,
|
||||
|
|
@ -877,6 +897,47 @@ mod tests {
|
|||
assert!(content.contains("echo: hello"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn tool_output_is_bounded_before_events_and_history() {
|
||||
let mut registry = ToolRegistry::new();
|
||||
registry.register(make_echo_tool());
|
||||
let text = "x".repeat(MAX_RETAINED_TOOL_OUTPUT_BYTES + 100);
|
||||
let tc = make_tool_call("echo", "call_large", serde_json::json!({"text": text}));
|
||||
let emitter = Emitter::new();
|
||||
let mut receiver = emitter.subscribe();
|
||||
|
||||
let result = execute_and_emit_one_tool(
|
||||
&tc,
|
||||
®istry,
|
||||
make_sandbox(),
|
||||
None,
|
||||
CancellationToken::new(),
|
||||
&SessionOptions::default(),
|
||||
&emitter,
|
||||
"test-session",
|
||||
"test-session",
|
||||
None,
|
||||
)
|
||||
.await;
|
||||
|
||||
let result_output = result.content.as_str().expect("string tool output");
|
||||
assert_eq!(result_output.len(), MAX_RETAINED_TOOL_OUTPUT_BYTES);
|
||||
|
||||
let completed_output = loop {
|
||||
let event = receiver.try_recv().expect("tool completion event");
|
||||
if let AgentEvent::ToolCallCompleted { output, .. } = event.event {
|
||||
break output;
|
||||
}
|
||||
};
|
||||
assert_eq!(
|
||||
completed_output
|
||||
.as_str()
|
||||
.expect("string event output")
|
||||
.len(),
|
||||
MAX_RETAINED_TOOL_OUTPUT_BYTES
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn post_tool_use_hook_fires_on_success() {
|
||||
let mut registry = ToolRegistry::new();
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ use tokio::task;
|
|||
use crate::config::NativeToolOptions;
|
||||
use crate::sandbox::{ExecStreamingResult, GrepOptions};
|
||||
use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource};
|
||||
use crate::truncation::{MAX_RETAINED_TOOL_OUTPUT_BYTES, RetainedToolOutput, retain_tool_output};
|
||||
use crate::types::AgentEvent;
|
||||
use crate::web_search::{SearchBackend, make_web_search_tool};
|
||||
|
||||
|
|
@ -303,6 +304,7 @@ pub(crate) async fn execute_shell_command(
|
|||
working_dir: cwd,
|
||||
env_vars: tool_env.as_ref(),
|
||||
cancel_token: Some(ctx.cancel.clone()),
|
||||
stream_output_bytes_cap: Some(MAX_RETAINED_TOOL_OUTPUT_BYTES),
|
||||
..crate::ExecStreamingRequest::new(command)
|
||||
})
|
||||
.await
|
||||
|
|
@ -318,7 +320,7 @@ pub(crate) async fn run_shell_command(
|
|||
cwd: Option<&str>,
|
||||
) -> Result<String, String> {
|
||||
let streaming = execute_shell_command(ctx, command, timeout_ms, cwd).await?;
|
||||
let text = render_shell_result(&streaming);
|
||||
let text = render_shell_result(&streaming).output;
|
||||
let is_success = streaming.result.is_success();
|
||||
emit_shell_process_completed(ctx, streaming).await;
|
||||
|
||||
|
|
@ -364,7 +366,7 @@ pub(crate) async fn emit_shell_process_completed(
|
|||
/// Renders the model-facing shell result: termination, exit code, duration,
|
||||
/// and provider-honest output sections. Metadata stays at the head and
|
||||
/// `stderr` at the tail so head/tail truncation preserves both.
|
||||
fn render_shell_result(streaming: &ExecStreamingResult) -> String {
|
||||
fn render_shell_result(streaming: &ExecStreamingResult) -> RetainedToolOutput {
|
||||
let result = &streaming.result;
|
||||
let mut output = format!(
|
||||
"Termination: {}\nExit code: {}\nDuration: {}ms\n",
|
||||
|
|
@ -384,7 +386,11 @@ fn render_shell_result(streaming: &ExecStreamingResult) -> String {
|
|||
} else if !result.stdout.is_empty() {
|
||||
let _ = write!(output, "output (combined):\n{}\n", result.stdout);
|
||||
}
|
||||
output
|
||||
retain_tool_output(
|
||||
&output,
|
||||
MAX_RETAINED_TOOL_OUTPUT_BYTES,
|
||||
streaming.output_capture().omitted_bytes,
|
||||
)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
|
|
|
|||
|
|
@ -1,6 +1,64 @@
|
|||
use crate::config::SessionOptions;
|
||||
use crate::sandbox::OutputCaptureStats;
|
||||
use crate::tool_permissions::canonical_tool_name;
|
||||
|
||||
pub(crate) const MAX_RETAINED_TOOL_OUTPUT_BYTES: usize = 1024 * 1024;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub(crate) struct RetainedToolOutput {
|
||||
pub output: String,
|
||||
pub stats: OutputCaptureStats,
|
||||
}
|
||||
|
||||
/// Keep an equal-sized UTF-8 prefix and suffix within a byte budget.
|
||||
///
|
||||
/// `previously_omitted_bytes` accounts for output a streaming provider
|
||||
/// discarded before the rendered result was assembled.
|
||||
#[must_use]
|
||||
pub(crate) fn retain_tool_output(
|
||||
output: &str,
|
||||
max_bytes: usize,
|
||||
previously_omitted_bytes: usize,
|
||||
) -> RetainedToolOutput {
|
||||
let observed_bytes = output.len().saturating_add(previously_omitted_bytes);
|
||||
if output.len() <= max_bytes {
|
||||
return RetainedToolOutput {
|
||||
output: output.to_string(),
|
||||
stats: OutputCaptureStats {
|
||||
observed_bytes,
|
||||
retained_bytes: output.len(),
|
||||
omitted_bytes: previously_omitted_bytes,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
let head_budget = max_bytes / 2;
|
||||
let tail_budget = max_bytes.saturating_sub(head_budget);
|
||||
let head_end = output.floor_char_boundary(head_budget);
|
||||
let tail_start = ceil_char_boundary(output, output.len().saturating_sub(tail_budget));
|
||||
let retained_bytes = head_end.saturating_add(output.len().saturating_sub(tail_start));
|
||||
let mut retained = String::with_capacity(retained_bytes);
|
||||
retained.push_str(&output[..head_end]);
|
||||
retained.push_str(&output[tail_start..]);
|
||||
|
||||
RetainedToolOutput {
|
||||
output: retained,
|
||||
stats: OutputCaptureStats {
|
||||
observed_bytes,
|
||||
retained_bytes,
|
||||
omitted_bytes: observed_bytes.saturating_sub(retained_bytes),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
fn ceil_char_boundary(output: &str, index: usize) -> usize {
|
||||
let mut index = index.min(output.len());
|
||||
while index < output.len() && !output.is_char_boundary(index) {
|
||||
index += 1;
|
||||
}
|
||||
index
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum TruncationMode {
|
||||
HeadTail,
|
||||
|
|
@ -121,6 +179,30 @@ pub fn truncate_tool_output(output: &str, tool_name: &str, config: &SessionOptio
|
|||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn retained_tool_output_keeps_equal_head_and_tail() {
|
||||
let retained = retain_tool_output("abcdefghijkl", 8, 0);
|
||||
|
||||
assert_eq!(retained.output, "abcdijkl");
|
||||
assert_eq!(retained.stats.observed_bytes, 12);
|
||||
assert_eq!(retained.stats.retained_bytes, 8);
|
||||
assert_eq!(retained.stats.omitted_bytes, 4);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn retained_tool_output_stays_within_budget_at_utf8_boundaries() {
|
||||
let retained = retain_tool_output("aa😀😀zz", 7, 3);
|
||||
|
||||
assert!(retained.output.len() <= 7, "{}", retained.output.len());
|
||||
assert!(retained.output.starts_with("aa"));
|
||||
assert!(retained.output.ends_with("zz"));
|
||||
assert_eq!(retained.stats.observed_bytes, "aa😀😀zz".len() + 3);
|
||||
assert_eq!(
|
||||
retained.stats.omitted_bytes,
|
||||
retained.stats.observed_bytes - retained.output.len()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn under_limit_passthrough_chars() {
|
||||
let output = "short output";
|
||||
|
|
|
|||
|
|
@ -32,8 +32,9 @@ use crate::git_retry::{self, CredentialContext, GitRetryReason};
|
|||
use crate::push_credentials::{self, PushCredentialState};
|
||||
use crate::redact::redact_auth_url;
|
||||
use crate::sandbox::{
|
||||
self, BASH_ENV_VAR, BASH_PROBE_MARKER, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH,
|
||||
REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, optional_timeout, resolve_path, validate_bash_probe,
|
||||
self, BASH_ENV_VAR, BASH_PROBE_MARKER, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS,
|
||||
OutputCaptureBuffer, REMOTE_BASH, REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, optional_timeout,
|
||||
resolve_path, validate_bash_probe,
|
||||
};
|
||||
use crate::{
|
||||
CommandOutputCallback, DirEntry, ExecResult, ExecStreamingRequest, ExecStreamingResult,
|
||||
|
|
@ -2212,6 +2213,7 @@ impl Sandbox for DaytonaSandbox {
|
|||
cancel_token,
|
||||
stdin,
|
||||
output_callback,
|
||||
stream_output_bytes_cap,
|
||||
} = request;
|
||||
let sandbox = self.sandbox()?;
|
||||
let start = Instant::now();
|
||||
|
|
@ -2252,8 +2254,12 @@ impl Sandbox for DaytonaSandbox {
|
|||
return Err(crate::Error::context("Failed to get process service", err));
|
||||
}
|
||||
};
|
||||
let stdout_seen = Arc::new(Mutex::new(Vec::new()));
|
||||
let stderr_seen = Arc::new(Mutex::new(Vec::new()));
|
||||
let stdout_seen = Arc::new(Mutex::new(OutputCaptureBuffer::new(
|
||||
stream_output_bytes_cap,
|
||||
)));
|
||||
let stderr_seen = Arc::new(Mutex::new(OutputCaptureBuffer::new(
|
||||
stream_output_bytes_cap,
|
||||
)));
|
||||
let saw_live_chunk = Arc::new(AtomicBool::new(false));
|
||||
|
||||
let stream_session_id = session.id().to_string();
|
||||
|
|
@ -2277,7 +2283,7 @@ impl Sandbox for DaytonaSandbox {
|
|||
let bytes = chunk.into_bytes();
|
||||
if !bytes.is_empty() {
|
||||
saw_live_chunk.store(true, Ordering::Relaxed);
|
||||
stdout_seen.lock().await.extend_from_slice(&bytes);
|
||||
stdout_seen.lock().await.push(&bytes);
|
||||
if let Some(callback) = callback {
|
||||
callback(CommandOutputStream::Stdout, bytes)
|
||||
.await
|
||||
|
|
@ -2295,7 +2301,7 @@ impl Sandbox for DaytonaSandbox {
|
|||
let bytes = chunk.into_bytes();
|
||||
if !bytes.is_empty() {
|
||||
saw_live_chunk.store(true, Ordering::Relaxed);
|
||||
stderr_seen.lock().await.extend_from_slice(&bytes);
|
||||
stderr_seen.lock().await.push(&bytes);
|
||||
if let Some(callback) = callback {
|
||||
callback(CommandOutputStream::Stderr, bytes)
|
||||
.await
|
||||
|
|
@ -2369,8 +2375,20 @@ impl Sandbox for DaytonaSandbox {
|
|||
.await?;
|
||||
}
|
||||
|
||||
let stdout = String::from_utf8_lossy(&stdout_seen.lock().await).into_owned();
|
||||
let stderr = String::from_utf8_lossy(&stderr_seen.lock().await).into_owned();
|
||||
let (stdout, stdout_capture) = {
|
||||
let stdout_seen = stdout_seen.lock().await;
|
||||
(
|
||||
String::from_utf8_lossy(&stdout_seen.to_bytes()).into_owned(),
|
||||
stdout_seen.stats(),
|
||||
)
|
||||
};
|
||||
let (stderr, stderr_capture) = {
|
||||
let stderr_seen = stderr_seen.lock().await;
|
||||
(
|
||||
String::from_utf8_lossy(&stderr_seen.to_bytes()).into_owned(),
|
||||
stderr_seen.stats(),
|
||||
)
|
||||
};
|
||||
|
||||
let result = ExecStreamingResult {
|
||||
result: ExecResult {
|
||||
|
|
@ -2384,6 +2402,8 @@ impl Sandbox for DaytonaSandbox {
|
|||
},
|
||||
streams_separated,
|
||||
live_streaming: saw_live_chunk.load(Ordering::Relaxed),
|
||||
stdout_capture,
|
||||
stderr_capture,
|
||||
};
|
||||
if let Some(stdin_file) = stdin_file.as_mut() {
|
||||
stdin_file.close().await;
|
||||
|
|
@ -2886,7 +2906,7 @@ async fn fetch_daytona_session_logs(
|
|||
async fn append_missing_log_suffix(
|
||||
stream: CommandOutputStream,
|
||||
final_bytes: &[u8],
|
||||
seen: &Arc<Mutex<Vec<u8>>>,
|
||||
seen: &Arc<Mutex<OutputCaptureBuffer>>,
|
||||
output_callback: Option<&CommandOutputCallback>,
|
||||
) -> crate::Result<()> {
|
||||
if final_bytes.is_empty() {
|
||||
|
|
@ -2894,13 +2914,13 @@ async fn append_missing_log_suffix(
|
|||
}
|
||||
|
||||
let mut seen = seen.lock().await;
|
||||
let offset = missing_log_suffix_offset(&seen, final_bytes);
|
||||
let offset = captured_log_suffix_offset(&seen, final_bytes);
|
||||
if offset >= final_bytes.len() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let missing = final_bytes[offset..].to_vec();
|
||||
seen.extend_from_slice(&missing);
|
||||
seen.push(&missing);
|
||||
drop(seen);
|
||||
match output_callback {
|
||||
Some(output_callback) => output_callback(stream, missing).await,
|
||||
|
|
@ -2908,6 +2928,42 @@ async fn append_missing_log_suffix(
|
|||
}
|
||||
}
|
||||
|
||||
fn captured_log_suffix_offset(seen: &OutputCaptureBuffer, final_bytes: &[u8]) -> usize {
|
||||
let stats = seen.stats();
|
||||
if stats.omitted_bytes == 0 {
|
||||
return missing_log_suffix_offset(&seen.to_bytes(), final_bytes);
|
||||
}
|
||||
|
||||
let observed_bytes = seen.observed_bytes();
|
||||
let head = seen.retained_head();
|
||||
let tail = seen.retained_tail();
|
||||
if final_bytes.len() >= observed_bytes
|
||||
&& final_bytes.starts_with(head)
|
||||
&& tail.iter().copied().eq(final_bytes
|
||||
[observed_bytes.saturating_sub(tail.len())..observed_bytes]
|
||||
.iter()
|
||||
.copied())
|
||||
{
|
||||
return observed_bytes;
|
||||
}
|
||||
if final_bytes.len() <= observed_bytes && final_bytes.starts_with(head) {
|
||||
return final_bytes.len();
|
||||
}
|
||||
|
||||
let max_overlap = tail.len().min(final_bytes.len());
|
||||
for overlap in (1..=max_overlap).rev() {
|
||||
if tail
|
||||
.iter()
|
||||
.skip(tail.len() - overlap)
|
||||
.copied()
|
||||
.eq(final_bytes[..overlap].iter().copied())
|
||||
{
|
||||
return overlap;
|
||||
}
|
||||
}
|
||||
0
|
||||
}
|
||||
|
||||
fn missing_log_suffix_offset(seen: &[u8], final_bytes: &[u8]) -> usize {
|
||||
if final_bytes.starts_with(seen) {
|
||||
return seen.len();
|
||||
|
|
@ -4730,6 +4786,16 @@ mod tests {
|
|||
assert_eq!(missing_log_suffix_offset(b"abc", b"def"), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn captured_log_suffix_offset_uses_observed_length_after_truncation() {
|
||||
let mut seen = OutputCaptureBuffer::new(Some(6));
|
||||
seen.push(b"abcdefgh");
|
||||
|
||||
assert_eq!(captured_log_suffix_offset(&seen, b"abcdefghij"), 8);
|
||||
assert_eq!(captured_log_suffix_offset(&seen, b"abcdefgh"), 8);
|
||||
assert_eq!(captured_log_suffix_offset(&seen, b"abcd"), 4);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn detect_git_remote_from_repo() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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,143 @@ pub struct ExecStreamingResult {
|
|||
pub result: ExecResult,
|
||||
pub streams_separated: bool,
|
||||
pub live_streaming: bool,
|
||||
pub stdout_capture: OutputCaptureStats,
|
||||
pub stderr_capture: OutputCaptureStats,
|
||||
}
|
||||
|
||||
impl ExecStreamingResult {
|
||||
#[must_use]
|
||||
pub fn output_capture(&self) -> OutputCaptureStats {
|
||||
self.stdout_capture.combine(self.stderr_capture)
|
||||
}
|
||||
}
|
||||
|
||||
/// Byte counts for output observed and retained while draining a process.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct OutputCaptureStats {
|
||||
pub observed_bytes: usize,
|
||||
pub retained_bytes: usize,
|
||||
pub omitted_bytes: usize,
|
||||
}
|
||||
|
||||
impl OutputCaptureStats {
|
||||
#[must_use]
|
||||
pub fn complete(byte_count: usize) -> Self {
|
||||
Self {
|
||||
observed_bytes: byte_count,
|
||||
retained_bytes: byte_count,
|
||||
omitted_bytes: 0,
|
||||
}
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn combine(self, other: Self) -> Self {
|
||||
Self {
|
||||
observed_bytes: self.observed_bytes.saturating_add(other.observed_bytes),
|
||||
retained_bytes: self.retained_bytes.saturating_add(other.retained_bytes),
|
||||
omitted_bytes: self.omitted_bytes.saturating_add(other.omitted_bytes),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A byte buffer that keeps an equal-sized stable prefix and rolling suffix.
|
||||
///
|
||||
/// The process is always drained. Once the optional cap is full, bytes from
|
||||
/// the middle are discarded while the newest suffix replaces the old tail.
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct OutputCaptureBuffer {
|
||||
max_bytes: Option<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 stats = self.stats();
|
||||
let mut bytes = Vec::with_capacity(stats.retained_bytes);
|
||||
bytes.extend_from_slice(&self.head);
|
||||
bytes.extend(self.tail.iter().copied());
|
||||
bytes
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub(crate) fn into_parts(self) -> (Vec<u8>, OutputCaptureStats) {
|
||||
let stats = self.stats();
|
||||
let mut bytes = Vec::with_capacity(stats.retained_bytes);
|
||||
bytes.extend(self.head);
|
||||
bytes.extend(self.tail);
|
||||
(bytes, stats)
|
||||
}
|
||||
|
||||
#[cfg(feature = "daytona")]
|
||||
#[must_use]
|
||||
pub(crate) fn observed_bytes(&self) -> usize {
|
||||
self.observed_bytes
|
||||
}
|
||||
|
||||
#[cfg(feature = "daytona")]
|
||||
#[must_use]
|
||||
pub(crate) fn retained_head(&self) -> &[u8] {
|
||||
&self.head
|
||||
}
|
||||
|
||||
#[cfg(feature = "daytona")]
|
||||
#[must_use]
|
||||
pub(crate) fn retained_tail(&self) -> &VecDeque<u8> {
|
||||
&self.tail
|
||||
}
|
||||
}
|
||||
|
||||
pub type CommandOutputCallback = Arc<
|
||||
|
|
@ -785,13 +922,16 @@ pub type CommandOutputCallback = Arc<
|
|||
/// type does not implement `Debug` because standard input can contain
|
||||
/// sensitive workflow data.
|
||||
pub struct ExecStreamingRequest<'a> {
|
||||
pub command: &'a str,
|
||||
pub timeout_ms: Option<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 +945,7 @@ impl<'a> ExecStreamingRequest<'a> {
|
|||
cancel_token: None,
|
||||
stdin: None,
|
||||
output_callback: None,
|
||||
stream_output_bytes_cap: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -846,9 +987,10 @@ where
|
|||
}
|
||||
|
||||
pub(crate) async fn replay_exec_result(
|
||||
result: ExecResult,
|
||||
mut result: ExecResult,
|
||||
streams_separated: bool,
|
||||
output_callback: Option<&CommandOutputCallback>,
|
||||
stream_output_bytes_cap: Option<usize>,
|
||||
) -> crate::Result<ExecStreamingResult> {
|
||||
if let Some(output_callback) = output_callback {
|
||||
if !result.stdout.is_empty() {
|
||||
|
|
@ -866,10 +1008,21 @@ pub(crate) async fn replay_exec_result(
|
|||
.await?;
|
||||
}
|
||||
}
|
||||
let mut stdout_capture = OutputCaptureBuffer::new(stream_output_bytes_cap);
|
||||
stdout_capture.push(result.stdout.as_bytes());
|
||||
let mut stderr_capture = OutputCaptureBuffer::new(stream_output_bytes_cap);
|
||||
stderr_capture.push(result.stderr.as_bytes());
|
||||
let (stdout, stdout_capture) = stdout_capture.into_parts();
|
||||
let (stderr, stderr_capture) = stderr_capture.into_parts();
|
||||
result.stdout = String::from_utf8_lossy(&stdout).into_owned();
|
||||
result.stderr = String::from_utf8_lossy(&stderr).into_owned();
|
||||
|
||||
Ok(ExecStreamingResult {
|
||||
result,
|
||||
streams_separated,
|
||||
live_streaming: false,
|
||||
stdout_capture,
|
||||
stderr_capture,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -1179,7 +1332,13 @@ pub trait Sandbox: Send + Sync {
|
|||
request.cancel_token,
|
||||
)
|
||||
.await?;
|
||||
replay_exec_result(result, true, request.output_callback.as_ref()).await
|
||||
replay_exec_result(
|
||||
result,
|
||||
true,
|
||||
request.output_callback.as_ref(),
|
||||
request.stream_output_bytes_cap,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Launch a long-lived process with bidirectional stdio attached.
|
||||
|
|
@ -2517,6 +2676,43 @@ mod push_tests {
|
|||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn output_capture_buffer_keeps_stable_head_and_rolling_tail() {
|
||||
let mut buffer = OutputCaptureBuffer::new(Some(8));
|
||||
buffer.push(b"abc");
|
||||
buffer.push(b"defghi");
|
||||
buffer.push(b"jkl");
|
||||
|
||||
let (bytes, stats) = buffer.into_parts();
|
||||
assert_eq!(bytes, b"abcdijkl");
|
||||
assert_eq!(stats.observed_bytes, 12);
|
||||
assert_eq!(stats.retained_bytes, 8);
|
||||
assert_eq!(stats.omitted_bytes, 4);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn output_capture_buffer_without_cap_retains_everything() {
|
||||
let mut buffer = OutputCaptureBuffer::new(None);
|
||||
buffer.push(b"abc");
|
||||
buffer.push(b"def");
|
||||
|
||||
let (bytes, stats) = buffer.into_parts();
|
||||
assert_eq!(bytes, b"abcdef");
|
||||
assert_eq!(stats, OutputCaptureStats::complete(6));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zero_byte_output_capture_buffer_still_counts_drained_bytes() {
|
||||
let mut buffer = OutputCaptureBuffer::new(Some(0));
|
||||
buffer.push(b"abcdef");
|
||||
|
||||
let (bytes, stats) = buffer.into_parts();
|
||||
assert!(bytes.is_empty());
|
||||
assert_eq!(stats.observed_bytes, 6);
|
||||
assert_eq!(stats.retained_bytes, 0);
|
||||
assert_eq!(stats.omitted_bytes, 6);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn exec_result_fields() {
|
||||
let result = ExecResult {
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue