diff --git a/lib/components/fabro-sandbox/src/error.rs b/lib/components/fabro-sandbox/src/error.rs index 76f096d4b..bf0d678fc 100644 --- a/lib/components/fabro-sandbox/src/error.rs +++ b/lib/components/fabro-sandbox/src/error.rs @@ -5,6 +5,7 @@ use bollard::errors::Error as BollardError; use fabro_util::error::{collect_causes, render_with_causes}; use crate::ExecResult; +use crate::sandbox::{DEFAULT_EXEC_OUTPUT_TAIL_BYTES, redacted_output_tail}; #[derive(Debug, thiserror::Error)] pub enum Error { @@ -48,6 +49,13 @@ pub enum Error { source: BollardError, }, + /// A sandbox-driver failure: provider, transport, or an operation whose + /// outcome is unknown. The driver's own variants stay reachable through + /// [`Error::driver`] so callers can act on `NotFound`, `Unsupported`, + /// `Transport`, and `Incomplete` without string matching. + #[error(transparent)] + Driver(Box), + #[error( "{label} failed (exit {exit}, termination={termination}, duration_ms={duration_ms}) - hint: {hint}", exit = format_exit_code(result.exit_code), @@ -118,11 +126,67 @@ impl Error { collect_causes(self) } + pub fn driver_error(source: sandbox_driver::Error) -> Self { + Self::Driver(Box::new(source)) + } + + /// The underlying sandbox-driver error, when this error carries one + /// anywhere in its chain. + pub fn driver(&self) -> Option<&sandbox_driver::Error> { + let mut current: Option<&(dyn std::error::Error + 'static)> = Some(self); + while let Some(err) = current { + if let Some(Self::Driver(driver)) = err.downcast_ref::() { + return Some(driver.as_ref()); + } + if let Some(driver) = err.downcast_ref::() { + return Some(driver); + } + current = err.source(); + } + None + } + + /// The facts established when a driver operation ended without a + /// complete outcome. A caller that sees `Some` must not replay the + /// operation: its effects may already have happened. + pub fn incomplete_operation(&self) -> Option<&sandbox_driver::IncompleteOperation> { + match self.driver()? { + sandbox_driver::Error::Incomplete(incomplete) => Some(incomplete), + _ => None, + } + } + + /// True when communication with an out-of-process provider failed. The + /// operation may or may not have run; fabro rebuilds handles through + /// `attach` rather than retrying blind. + pub fn is_transport(&self) -> bool { + matches!(self.driver(), Some(sandbox_driver::Error::Transport(_))) + } + + /// True when the driver reported the resource missing. + pub fn is_not_found(&self) -> bool { + matches!(self.driver(), Some(sandbox_driver::Error::NotFound { .. })) + } + + /// True when the provider does not support the requested capability. + pub fn is_unsupported(&self) -> bool { + matches!( + self.driver(), + Some(sandbox_driver::Error::Unsupported { .. }) + ) + } + pub fn display_with_causes(&self) -> String { render_with_causes(&self.to_string(), &self.causes()) } } +impl From for Error { + fn from(value: sandbox_driver::Error) -> Self { + Self::driver_error(value) + } +} + impl From for Error { fn from(value: String) -> Self { Self::Message(value) @@ -183,6 +247,15 @@ pub fn default_redacted_output_tail( if let Some(Error::Exec { result, .. }) = err.downcast_ref::() { return result.default_redacted_output_tail(); } + if let Some(sandbox_driver::Error::Exec(failure)) = + err.downcast_ref::() + { + return redacted_output_tail( + &String::from_utf8_lossy(failure.stdout()), + &String::from_utf8_lossy(failure.stderr()), + DEFAULT_EXEC_OUTPUT_TAIL_BYTES, + ); + } current = err.source(); } None diff --git a/lib/components/fabro-sandbox/src/exec.rs b/lib/components/fabro-sandbox/src/exec.rs new file mode 100644 index 000000000..64931e90b --- /dev/null +++ b/lib/components/fabro-sandbox/src/exec.rs @@ -0,0 +1,748 @@ +//! Fabro's command execution policy over the sandbox-driver [`Exec`] facet. +//! +//! The driver sends signals; fabro decides when. A command runs as Bash +//! source under `bash -c` with `BASH_ENV` blanked, and ends in one of three +//! ways fabro controls: +//! +//! - **timeout**: fabro's own timer fires, the process group gets `TERM`, and +//! after [`SandboxExec::stop_grace`] it gets `KILL`. The result reports +//! [`CommandTermination::TimedOut`]. The driver's hard timeout is disabled so +//! the graceful ladder always runs first. +//! - **cancellation**: the caller's [`CancellationToken`] runs the same ladder +//! and reports [`CommandTermination::Cancelled`]. +//! - **exit**: the process ended on its own. +//! +//! Output is drained regardless of the retention cap, redacted only when a +//! tail is rendered for events or logs, and delivered live through the +//! caller's callback. Explicit environment variables pass through a +//! fail-closed secret filter under [`ExplicitEnvPolicy::FilterSensitive`], +//! matching what the Host provider already does for inherited variables. + +use std::collections::HashMap; +use std::future; +use std::pin::pin; +use std::sync::{Arc, Mutex, PoisonError}; +use std::time::{Duration, Instant}; + +use async_trait::async_trait; +use fabro_static::EnvVars; +use fabro_types::{CommandOutputStream, CommandTermination}; +use sandbox_driver::{ + BASH_ENV_VAR, CaptureStats, Exec, ExecControls, ExecSpec, OutputStream, SpawnSpec, + StdioProcessHandle as DriverStdioProcessHandle, Termination, TransportError, +}; +use tokio::time; +use tokio_util::sync::CancellationToken; + +use crate::sandbox::{ + CommandOutputCallback, ExecResult, ExecStreamingRequest, ExecStreamingResult, + OutputCaptureStats, StderrCollector, StdioProcess, StdioProcessControl, StdioProcessHandle, + StdioProcessTermination, +}; + +/// Time between `TERM` and `KILL` when fabro stops a command. +pub const DEFAULT_STOP_GRACE: Duration = Duration::from_secs(2); + +/// Retention when a caller sets no cap: enough for any build log fabro +/// renders, bounded so a runaway command cannot exhaust memory. +pub const DEFAULT_RETAINED_OUTPUT_BYTES: usize = sandbox_driver::DEFAULT_BUFFER_BYTES; + +/// How explicit per-command environment variables are treated. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ExplicitEnvPolicy { + /// Drop variables whose names look like credentials unless safelisted. + /// Used where the command runs on the worker host and the caller's env + /// may carry worker secrets. + FilterSensitive, + /// Pass every variable through. Used for isolated providers, where the + /// caller composed the environment deliberately. + TrustCaller, +} + +/// Variables that look like credentials but are needed by ordinary tools. +const ENV_SAFELIST: &[&str] = &[ + EnvVars::PATH, + EnvVars::HOME, + EnvVars::USER, + EnvVars::SHELL, + EnvVars::LANG, + EnvVars::TERM, + EnvVars::TMPDIR, + EnvVars::GOPATH, + EnvVars::CARGO_HOME, + EnvVars::NVM_DIR, +]; + +/// Whether an environment variable name looks like a credential. +#[must_use] +pub fn is_sensitive_env_var(key: &str) -> bool { + if ENV_SAFELIST.contains(&key) { + return false; + } + let lower = key.to_lowercase(); + lower.ends_with("_api_key") + || lower.ends_with("_secret") + || lower.ends_with("_token") + || lower.ends_with("_password") + || lower.ends_with("_credential") +} + +/// Fabro's exec policy bound to one driver [`Exec`] facet. +pub struct SandboxExec<'a> { + exec: &'a dyn Exec, + env_policy: ExplicitEnvPolicy, + stop_grace: Duration, +} + +impl<'a> SandboxExec<'a> { + #[must_use] + pub fn new(exec: &'a dyn Exec, env_policy: ExplicitEnvPolicy) -> Self { + Self { + exec, + env_policy, + stop_grace: DEFAULT_STOP_GRACE, + } + } + + /// Time between `TERM` and `KILL` when a command is stopped. + #[must_use] + pub fn with_stop_grace(mut self, stop_grace: Duration) -> Self { + self.stop_grace = stop_grace; + self + } + + #[must_use] + pub fn stop_grace(&self) -> Duration { + self.stop_grace + } + + /// Runs Bash source to completion and returns its captured output. + /// + /// Equivalent to `bash -c ` with a clean, non-login shell: no + /// `errexit`, no `pipefail`, `BASH_ENV` blanked. A caller that wants + /// different semantics writes them into the command. + pub async fn run( + &self, + command: &str, + timeout: Option, + working_dir: Option<&str>, + env_vars: Option<&HashMap>, + cancel_token: Option, + ) -> crate::Result { + let streaming = self + .run_streaming(ExecStreamingRequest { + timeout_ms: timeout + .map(|timeout| u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX)), + working_dir, + env_vars, + cancel_token, + ..ExecStreamingRequest::new(command) + }) + .await?; + Ok(streaming.result) + } + + /// Runs Bash source, delivering output through `request.output_callback` + /// as it arrives. Same interpreter contract as [`Self::run`]. + pub async fn run_streaming( + &self, + request: ExecStreamingRequest<'_>, + ) -> crate::Result { + let ExecStreamingRequest { + command, + timeout_ms, + working_dir, + env_vars, + cancel_token, + stdin, + output_callback, + stream_output_bytes_cap, + } = request; + let started = Instant::now(); + + let mut spec = ExecSpec::bash(command).no_timeout(); + if let Some(dir) = working_dir { + spec = spec.working_dir(dir); + } + for (key, value) in self.explicit_env(env_vars) { + spec = spec.env_var(key, value); + } + if let Some(bytes) = stdin { + spec = spec.stdin(bytes); + } + + let ladder = StopLadder::new(self.stop_grace); + let controls = ExecControls { + term: Some(ladder.term.clone()), + kill: Some(ladder.kill.clone()), + stdin: None, + sink: output_callback.map(adapt_output_callback), + retained_output_limit: Some( + stream_output_bytes_cap.unwrap_or(DEFAULT_RETAINED_OUTPUT_BYTES), + ), + }; + + let mut escalation = + pin!(ladder.drive(timeout_ms.map(Duration::from_millis), cancel_token)); + let mut running = pin!(self.exec.run_streaming(&spec, controls)); + let streaming = loop { + tokio::select! { + result = &mut running => break result?, + () = &mut escalation => {} + } + }; + + let termination = map_termination(streaming.result.termination, ladder.cause()); + let duration_ms = elapsed_ms(started); + Ok(ExecStreamingResult { + result: ExecResult { + stdout: String::from_utf8_lossy(&streaming.result.stdout).into_owned(), + stderr: String::from_utf8_lossy(&streaming.result.stderr).into_owned(), + exit_code: exit_code_for(termination, streaming.result.exit_code), + termination, + duration_ms, + }, + streams_separated: streaming.streams_separated, + live_streaming: streaming.live_streaming, + stdout_capture: capture_stats(streaming.stdout_capture), + stderr_capture: capture_stats(streaming.stderr_capture), + }) + } + + /// Launches a long-lived process with bidirectional stdio. + /// + /// `command` is evaluated under the same non-login Bash contract before + /// the shell replaces itself with the requested process. Cancelling + /// `cancel_token` terminates the process. + pub async fn spawn_stdio( + &self, + command: &str, + working_dir: Option<&str>, + env_vars: Option<&HashMap>, + cancel_token: Option, + ) -> crate::Result { + let mut spec = SpawnSpec::bash(format!("exec {command}")); + if let Some(dir) = working_dir { + spec = spec.working_dir(dir); + } + for (key, value) in self.explicit_env(env_vars) { + spec = spec.env_var(key, value); + } + let process = self.exec.spawn_stdio(&spec).await?; + let handle = StdioProcessHandle::new(DriverStdioControl { + handle: Arc::from(process.handle), + }); + if let Some(token) = cancel_token { + let handle = handle.clone(); + tokio::spawn(async move { + token.cancelled().await; + if let Err(error) = handle.terminate().await { + tracing::warn!(error = %error, "failed to terminate stdio process on cancel"); + } + }); + } + Ok(StdioProcess { + stdin: process.stdin, + stdout: process.stdout, + stderr: StderrCollector::from_driver_tail(process.stderr_tail), + handle, + }) + } + + /// The explicit environment after policy: `BASH_ENV` never passes, + /// because the Bash helper blanks it and a caller value would override + /// that; credential-shaped names pass only under `TrustCaller`. + fn explicit_env(&self, env_vars: Option<&HashMap>) -> Vec<(String, String)> { + let mut entries: Vec<(String, String)> = env_vars + .into_iter() + .flatten() + .filter(|(key, _)| key.as_str() != BASH_ENV_VAR) + .filter(|(key, _)| { + self.env_policy == ExplicitEnvPolicy::TrustCaller || !is_sensitive_env_var(key) + }) + .map(|(key, value)| (key.clone(), value.clone())) + .collect(); + entries.sort(); + entries + } +} + +/// Why fabro stopped a command, recorded when the ladder starts. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum StopCause { + TimedOut, + Cancelled, +} + +/// The TERM, grace, KILL escalation. Fabro fires `term`, waits `grace`, +/// then fires `kill`; the provider only delivers the signals. +struct StopLadder { + term: CancellationToken, + kill: CancellationToken, + grace: Duration, + cause: Mutex>, +} + +impl StopLadder { + fn new(grace: Duration) -> Self { + Self { + term: CancellationToken::new(), + kill: CancellationToken::new(), + grace, + cause: Mutex::new(None), + } + } + + fn cause(&self) -> Option { + *self.cause.lock().unwrap_or_else(PoisonError::into_inner) + } + + /// Waits for the timeout or the caller's cancellation, runs the ladder, + /// then never resolves so it can sit in a `select!` beside the command. + async fn drive(&self, timeout: Option, cancel_token: Option) { + let cause = tokio::select! { + () = sleep_or_never(timeout) => StopCause::TimedOut, + () = cancelled_or_never(cancel_token.as_ref()) => StopCause::Cancelled, + }; + *self.cause.lock().unwrap_or_else(PoisonError::into_inner) = Some(cause); + self.term.cancel(); + time::sleep(self.grace).await; + self.kill.cancel(); + future::pending::<()>().await; + } +} + +async fn sleep_or_never(timeout: Option) { + match timeout { + Some(timeout) => time::sleep(timeout).await, + None => future::pending().await, + } +} + +async fn cancelled_or_never(token: Option<&CancellationToken>) { + match token { + Some(token) => token.cancelled().await, + None => future::pending().await, + } +} + +/// The driver reports which signal ended the command; fabro reports why it +/// sent it. A stop the driver saw without fabro asking for one (a foreign +/// `kill`, a provider-side abort) reads as cancelled: the command did not +/// finish and fabro did not time it out. +fn map_termination(termination: Termination, cause: Option) -> CommandTermination { + match termination { + Termination::TimedOut => CommandTermination::TimedOut, + Termination::Cancelled | Termination::Killed => match cause { + Some(StopCause::TimedOut) => CommandTermination::TimedOut, + Some(StopCause::Cancelled) | None => CommandTermination::Cancelled, + }, + // `Exited`, or a provider that could not tell how the command ended. + // Nothing asserts success here: `exit_code` is whatever was observed + // and `is_success` still requires `Some(0)`. + _ => CommandTermination::Exited, + } +} + +/// An exit code is only the command's own when it exited on its own. A +/// stopped command may still report the shell's `128 + signal` (143 for a +/// trapped `TERM`), which callers must not mistake for a program result. +fn exit_code_for(termination: CommandTermination, exit_code: Option) -> Option { + match termination { + CommandTermination::Exited => exit_code, + CommandTermination::TimedOut | CommandTermination::Cancelled => None, + } +} + +fn capture_stats(stats: CaptureStats) -> OutputCaptureStats { + OutputCaptureStats { + observed_bytes: stats.observed_bytes, + retained_bytes: stats.retained_bytes, + omitted_bytes: stats.omitted_bytes, + } +} + +fn adapt_output_callback(callback: CommandOutputCallback) -> sandbox_driver::OutputSink { + Arc::new(move |stream, chunk| { + let stream = match stream { + OutputStream::Stdout => CommandOutputStream::Stdout, + OutputStream::Stderr => CommandOutputStream::Stderr, + }; + let callback = Arc::clone(&callback); + Box::pin(async move { + callback(stream, chunk).await.map_err(|error| { + sandbox_driver::Error::Transport(TransportError::with_source( + "command output callback failed", + error, + )) + }) + }) + }) +} + +fn elapsed_ms(started: Instant) -> u64 { + u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX) +} + +struct DriverStdioControl { + handle: Arc, +} + +#[async_trait] +impl StdioProcessControl for DriverStdioControl { + async fn terminate(&self) -> crate::Result<()> { + self.handle.terminate().await; + Ok(()) + } + + async fn wait(&self) -> crate::Result { + let (termination, exit_code) = self.handle.wait().await; + let termination = map_termination(termination, None); + Ok(StdioProcessTermination { + termination, + exit_code: exit_code_for(termination, exit_code), + }) + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use sandbox_driver::{SandboxProvider as _, SandboxSource, SandboxSpec}; + use sandbox_driver_host::HostProvider; + use tokio::fs; + use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; + + use super::*; + + struct HostFixture { + workspace: tempfile::TempDir, + provider: HostProvider, + sandbox: Arc, + } + + impl HostFixture { + async fn new() -> Self { + let workspace = tempfile::tempdir().unwrap(); + let provider = HostProvider::new(); + let sandbox = provider + .create( + &SandboxSpec::new(SandboxSource::HostDirectory) + .working_directory(workspace.path().display().to_string()), + None, + ) + .await + .unwrap(); + Self { + workspace, + provider, + sandbox, + } + } + + fn exec(&self, policy: ExplicitEnvPolicy) -> SandboxExec<'_> { + let _ = &self.provider; + SandboxExec::new(self.sandbox.exec(), policy) + } + } + + async fn run(fixture: &HostFixture, command: &str) -> ExecResult { + fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run(command, Some(Duration::from_secs(10)), None, None, None) + .await + .unwrap() + } + + #[tokio::test] + async fn runs_bash_source_and_reports_exit_code_and_streams() { + let fixture = HostFixture::new().await; + let result = run(&fixture, "echo out; echo err >&2; exit 3").await; + assert_eq!(result.stdout, "out\n"); + assert_eq!(result.stderr, "err\n"); + assert_eq!(result.exit_code, Some(3)); + assert_eq!(result.termination, CommandTermination::Exited); + assert!(!result.is_success()); + assert!(run(&fixture, "true").await.is_success()); + } + + #[tokio::test] + async fn runs_bash_only_syntax_in_a_clean_non_login_shell() { + let fixture = HostFixture::new().await; + let result = run( + &fixture, + "[[ -n ${BASH_VERSION:-} ]] && shopt -q login_shell && echo login || echo nonlogin; \ + set -o | grep -E '^(errexit|pipefail)' | awk '{print $2}' | sort -u", + ) + .await; + assert_eq!(result.stdout, "nonlogin\noff\n", "{result:?}"); + } + + #[tokio::test] + async fn a_caller_supplied_bash_env_never_runs() { + let fixture = HostFixture::new().await; + let startup = fixture.workspace.path().join("startup.sh"); + fs::write(&startup, "echo startup-source-loaded\n") + .await + .unwrap(); + let env = HashMap::from([(BASH_ENV_VAR.to_string(), startup.display().to_string())]); + let result = fixture + .exec(ExplicitEnvPolicy::TrustCaller) + .run( + "echo body", + Some(Duration::from_secs(10)), + None, + Some(&env), + None, + ) + .await + .unwrap(); + assert_eq!(result.stdout, "body\n"); + } + + #[tokio::test] + async fn filter_sensitive_drops_credential_shaped_explicit_variables() { + let fixture = HostFixture::new().await; + let env = HashMap::from([ + ("FABRO_WORKER_TOKEN".to_string(), "leaked".to_string()), + ("MY_VAR".to_string(), "ok".to_string()), + ]); + let filtered = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run("env", Some(Duration::from_secs(10)), None, Some(&env), None) + .await + .unwrap(); + assert!(!filtered.stdout.contains("FABRO_WORKER_TOKEN=leaked")); + assert!(filtered.stdout.contains("MY_VAR=ok")); + + let trusted = fixture + .exec(ExplicitEnvPolicy::TrustCaller) + .run("env", Some(Duration::from_secs(10)), None, Some(&env), None) + .await + .unwrap(); + assert!(trusted.stdout.contains("FABRO_WORKER_TOKEN=leaked")); + } + + #[test] + fn sensitive_name_classification_matches_the_worker_policy() { + for key in [ + "OPENAI_API_KEY", + "DB_PASSWORD", + "AWS_SECRET", + "AUTH_TOKEN", + "MY_CREDENTIAL", + "FABRO_WORKER_TOKEN", + ] { + assert!(is_sensitive_env_var(key), "{key}"); + } + for key in ["PATH", "HOME", "MY_VAR", "GITHUB_ACTOR"] { + assert!(!is_sensitive_env_var(key), "{key}"); + } + } + + #[tokio::test] + async fn timeout_runs_the_ladder_and_reports_timed_out() { + let fixture = HostFixture::new().await; + let started = Instant::now(); + let result = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run( + "sleep 10", + Some(Duration::from_millis(200)), + None, + None, + None, + ) + .await + .unwrap(); + assert_eq!(result.termination, CommandTermination::TimedOut); + assert_eq!(result.exit_code, None); + assert!( + started.elapsed() < Duration::from_secs(5), + "sleep honours TERM, so KILL should not have been needed" + ); + } + + #[tokio::test] + async fn a_command_that_ignores_term_is_killed_after_the_grace_period() { + let fixture = HostFixture::new().await; + let started = Instant::now(); + let result = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .with_stop_grace(Duration::from_millis(300)) + .run( + "trap '' TERM; sleep 10", + Some(Duration::from_millis(100)), + None, + None, + None, + ) + .await + .unwrap(); + assert_eq!(result.termination, CommandTermination::TimedOut); + let elapsed = started.elapsed(); + assert!(elapsed >= Duration::from_millis(400), "{elapsed:?}"); + assert!(elapsed < Duration::from_secs(5), "{elapsed:?}"); + } + + #[tokio::test] + async fn cancellation_reports_cancelled() { + let fixture = HostFixture::new().await; + let token = CancellationToken::new(); + let cancel = token.clone(); + tokio::spawn(async move { + time::sleep(Duration::from_millis(100)).await; + cancel.cancel(); + }); + let result = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run( + "sleep 10", + Some(Duration::from_secs(30)), + None, + None, + Some(token), + ) + .await + .unwrap(); + assert_eq!(result.termination, CommandTermination::Cancelled); + assert_eq!(result.exit_code, None); + } + + #[tokio::test] + async fn streaming_delivers_live_chunks_and_drains_past_the_retention_cap() { + let fixture = HostFixture::new().await; + let seen = Arc::new(Mutex::new(Vec::::new())); + let sink_seen = Arc::clone(&seen); + let callback: CommandOutputCallback = Arc::new(move |stream, chunk| { + let seen = Arc::clone(&sink_seen); + Box::pin(async move { + assert_eq!(stream, CommandOutputStream::Stdout); + seen.lock().unwrap().extend_from_slice(&chunk); + Ok(()) + }) + }); + let streaming = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run_streaming(ExecStreamingRequest { + timeout_ms: Some(10_000), + output_callback: Some(callback), + stream_output_bytes_cap: Some(64), + ..ExecStreamingRequest::new("for i in $(seq 1 200); do echo line-$i; done") + }) + .await + .unwrap(); + assert!(streaming.result.is_success()); + assert!(streaming.live_streaming); + assert!(streaming.streams_separated); + let delivered = seen.lock().unwrap().len(); + assert_eq!(streaming.stdout_capture.observed_bytes, delivered); + assert!(streaming.stdout_capture.omitted_bytes > 0); + assert!(streaming.result.stdout.len() <= 64); + assert!(streaming.result.stdout.starts_with("line-1\n")); + assert!(streaming.result.stdout.ends_with("line-200\n")); + } + + #[tokio::test] + async fn stdin_bytes_are_written_exactly_then_closed() { + let fixture = HostFixture::new().await; + let stdin = b"first line\n$(touch must-not-run)\nlast line".to_vec(); + let streaming = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run_streaming(ExecStreamingRequest { + timeout_ms: Some(10_000), + stdin: Some(stdin.clone()), + ..ExecStreamingRequest::new("cat; test -e must-not-run && echo RAN") + }) + .await + .unwrap(); + assert_eq!(streaming.result.stdout.as_bytes(), stdin.as_slice()); + } + + #[tokio::test] + async fn a_failing_output_callback_stops_the_command_with_an_error() { + let fixture = HostFixture::new().await; + let callback: CommandOutputCallback = + Arc::new(|_, _| Box::pin(async { Err(crate::Error::message("consumer gave up")) })); + let error = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .run_streaming(ExecStreamingRequest { + timeout_ms: Some(10_000), + output_callback: Some(callback), + ..ExecStreamingRequest::new("echo hello; sleep 5") + }) + .await + .map(|streaming| streaming.result.termination); + // The driver either surfaces the sink failure or reports the command + // cancelled by it; both keep the consumer's error visible. + match error { + Ok(termination) => assert_eq!(termination, CommandTermination::Cancelled), + Err(error) => assert!(error.to_string().contains("consumer gave up"), "{error}"), + } + } + + #[tokio::test] + async fn stdio_process_round_trips_lines_and_reports_exit() { + let fixture = HostFixture::new().await; + let process = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .spawn_stdio("cat", None, None, None) + .await + .unwrap(); + let mut stdin = process.stdin; + let mut stdout = BufReader::new(process.stdout); + stdin.write_all(b"ping\n").await.unwrap(); + let mut line = String::new(); + stdout.read_line(&mut line).await.unwrap(); + assert_eq!(line, "ping\n"); + drop(stdin); + let termination = process.handle.wait().await.unwrap(); + assert_eq!(termination.termination, CommandTermination::Exited); + assert_eq!(termination.exit_code, Some(0)); + } + + #[tokio::test] + async fn stdio_process_terminates_on_cancel_and_keeps_a_stderr_tail() { + let fixture = HostFixture::new().await; + let token = CancellationToken::new(); + let process = fixture + .exec(ExplicitEnvPolicy::FilterSensitive) + .spawn_stdio( + "sh -c 'echo diag >&2; sleep 30'", + None, + None, + Some(token.clone()), + ) + .await + .unwrap(); + time::sleep(Duration::from_millis(200)).await; + token.cancel(); + let termination = time::timeout(Duration::from_secs(5), process.handle.wait()) + .await + .expect("cancel terminates the process") + .unwrap(); + assert_ne!(termination.termination, CommandTermination::Exited); + assert_eq!(process.stderr.tail_string().await, "diag\n"); + } + + #[test] + fn termination_mapping_reports_fabro_intent_over_driver_signal() { + assert_eq!( + map_termination(Termination::Killed, Some(StopCause::TimedOut)), + CommandTermination::TimedOut + ); + assert_eq!( + map_termination(Termination::Cancelled, Some(StopCause::Cancelled)), + CommandTermination::Cancelled + ); + assert_eq!( + map_termination(Termination::Killed, None), + CommandTermination::Cancelled + ); + assert_eq!( + map_termination(Termination::Exited, Some(StopCause::TimedOut)), + CommandTermination::Exited + ); + } +} diff --git a/lib/components/fabro-sandbox/src/lib.rs b/lib/components/fabro-sandbox/src/lib.rs index 98ae0473e..3b314ff30 100644 --- a/lib/components/fabro-sandbox/src/lib.rs +++ b/lib/components/fabro-sandbox/src/lib.rs @@ -22,6 +22,8 @@ pub mod details; pub mod driver; +pub mod exec; + pub mod reconnect; pub mod terminal; @@ -41,6 +43,7 @@ pub use details::sandbox_details; #[cfg(feature = "docker")] pub use docker::{DockerSandbox, DockerSandboxOptions}; pub use error::{Error, Result, default_redacted_output_tail, display_for_log}; +pub use exec::{ExplicitEnvPolicy, SandboxExec, is_sensitive_env_var}; pub use fabro_github::token_source::{ InstallationTokenSource, ResolvedToken, TokenProvenance, TokenSnapshot, }; diff --git a/lib/components/fabro-sandbox/src/local.rs b/lib/components/fabro-sandbox/src/local.rs index d2174af6c..e0de7b588 100644 --- a/lib/components/fabro-sandbox/src/local.rs +++ b/lib/components/fabro-sandbox/src/local.rs @@ -12,6 +12,7 @@ use tokio::task::spawn_blocking; use tokio::{fs, time}; use tokio_util::sync::CancellationToken; +use crate::exec::is_sensitive_env_var; use crate::sandbox::{ self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, OutputCaptureBuffer, StdioProcessControl, optional_timeout, validate_bash_probe, write_process_stdin, @@ -142,29 +143,8 @@ impl LocalSandbox { } } - const ENV_SAFELIST: &'static [&'static str] = &[ - EnvVars::PATH, - EnvVars::HOME, - EnvVars::USER, - EnvVars::SHELL, - EnvVars::LANG, - EnvVars::TERM, - EnvVars::TMPDIR, - EnvVars::GOPATH, - EnvVars::CARGO_HOME, - EnvVars::NVM_DIR, - ]; - fn should_filter_env_var(key: &str) -> bool { - if Self::ENV_SAFELIST.contains(&key) { - return false; - } - let lower = key.to_lowercase(); - lower.ends_with("_api_key") - || lower.ends_with("_secret") - || lower.ends_with("_token") - || lower.ends_with("_password") - || lower.ends_with("_credential") + is_sensitive_env_var(key) } fn resolve_path(&self, path: &str) -> PathBuf { diff --git a/lib/components/fabro-sandbox/src/sandbox.rs b/lib/components/fabro-sandbox/src/sandbox.rs index 3cb2cb221..fc7cb55fd 100644 --- a/lib/components/fabro-sandbox/src/sandbox.rs +++ b/lib/components/fabro-sandbox/src/sandbox.rs @@ -1040,31 +1040,63 @@ pub struct StdioProcess { #[derive(Debug, Clone)] pub struct StderrCollector { - inner: Arc>>, - max_bytes: usize, + inner: StderrCollectorInner, +} + +#[derive(Debug, Clone)] +enum StderrCollectorInner { + Buffer { + bytes: Arc>>, + max_bytes: usize, + }, + /// A tail the sandbox driver already keeps for a spawned process. + Driver(sandbox_driver::StderrTail), } impl StderrCollector { #[must_use] pub fn new(max_bytes: usize) -> Self { Self { - inner: Arc::new(TokioMutex::new(Vec::new())), - max_bytes, + inner: StderrCollectorInner::Buffer { + bytes: Arc::new(TokioMutex::new(Vec::new())), + max_bytes, + }, + } + } + + /// Wraps the rolling stderr tail of a driver-spawned process. + #[must_use] + pub fn from_driver_tail(tail: sandbox_driver::StderrTail) -> Self { + Self { + inner: StderrCollectorInner::Driver(tail), } } pub async fn push(&self, bytes: &[u8]) { - let mut tail = self.inner.lock().await; - tail.extend_from_slice(bytes); - if tail.len() > self.max_bytes { - let excess = tail.len() - self.max_bytes; - tail.drain(..excess); + match &self.inner { + StderrCollectorInner::Buffer { + bytes: buffer, + max_bytes, + } => { + let mut tail = buffer.lock().await; + tail.extend_from_slice(bytes); + if tail.len() > *max_bytes { + let excess = tail.len() - max_bytes; + tail.drain(..excess); + } + } + StderrCollectorInner::Driver(tail) => tail.push(bytes), } } pub async fn tail_string(&self) -> String { - let tail = self.inner.lock().await; - String::from_utf8_lossy(&tail).into_owned() + match &self.inner { + StderrCollectorInner::Buffer { bytes, .. } => { + let tail = bytes.lock().await; + String::from_utf8_lossy(&tail).into_owned() + } + StderrCollectorInner::Driver(tail) => tail.to_string_lossy(), + } } pub fn spawn_reader(&self, mut reader: R) -> JoinHandle<()>