fabro(01KQT1V2W1R6ZH72CFT2QDJ39Q): simplify_opus (succeeded)

Fabro-Run: 01KQT1V2W1R6ZH72CFT2QDJ39Q
Fabro-Completed: 6
Fabro-Checkpoint: d0025d74b4

⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
Fabro 2026-05-04 19:31:25 +00:00
parent 1c3eea1053
commit 4cbf76d4b0
6 changed files with 107 additions and 80 deletions

View file

@ -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! {

View file

@ -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;

View file

@ -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();

View file

@ -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<u64>) {
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 {

View file

@ -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 "<stderr-tail>\nstdout: <stdout-tail>" 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<TokioMutex<Vec<u8>>> = Arc::new(TokioMutex::new(Vec::new()));
let stderr_buffer: Arc<TokioMutex<Vec<u8>>> = 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<Mutex<Vec<u8>>> = Arc::new(Mutex::new(Vec::new()));
let stderr_buffer: Arc<Mutex<Vec<u8>>> = 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::<Vec<_>>()
.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::<Vec<_>>()
.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

View file

@ -66,7 +66,7 @@ impl Handler for PromptHandler {
.provider()
.and_then(|s| s.parse::<Provider>().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