From 4cbf76d4b0bbee27465ab8430a772b2098051828 Mon Sep 17 00:00:00 2001 From: Fabro Date: Mon, 4 May 2026 19:31:25 +0000 Subject: [PATCH] fabro(01KQT1V2W1R6ZH72CFT2QDJ39Q): simplify_opus (succeeded) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fabro-Run: 01KQT1V2W1R6ZH72CFT2QDJ39Q Fabro-Completed: 6 Fabro-Checkpoint: d0025d74b4fb913a3e39642f3b3f0760bf1c52fc ⚒️ Generated with [Fabro](https://fabro.sh) --- lib/crates/fabro-sandbox/src/daytona/mod.rs | 46 ++++++--- lib/crates/fabro-sandbox/src/docker.rs | 12 +-- lib/crates/fabro-sandbox/src/local.rs | 10 +- lib/crates/fabro-sandbox/src/sandbox.rs | 10 ++ .../fabro-workflow/src/handler/llm/cli.rs | 99 ++++++++++--------- .../fabro-workflow/src/handler/prompt.rs | 10 +- 6 files changed, 107 insertions(+), 80 deletions(-) diff --git a/lib/crates/fabro-sandbox/src/daytona/mod.rs b/lib/crates/fabro-sandbox/src/daytona/mod.rs index 51a45e92d..8230bcaee 100644 --- a/lib/crates/fabro-sandbox/src/daytona/mod.rs +++ b/lib/crates/fabro-sandbox/src/daytona/mod.rs @@ -26,7 +26,7 @@ use tokio_util::sync::CancellationToken; use crate::clone_source::{self, CloneDecision, EmptyWorkspaceReason}; use crate::redact::redact_auth_url; -use crate::sandbox::resolve_path; +use crate::sandbox::{optional_timeout, resolve_path}; use crate::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingResult, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback, format_lines_numbered, shell_quote, @@ -37,6 +37,9 @@ const DEFAULT_SNAPSHOT: &str = "daytona-medium"; pub const DEFAULT_DAYTONA_API_URL: &str = "https://app.daytona.io/api"; const FABRO_SANDBOX_USER_AGENT: &str = concat!("fabro-sandbox/", env!("CARGO_PKG_VERSION")); const DAYTONA_PROBE_TIMEOUT: Duration = Duration::from_secs(20); +/// Upper bound on `DaytonaSession::close` so a stalled Daytona REST call cannot +/// block cancellation/timeout paths from returning. +const DAYTONA_SESSION_CLOSE_TIMEOUT: Duration = Duration::from_secs(10); /// Permissions a Daytona API key needs for Fabro's snapshot and sandbox flow. pub const REQUIRED_DAYTONA_PERMISSIONS: &[Permissions] = &[ @@ -1673,19 +1676,39 @@ impl DaytonaSession { } /// Idempotent: a second call after `active=false` is a no-op. + /// + /// `delete_session` is bounded by [`DAYTONA_SESSION_CLOSE_TIMEOUT`] so a + /// stalled Daytona REST call cannot block cancellation paths indefinitely. async fn close(&mut self, reason: &'static str) { if !self.active { return; } self.active = false; if let Some(svc) = self.process_svc.take() { - if let Err(err) = svc.delete_session(&self.session_id).await { - tracing::warn!( - error = %err, - session_id = %self.session_id, - reason, - "failed to delete Daytona session" - ); + match time::timeout( + DAYTONA_SESSION_CLOSE_TIMEOUT, + svc.delete_session(&self.session_id), + ) + .await + { + Ok(Ok(())) => {} + Ok(Err(err)) => { + tracing::warn!( + error = %err, + session_id = %self.session_id, + reason, + "failed to delete Daytona session" + ); + } + Err(_) => { + tracing::warn!( + session_id = %self.session_id, + reason, + timeout_ms = u64::try_from(DAYTONA_SESSION_CLOSE_TIMEOUT.as_millis()) + .unwrap_or(u64::MAX), + "timed out deleting Daytona session" + ); + } } } } @@ -1742,12 +1765,7 @@ async fn wait_for_completion( }); } - let timeout_future = async { - match timeout_ms { - Some(ms) => time::sleep(Duration::from_millis(ms)).await, - None => std::future::pending::<()>().await, - } - }; + let timeout_future = optional_timeout(timeout_ms); tokio::pin!(timeout_future); loop { tokio::select! { diff --git a/lib/crates/fabro-sandbox/src/docker.rs b/lib/crates/fabro-sandbox/src/docker.rs index 27c25732d..5e9dcce85 100644 --- a/lib/crates/fabro-sandbox/src/docker.rs +++ b/lib/crates/fabro-sandbox/src/docker.rs @@ -2,7 +2,7 @@ use std::collections::HashMap; use std::fmt::Write as _; use std::io::Cursor; use std::sync::atomic::{AtomicU64, Ordering}; -use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; +use std::time::{Instant, SystemTime, UNIX_EPOCH}; use async_trait::async_trait; use bollard::Docker; @@ -24,7 +24,7 @@ use tokio_util::sync::CancellationToken; use crate::clone_source::{self, CloneDecision, EmptyWorkspaceReason}; use crate::redact::redact_auth_url; -use crate::sandbox::resolve_path; +use crate::sandbox::{optional_timeout, resolve_path}; use crate::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingResult, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback, format_lines_numbered, shell_quote, @@ -380,12 +380,7 @@ impl DockerSandbox { controlled_command, ]; - let timeout_future = async { - match timeout_ms { - Some(ms) => time::sleep(Duration::from_millis(ms)).await, - None => std::future::pending::<()>().await, - } - }; + let timeout_future = optional_timeout(timeout_ms); tokio::pin!(timeout_future); let token = cancel_token.unwrap_or_default(); @@ -1545,6 +1540,7 @@ mod tests { reason = "unit test reads an in-memory tar entry synchronously" )] use std::io::Read as _; + use std::time::Duration; use tokio::process::Command; diff --git a/lib/crates/fabro-sandbox/src/local.rs b/lib/crates/fabro-sandbox/src/local.rs index 5c5a6e1ba..44734e9f1 100644 --- a/lib/crates/fabro-sandbox/src/local.rs +++ b/lib/crates/fabro-sandbox/src/local.rs @@ -1,5 +1,5 @@ use std::path::{Path, PathBuf}; -use std::time::{Duration, Instant}; +use std::time::Instant; use async_trait::async_trait; use fabro_static::EnvVars; @@ -10,6 +10,7 @@ use tokio::task::spawn_blocking; use tokio::{fs, time}; use tokio_util::sync::CancellationToken; +use crate::sandbox::optional_timeout; use crate::{ CommandOutputCallback, DirEntry, ExecResult, ExecStreamingResult, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback, format_lines_numbered, @@ -367,12 +368,7 @@ impl Sandbox for LocalSandbox { .spawn() .map_err(|e| crate::Error::context("Failed to spawn command", e))?; - let timeout_future = async { - match timeout_ms { - Some(ms) => time::sleep(Duration::from_millis(ms)).await, - None => std::future::pending::<()>().await, - } - }; + let timeout_future = optional_timeout(timeout_ms); tokio::pin!(timeout_future); let token = cancel_token.unwrap_or_default(); diff --git a/lib/crates/fabro-sandbox/src/sandbox.rs b/lib/crates/fabro-sandbox/src/sandbox.rs index 2ca67bb16..bfd31d91b 100644 --- a/lib/crates/fabro-sandbox/src/sandbox.rs +++ b/lib/crates/fabro-sandbox/src/sandbox.rs @@ -17,6 +17,16 @@ const GIT: &str = "git -c maintenance.auto=0 -c gc.auto=0"; pub const DEFAULT_EXEC_OUTPUT_TAIL_BYTES: usize = 8 * 1024; +/// Sleep for `timeout_ms` if `Some`, otherwise never resolves. Used by +/// streaming `exec_command` impls to model "no timeout" without scheduling a +/// `Duration::from_millis(u64::MAX)` sleep. +pub(crate) async fn optional_timeout(timeout_ms: Option) { + match timeout_ms { + Some(ms) => time::sleep(Duration::from_millis(ms)).await, + None => std::future::pending::<()>().await, + } +} + /// Information returned when a sandbox sets up git for a workflow run. #[derive(Debug, Clone)] pub struct GitRunInfo { diff --git a/lib/crates/fabro-workflow/src/handler/llm/cli.rs b/lib/crates/fabro-workflow/src/handler/llm/cli.rs index 5b97c1a04..22064a100 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/cli.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/cli.rs @@ -1,5 +1,5 @@ use std::collections::HashMap; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use async_trait::async_trait; use fabro_agent::Sandbox; @@ -9,9 +9,30 @@ use fabro_llm::types::TokenCounts; use fabro_model::Provider; use fabro_types::{CommandOutputStream, CommandTermination}; use fabro_util::time::elapsed_ms; -use tokio::sync::Mutex as TokioMutex; use tokio_util::sync::CancellationToken; +/// Returns up to the last `n` characters of `s`, preserving char boundaries. +fn tail_chars(s: &str, n: usize) -> String { + let total = s.chars().count(); + if total <= n { + return s.to_string(); + } + s.chars().skip(total - n).collect() +} + +/// Build a "\nstdout: " detail string for CLI failure +/// messages, falling back to the original command when both streams are empty. +fn cli_failure_detail(stdout: &str, stderr: &str, command: &str) -> String { + let stderr_tail = tail_chars(stderr, 500); + let stdout_tail = tail_chars(stdout, 500); + match (stderr_tail.is_empty(), stdout_tail.is_empty()) { + (false, false) => format!("{stderr_tail}\nstdout: {stdout_tail}"), + (false, true) => stderr_tail, + (true, false) => format!("stdout: {stdout_tail}"), + (true, true) => format!("command: {command}"), + } +} + use super::super::agent::{CodergenBackend, CodergenResult}; use crate::context::Context; use crate::error::Error; @@ -598,8 +619,11 @@ impl CodergenBackend for AgentCliBackend { // `exec_command_streaming` the run-level cancel token (and node // timeout, when set) terminate the CLI and its descendants. let outer_command = format!(". {env_path} && {command}"); - let stdout_buffer: Arc>> = Arc::new(TokioMutex::new(Vec::new())); - let stderr_buffer: Arc>> = Arc::new(TokioMutex::new(Vec::new())); + // Use a synchronous Mutex: each callback invocation only does a short + // `extend_from_slice` with no awaits while the lock is held, so an + // async Mutex would just add per-chunk scheduling overhead. + let stdout_buffer: Arc>> = Arc::new(Mutex::new(Vec::new())); + let stderr_buffer: Arc>> = Arc::new(Mutex::new(Vec::new())); let stdout_buf_cb = Arc::clone(&stdout_buffer); let stderr_buf_cb = Arc::clone(&stderr_buffer); let emitter_for_callback = Arc::clone(emitter); @@ -611,14 +635,13 @@ impl CodergenBackend for AgentCliBackend { // Touch the stall watchdog whenever the CLI emits output // so long-running invocations don't trip stall timeout. emitter.touch(); - match stream { - CommandOutputStream::Stdout => { - stdout_buf.lock().await.extend_from_slice(&bytes); - } - CommandOutputStream::Stderr => { - stderr_buf.lock().await.extend_from_slice(&bytes); - } - } + let buf = match stream { + CommandOutputStream::Stdout => stdout_buf, + CommandOutputStream::Stderr => stderr_buf, + }; + buf.lock() + .expect("CLI output buffer mutex poisoned") + .extend_from_slice(&bytes); Ok(()) }) }); @@ -664,8 +687,18 @@ impl CodergenBackend for AgentCliBackend { let result = streaming.result; // Prefer the buffered streaming output (live chunks); fall back to the // result struct for sandboxes that bundle output at the end. - let buffered_stdout = String::from_utf8_lossy(&stdout_buffer.lock().await).into_owned(); - let buffered_stderr = String::from_utf8_lossy(&stderr_buffer.lock().await).into_owned(); + let buffered_stdout = { + let buf = stdout_buffer + .lock() + .expect("CLI stdout buffer mutex poisoned"); + String::from_utf8_lossy(&buf).into_owned() + }; + let buffered_stderr = { + let buf = stderr_buffer + .lock() + .expect("CLI stderr buffer mutex poisoned"); + String::from_utf8_lossy(&buf).into_owned() + }; let stdout = if buffered_stdout.is_empty() { result.stdout.clone() } else { @@ -676,7 +709,7 @@ impl CodergenBackend for AgentCliBackend { } else { buffered_stderr }; - let duration_ms = u64::try_from(launch_start.elapsed().as_millis()).unwrap_or(u64::MAX); + let duration_ms = elapsed_ms(launch_start); match result.termination { CommandTermination::Cancelled => { @@ -703,23 +736,7 @@ impl CodergenBackend for AgentCliBackend { &stage_scope, ); cleanup_temp_files().await; - let tail = |s: &str, n: usize| -> String { - s.chars() - .rev() - .take(n) - .collect::>() - .into_iter() - .rev() - .collect() - }; - let stderr_tail = tail(&stderr, 500); - let stdout_tail = tail(&stdout, 500); - let detail = match (stderr_tail.is_empty(), stdout_tail.is_empty()) { - (false, false) => format!("{stderr_tail}\nstdout: {stdout_tail}"), - (false, true) => stderr_tail, - (true, false) => format!("stdout: {stdout_tail}"), - (true, true) => format!("command: {command}"), - }; + let detail = cli_failure_detail(&stdout, &stderr, &command); return Err(Error::handler(format!( "CLI command timed out after {duration_ms} ms: {detail}" ))); @@ -744,23 +761,7 @@ impl CodergenBackend for AgentCliBackend { let exited_success = result.termination == CommandTermination::Exited && result.exit_code == Some(0); if !exited_success { - let tail = |s: &str, n: usize| -> String { - s.chars() - .rev() - .take(n) - .collect::>() - .into_iter() - .rev() - .collect() - }; - let stderr_tail = tail(&stderr, 500); - let stdout_tail = tail(&stdout, 500); - let detail = match (stderr_tail.is_empty(), stdout_tail.is_empty()) { - (false, false) => format!("{stderr_tail}\nstdout: {stdout_tail}"), - (false, true) => stderr_tail, - (true, false) => format!("stdout: {stdout_tail}"), - (true, true) => format!("command: {command}"), - }; + let detail = cli_failure_detail(&stdout, &stderr, &command); return Err(Error::handler(format!( "CLI command exited with code {}: {detail}", result diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index 9d1bad3ad..5d241923e 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -66,7 +66,7 @@ impl Handler for PromptHandler { .provider() .and_then(|s| s.parse::().ok()) .unwrap_or(services.run.provider); - let docs = fabro_agent::discover_memory( + let docs = match fabro_agent::discover_memory( &*services.run.sandbox, working_dir, working_dir, @@ -74,7 +74,13 @@ impl Handler for PromptHandler { &services.run.cancel_token(), ) .await - .unwrap_or_default(); + { + Ok(docs) => docs, + Err(fabro_agent::Error::Interrupted(fabro_agent::InterruptReason::Cancelled)) => { + return Err(Error::Cancelled); + } + Err(_) => Vec::new(), + }; if docs.is_empty() { None