From 5df7e178ac1bedc1a9690f87f5b305ebe4a46558 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 7 Mar 2026 09:18:00 -0500 Subject: [PATCH] Add SandboxProvider::Exe for exe.dev VM sandboxes New `arc-exe` crate that runs agent tool operations inside ephemeral exe.dev VMs via SSH. Uses two SSH connections: a management plane (`ssh exe.dev`) for VM lifecycle and a data plane (`ssh vmname.exe.xyz`) for command execution and file I/O. Includes SshRunner trait with MockSshRunner for unit tests and OpensshRunner for real SSH, with raw_mode for the exe.dev management plane's custom command handler. Co-Authored-By: Claude Opus 4.6 (1M context) --- Cargo.lock | 49 +- Cargo.toml | 1 + crates/arc-api/src/demo/mod.rs | 6 + crates/arc-exe/Cargo.toml | 24 + crates/arc-exe/src/lib.rs | 1124 ++++++++++++++++++++ crates/arc-exe/src/openssh_runner.rs | 123 +++ crates/arc-exe/tests/integration.rs | 49 + crates/arc-workflows/Cargo.toml | 1 + crates/arc-workflows/src/cli/mod.rs | 13 + crates/arc-workflows/src/cli/run.rs | 25 +- crates/arc-workflows/src/cli/run_config.rs | 10 + 11 files changed, 1415 insertions(+), 10 deletions(-) create mode 100644 crates/arc-exe/Cargo.toml create mode 100644 crates/arc-exe/src/lib.rs create mode 100644 crates/arc-exe/src/openssh_runner.rs create mode 100644 crates/arc-exe/tests/integration.rs diff --git a/Cargo.lock b/Cargo.lock index 90af8fdf6..d776213db 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -236,6 +236,22 @@ dependencies = [ "tracing", ] +[[package]] +name = "arc-exe" +version = "0.1.0" +dependencies = [ + "arc-agent", + "async-trait", + "base64", + "openssh", + "serde", + "serde_json", + "tempfile", + "tokio", + "tokio-util", + "tracing", +] + [[package]] name = "arc-git-storage" version = "0.1.0" @@ -339,6 +355,7 @@ version = "0.1.0" dependencies = [ "anyhow", "arc-agent", + "arc-exe", "arc-git-storage", "arc-llm", "arc-util", @@ -806,7 +823,7 @@ version = "3.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "faf9468729b8cbcea668e36183cb69d317348c2e08e994829fb56ebfdfbaac34" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] @@ -1273,7 +1290,7 @@ dependencies = [ "libc", "option-ext", "redox_users", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -1360,7 +1377,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -2672,7 +2689,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -2809,6 +2826,20 @@ dependencies = [ "serde_json", ] +[[package]] +name = "openssh" +version = "0.11.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d534c4bfecb0ed71dea4db444a5922a294d15cf40e700548f27295e1feb0ef18" +dependencies = [ + "libc", + "once_cell", + "shell-escape", + "tempfile", + "thiserror 2.0.18", + "tokio", +] + [[package]] name = "openssl" version = "0.10.75" @@ -3217,7 +3248,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.59.0", + "windows-sys 0.60.2", ] [[package]] @@ -3606,7 +3637,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -3674,7 +3705,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -4484,7 +4515,7 @@ dependencies = [ "getrandom 0.4.1", "once_cell", "rustix", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -5363,7 +5394,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index ae4775383..c1cc4bc9d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -48,6 +48,7 @@ toml = "0.8" jsonwebtoken = { version = "10", features = ["aws_lc_rs"] } tokio-tungstenite = { version = "0.26", features = ["rustls-tls-webpki-roots"] } futures-util = "0.3" +openssh = "0.11" daytona-sdk = { git = "https://github.com/brynary/daytona-sdk-rust", package = "daytona-sdk" } daytona-api-client = { git = "https://github.com/brynary/daytona-sdk-rust", package = "daytona-api-client" } diff --git a/crates/arc-api/src/demo/mod.rs b/crates/arc-api/src/demo/mod.rs index 907460dd6..2439bac6a 100644 --- a/crates/arc-api/src/demo/mod.rs +++ b/crates/arc-api/src/demo/mod.rs @@ -933,6 +933,7 @@ mod runs { }), network: Some(arc_workflows::daytona_sandbox::DaytonaNetwork::Block), }), + exe: None, }), vars: Some(std::collections::HashMap::from([ ("repo_url".into(), "https://github.com/org/api-server".into()), @@ -1051,6 +1052,7 @@ mod workflows { }), network: None, }), + exe: None, }), vars: Some(std::collections::HashMap::from([ ("repo_url".into(), "https://github.com/org/service".into()), @@ -1114,6 +1116,7 @@ mod workflows { }), network: None, }), + exe: None, }), vars: Some(std::collections::HashMap::from([ ("spec_path".into(), "specs/feature.md".into()), @@ -1188,6 +1191,7 @@ mod workflows { }), network: None, }), + exe: None, }), vars: Some(std::collections::HashMap::from([ ("source_env".into(), "production".into()), @@ -1253,6 +1257,7 @@ mod workflows { }), network: None, }), + exe: None, }), vars: Some(std::collections::HashMap::from([ ("analytics_window".into(), "30d".into()), @@ -2590,6 +2595,7 @@ mod settings { snapshot: None, network: Some(arc_workflows::daytona_sandbox::DaytonaNetwork::Block), }), + exe: None, }), vars: None, }, diff --git a/crates/arc-exe/Cargo.toml b/crates/arc-exe/Cargo.toml new file mode 100644 index 000000000..ae4419c92 --- /dev/null +++ b/crates/arc-exe/Cargo.toml @@ -0,0 +1,24 @@ +[package] +name = "arc-exe" +edition.workspace = true +version.workspace = true +license.workspace = true +description = "exe.dev VM sandbox for Arc agent tool operations" + +[lib] +doctest = false + +[dependencies] +arc-agent = { path = "../arc-agent" } +async-trait.workspace = true +tokio.workspace = true +tokio-util.workspace = true +openssh.workspace = true +serde_json.workspace = true +base64.workspace = true +tracing.workspace = true +serde.workspace = true + +[dev-dependencies] +tokio = { workspace = true, features = ["test-util", "macros"] } +tempfile = "3" diff --git a/crates/arc-exe/src/lib.rs b/crates/arc-exe/src/lib.rs new file mode 100644 index 000000000..cf5e94099 --- /dev/null +++ b/crates/arc-exe/src/lib.rs @@ -0,0 +1,1124 @@ +mod openssh_runner; + +use std::collections::HashMap; +use std::path::Path; +use std::time::Instant; + +use arc_agent::sandbox::{ + format_lines_numbered, DirEntry, ExecResult, GrepOptions, Sandbox, SandboxEvent, + SandboxEventCallback, +}; +use async_trait::async_trait; +use serde::{Deserialize, Serialize}; +use tokio_util::sync::CancellationToken; + +pub use openssh_runner::OpensshRunner; + +const WORKING_DIRECTORY: &str = "/home/exedev"; +const PROVIDER: &str = "exe"; + +/// Output from an SSH command execution. +pub struct SshOutput { + pub stdout: Vec, + pub stderr: Vec, + pub exit_code: i32, +} + +/// Trait abstracting SSH operations for testability. +#[async_trait] +pub trait SshRunner: Send + Sync { + async fn run_command(&self, command: &str) -> Result; + + async fn run_command_with_timeout( + &self, + command: &str, + timeout: std::time::Duration, + ) -> Result; + + async fn upload_file(&self, path: &str, content: &[u8]) -> Result<(), String>; + + async fn download_file(&self, path: &str) -> Result, String>; +} + +/// Configuration for an exe.dev sandbox (TOML target for `[sandbox.exe]`). +#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)] +pub struct ExeConfig {} + +/// Sandbox that runs all operations inside an exe.dev VM via SSH. +/// +/// Uses two SSH connections: +/// - Management plane (`ssh exe.dev`) for VM lifecycle (create/destroy) +/// - Data plane (`ssh .exe.xyz`) for command execution and file I/O +pub struct ExeSandbox { + mgmt_ssh: Box, + data_ssh: tokio::sync::OnceCell>, + vm_name: tokio::sync::OnceCell, + data_host: tokio::sync::OnceCell, + rg_available: tokio::sync::OnceCell, + event_callback: Option, + /// Factory for creating data-plane SSH runners, used during initialize(). + /// In production, this connects to the VM host via OpensshRunner. + /// In tests, this is replaced with a closure that returns a MockSshRunner. + data_ssh_factory: Box std::pin::Pin, String>> + Send>> + Send + Sync>, +} + +impl ExeSandbox { + /// Creates a new `ExeSandbox` with a management-plane SSH runner. + pub fn new(mgmt_ssh: Box) -> Self { + Self { + mgmt_ssh, + data_ssh: tokio::sync::OnceCell::new(), + vm_name: tokio::sync::OnceCell::new(), + data_host: tokio::sync::OnceCell::new(), + rg_available: tokio::sync::OnceCell::const_new(), + event_callback: None, + data_ssh_factory: Box::new(|host: &str| { + let host = host.to_string(); + Box::pin(async move { + OpensshRunner::connect(&host) + .await + .map(|r| Box::new(r) as Box) + }) + }), + } + } + + pub fn set_event_callback(&mut self, cb: SandboxEventCallback) { + self.event_callback = Some(cb); + } + + fn emit(&self, event: SandboxEvent) { + event.trace(); + if let Some(ref cb) = self.event_callback { + cb(event); + } + } + + /// Get the data-plane SSH runner, returning an error if not yet initialized. + fn data_ssh(&self) -> Result<&dyn SshRunner, String> { + self.data_ssh + .get() + .map(|b| b.as_ref()) + .ok_or_else(|| "Exe sandbox not initialized — call initialize() first".to_string()) + } + + /// Resolve a path: relative paths are prepended with the working directory. + fn resolve_path(&self, path: &str) -> String { + if Path::new(path).is_absolute() { + path.to_string() + } else { + format!("{WORKING_DIRECTORY}/{path}") + } + } +} + +#[async_trait] +impl Sandbox for ExeSandbox { + async fn initialize(&self) -> Result<(), String> { + self.emit(SandboxEvent::Initializing { + provider: PROVIDER.into(), + }); + let init_start = Instant::now(); + + // Create a new VM via the management plane + let output = self + .mgmt_ssh + .run_command("new --json") + .await + .map_err(|e| { + let err = format!("Failed to create exe.dev VM: {e}"); + let duration_ms = + u64::try_from(init_start.elapsed().as_millis()).unwrap_or(u64::MAX); + self.emit(SandboxEvent::InitializeFailed { + provider: PROVIDER.into(), + error: err.clone(), + duration_ms, + }); + err + })?; + + if output.exit_code != 0 { + let stderr = String::from_utf8_lossy(&output.stderr); + let err = format!("exe.dev VM creation failed (exit {}): {stderr}", output.exit_code); + let duration_ms = + u64::try_from(init_start.elapsed().as_millis()).unwrap_or(u64::MAX); + self.emit(SandboxEvent::InitializeFailed { + provider: PROVIDER.into(), + error: err.clone(), + duration_ms, + }); + return Err(err); + } + + // Parse JSON response to get VM name and host + let stdout = String::from_utf8_lossy(&output.stdout); + let json: serde_json::Value = serde_json::from_str(stdout.trim()) + .map_err(|e| { + let err = format!("Failed to parse exe.dev response: {e}"); + let duration_ms = + u64::try_from(init_start.elapsed().as_millis()).unwrap_or(u64::MAX); + self.emit(SandboxEvent::InitializeFailed { + provider: PROVIDER.into(), + error: err.clone(), + duration_ms, + }); + err + })?; + + let vm_name = json["vm_name"] + .as_str() + .ok_or_else(|| "Missing 'vm_name' in exe.dev response".to_string())? + .to_string(); + let data_host = json["ssh_dest"] + .as_str() + .ok_or_else(|| "Missing 'ssh_dest' in exe.dev response".to_string())? + .to_string(); + + self.vm_name + .set(vm_name.clone()) + .map_err(|_| "Exe sandbox already initialized".to_string())?; + self.data_host + .set(data_host.clone()) + .map_err(|_| "Exe sandbox already initialized".to_string())?; + + // Create data-plane SSH connection + let runner = (self.data_ssh_factory)(&data_host).await.map_err(|e| { + let err = format!("Failed to connect to exe.dev VM {data_host}: {e}"); + let duration_ms = + u64::try_from(init_start.elapsed().as_millis()).unwrap_or(u64::MAX); + self.emit(SandboxEvent::InitializeFailed { + provider: PROVIDER.into(), + error: err.clone(), + duration_ms, + }); + err + })?; + self.data_ssh + .set(runner) + .map_err(|_| "Exe sandbox data SSH already set".to_string())?; + + let init_duration = u64::try_from(init_start.elapsed().as_millis()).unwrap_or(u64::MAX); + self.emit(SandboxEvent::Ready { + provider: PROVIDER.into(), + duration_ms: init_duration, + name: Some(vm_name), + cpu: None, + memory: None, + url: None, + }); + + Ok(()) + } + + async fn cleanup(&self) -> Result<(), String> { + self.emit(SandboxEvent::CleanupStarted { + provider: PROVIDER.into(), + }); + let start = Instant::now(); + + if let Some(vm_name) = self.vm_name.get() { + let cmd = format!("rm {vm_name}"); + if let Err(e) = self.mgmt_ssh.run_command(&cmd).await { + let err = format!("Failed to destroy exe.dev VM: {e}"); + self.emit(SandboxEvent::CleanupFailed { + provider: PROVIDER.into(), + error: err.clone(), + }); + return Err(err); + } + } + + let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX); + self.emit(SandboxEvent::CleanupCompleted { + provider: PROVIDER.into(), + duration_ms, + }); + Ok(()) + } + + async fn exec_command( + &self, + command: &str, + timeout_ms: u64, + working_dir: Option<&str>, + env_vars: Option<&HashMap>, + _cancel_token: Option, + ) -> Result { + let ssh = self.data_ssh()?; + let start = Instant::now(); + + // Build the shell command with optional cd and env vars + let mut parts = Vec::new(); + + if let Some(vars) = env_vars { + for (key, value) in vars { + parts.push(format!("{}='{}'", key, value.replace('\'', "'\\''"))); + } + } + + if let Some(dir) = working_dir { + let resolved = self.resolve_path(dir); + parts.push(format!("cd '{}'", resolved.replace('\'', "'\\''"))); + parts.push("&&".to_string()); + } else { + parts.push(format!("cd '{WORKING_DIRECTORY}'")); + parts.push("&&".to_string()); + } + + parts.push(command.to_string()); + let full_cmd = parts.join(" "); + + let timeout = std::time::Duration::from_millis(timeout_ms); + let output = ssh.run_command_with_timeout(&full_cmd, timeout).await; + + let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX); + + match output { + Ok(out) => Ok(ExecResult { + stdout: String::from_utf8_lossy(&out.stdout).to_string(), + stderr: String::from_utf8_lossy(&out.stderr).to_string(), + exit_code: out.exit_code, + timed_out: false, + duration_ms, + }), + Err(e) if e.contains("timed out") => Ok(ExecResult { + stdout: String::new(), + stderr: "Command timed out".to_string(), + exit_code: -1, + timed_out: true, + duration_ms, + }), + Err(e) => Err(e), + } + } + + async fn read_file( + &self, + path: &str, + offset: Option, + limit: Option, + ) -> Result { + let ssh = self.data_ssh()?; + let resolved = self.resolve_path(path); + + let output = ssh + .run_command(&format!("cat '{}'", resolved.replace('\'', "'\\''"))) + .await?; + + if output.exit_code != 0 { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(format!("Failed to read {resolved}: {stderr}")); + } + + let content = + String::from_utf8(output.stdout).map_err(|e| format!("File is not valid UTF-8: {e}"))?; + + Ok(format_lines_numbered(&content, offset, limit)) + } + + async fn write_file(&self, path: &str, content: &str) -> Result<(), String> { + let ssh = self.data_ssh()?; + let resolved = self.resolve_path(path); + + // Ensure parent directory exists + if let Some(parent) = Path::new(&resolved).parent() { + let parent_str = parent.to_string_lossy(); + if parent_str != "/" { + ssh.run_command(&format!("mkdir -p '{}'", parent_str.replace('\'', "'\\''"))) + .await?; + } + } + + ssh.upload_file(&resolved, content.as_bytes()).await + } + + async fn delete_file(&self, path: &str) -> Result<(), String> { + let ssh = self.data_ssh()?; + let resolved = self.resolve_path(path); + + let output = ssh + .run_command(&format!("rm -f '{}'", resolved.replace('\'', "'\\''"))) + .await?; + + if output.exit_code != 0 { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(format!("Failed to delete {resolved}: {stderr}")); + } + Ok(()) + } + + async fn file_exists(&self, path: &str) -> Result { + let ssh = self.data_ssh()?; + let resolved = self.resolve_path(path); + + let output = ssh + .run_command(&format!("test -e '{}'", resolved.replace('\'', "'\\''"))) + .await?; + + Ok(output.exit_code == 0) + } + + async fn list_directory( + &self, + path: &str, + depth: Option, + ) -> Result, String> { + let resolved = self.resolve_path(path); + let max_depth = depth.unwrap_or(1); + + let cmd = format!( + "find '{}' -mindepth 1 -maxdepth {} -printf '%y\\t%s\\t%P\\n'", + resolved.replace('\'', "'\\''"), + max_depth, + ); + + let result = self + .exec_command(&cmd, 30_000, None, None, None) + .await?; + + if result.exit_code != 0 { + return Err(format!( + "Failed to list directory {resolved}: {}", + result.stderr + )); + } + + let mut entries: Vec = result + .stdout + .lines() + .filter(|line| !line.is_empty()) + .filter_map(|line| { + let parts: Vec<&str> = line.splitn(3, '\t').collect(); + if parts.len() < 3 { + return None; + } + let file_type = parts[0]; + let size: Option = parts[1].parse().ok(); + let name = parts[2].to_string(); + let is_dir = file_type == "d"; + Some(DirEntry { + name, + is_dir, + size: if is_dir { None } else { size }, + }) + }) + .collect(); + + entries.sort_by(|a, b| a.name.cmp(&b.name)); + Ok(entries) + } + + async fn grep( + &self, + pattern: &str, + path: &str, + options: &GrepOptions, + ) -> Result, String> { + let resolved = self.resolve_path(path); + + // Detect ripgrep availability (cached) + let use_rg = *self + .rg_available + .get_or_init(|| async { + let result = self + .exec_command("rg --version", 10_000, None, None, None) + .await; + matches!(result, Ok(r) if r.exit_code == 0) + }) + .await; + + let cmd = if use_rg { + let mut cmd = "rg --line-number --no-heading".to_string(); + if options.case_insensitive { + cmd.push_str(" -i"); + } + if let Some(ref glob_filter) = options.glob_filter { + cmd.push_str(&format!(" --glob '{glob_filter}'")); + } + if let Some(max) = options.max_results { + cmd.push_str(&format!(" --max-count {max}")); + } + cmd.push_str(&format!( + " -- '{}' '{}'", + pattern.replace('\'', "'\\''"), + resolved + )); + cmd + } else { + let mut cmd = "grep -rn".to_string(); + if options.case_insensitive { + cmd.push_str(" -i"); + } + if let Some(ref glob_filter) = options.glob_filter { + cmd.push_str(&format!(" --include '{glob_filter}'")); + } + if let Some(max) = options.max_results { + cmd.push_str(&format!(" -m {max}")); + } + cmd.push_str(&format!( + " -- '{}' '{}'", + pattern.replace('\'', "'\\''"), + resolved + )); + cmd + }; + + let result = self.exec_command(&cmd, 30_000, None, None, None).await?; + + if result.exit_code == 1 { + return Ok(Vec::new()); + } + if result.exit_code != 0 { + return Err(format!( + "grep failed (exit {}): {}", + result.exit_code, result.stderr + )); + } + + Ok(result.stdout.lines().map(String::from).collect()) + } + + async fn glob(&self, pattern: &str, path: Option<&str>) -> Result, String> { + let base = path + .map(|p| self.resolve_path(p)) + .unwrap_or_else(|| WORKING_DIRECTORY.to_string()); + + let cmd = format!( + "find '{}' -name '{}' -type f | sort", + base.replace('\'', "'\\''"), + pattern.replace('\'', "'\\''"), + ); + + let result = self.exec_command(&cmd, 30_000, None, None, None).await?; + + if result.exit_code != 0 { + return Err(format!( + "glob failed (exit {}): {}", + result.exit_code, result.stderr + )); + } + + Ok(result + .stdout + .lines() + .filter(|l| !l.is_empty()) + .map(String::from) + .collect()) + } + + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &Path, + ) -> Result<(), String> { + let ssh = self.data_ssh()?; + let resolved = self.resolve_path(remote_path); + + let bytes = ssh.download_file(&resolved).await?; + + if let Some(parent) = local_path.parent() { + tokio::fs::create_dir_all(parent) + .await + .map_err(|e| format!("Failed to create parent dirs: {e}"))?; + } + tokio::fs::write(local_path, &bytes) + .await + .map_err(|e| format!("Failed to write {}: {e}", local_path.display()))?; + + Ok(()) + } + + fn working_directory(&self) -> &str { + WORKING_DIRECTORY + } + + fn platform(&self) -> &str { + "linux" + } + + fn os_version(&self) -> String { + "Linux (exe.dev)".to_string() + } + + fn sandbox_info(&self) -> String { + self.vm_name + .get() + .cloned() + .unwrap_or_default() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::{Arc, Mutex}; + + /// A recorded command sent to the mock SSH runner. + #[derive(Debug, Clone)] + struct RecordedCommand { + command: String, + } + + /// A queued response for MockSshRunner. + struct MockResponse { + stdout: Vec, + stderr: Vec, + exit_code: i32, + } + + /// Mock upload record. + #[derive(Debug, Clone)] + struct RecordedUpload { + path: String, + content: Vec, + } + + /// Mock download response. + struct MockDownload { + content: Vec, + } + + /// Mock SSH runner for unit tests. + struct MockSshRunner { + commands: Arc>>, + responses: Arc>>, + uploads: Arc>>, + downloads: Arc>>, + } + + impl MockSshRunner { + fn new() -> Self { + Self { + commands: Arc::new(Mutex::new(Vec::new())), + responses: Arc::new(Mutex::new(Vec::new())), + uploads: Arc::new(Mutex::new(Vec::new())), + downloads: Arc::new(Mutex::new(Vec::new())), + } + } + + fn queue_response(&self, stdout: &str, stderr: &str, exit_code: i32) { + self.responses.lock().unwrap().push(MockResponse { + stdout: stdout.as_bytes().to_vec(), + stderr: stderr.as_bytes().to_vec(), + exit_code, + }); + } + + fn queue_response_bytes(&self, stdout: Vec, stderr: &str, exit_code: i32) { + self.responses.lock().unwrap().push(MockResponse { + stdout, + stderr: stderr.as_bytes().to_vec(), + exit_code, + }); + } + + fn queue_download(&self, content: Vec) { + self.downloads + .lock() + .unwrap() + .push(MockDownload { content }); + } + + fn pop_response(&self) -> MockResponse { + let mut responses = self.responses.lock().unwrap(); + if responses.is_empty() { + MockResponse { + stdout: Vec::new(), + stderr: b"no mock response queued".to_vec(), + exit_code: 1, + } + } else { + responses.remove(0) + } + } + } + + #[async_trait] + impl SshRunner for MockSshRunner { + async fn run_command(&self, command: &str) -> Result { + self.commands.lock().unwrap().push(RecordedCommand { + command: command.to_string(), + }); + let resp = self.pop_response(); + Ok(SshOutput { + stdout: resp.stdout, + stderr: resp.stderr, + exit_code: resp.exit_code, + }) + } + + async fn run_command_with_timeout( + &self, + command: &str, + _timeout: std::time::Duration, + ) -> Result { + self.commands.lock().unwrap().push(RecordedCommand { + command: command.to_string(), + }); + let resp = self.pop_response(); + if resp.exit_code == -99 { + return Err("Command timed out".to_string()); + } + Ok(SshOutput { + stdout: resp.stdout, + stderr: resp.stderr, + exit_code: resp.exit_code, + }) + } + + async fn upload_file(&self, path: &str, content: &[u8]) -> Result<(), String> { + self.uploads.lock().unwrap().push(RecordedUpload { + path: path.to_string(), + content: content.to_vec(), + }); + Ok(()) + } + + async fn download_file(&self, _path: &str) -> Result, String> { + let mut downloads = self.downloads.lock().unwrap(); + if downloads.is_empty() { + Err("no mock download queued".to_string()) + } else { + Ok(downloads.remove(0).content) + } + } + } + + /// Helper: create an ExeSandbox with mock data SSH already initialized (skipping lifecycle). + fn sandbox_with_mock_data(data_ssh: MockSshRunner) -> ExeSandbox { + let mgmt = MockSshRunner::new(); + let sandbox = ExeSandbox::new(Box::new(mgmt)); + let _ = sandbox.vm_name.set("test-vm".to_string()); + let _ = sandbox.data_host.set("test-vm.exe.xyz".to_string()); + let _ = sandbox.data_ssh.set(Box::new(data_ssh)); + sandbox + } + + // ---- Step 1: Metadata accessors ---- + + #[test] + fn working_directory_returns_home_user() { + let sandbox = sandbox_with_mock_data(MockSshRunner::new()); + assert_eq!(sandbox.working_directory(), "/home/exedev"); + } + + #[test] + fn platform_returns_linux() { + let sandbox = sandbox_with_mock_data(MockSshRunner::new()); + assert_eq!(sandbox.platform(), "linux"); + } + + #[test] + fn sandbox_info_returns_vm_name() { + let sandbox = sandbox_with_mock_data(MockSshRunner::new()); + assert_eq!(sandbox.sandbox_info(), "test-vm"); + } + + #[test] + fn os_version_returns_linux_exe() { + let sandbox = sandbox_with_mock_data(MockSshRunner::new()); + assert_eq!(sandbox.os_version(), "Linux (exe.dev)"); + } + + // ---- Step 2: exec_command ---- + + #[tokio::test] + async fn exec_command_runs_via_ssh() { + let data = MockSshRunner::new(); + data.queue_response("hello world\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let result = sandbox + .exec_command("echo hello world", 5000, None, None, None) + .await + .unwrap(); + + assert_eq!(result.stdout.trim(), "hello world"); + assert_eq!(result.exit_code, 0); + assert!(!result.timed_out); + } + + #[tokio::test] + async fn exec_command_with_working_dir() { + let data = MockSshRunner::new(); + let commands = data.commands.clone(); + data.queue_response("", "", 0); + let sandbox = sandbox_with_mock_data(data); + + sandbox + .exec_command("ls", 5000, Some("/tmp/work"), None, None) + .await + .unwrap(); + + let recorded = commands.lock().unwrap(); + assert!( + recorded[0].command.contains("cd '/tmp/work'"), + "expected cd to working dir, got: {}", + recorded[0].command, + ); + } + + #[tokio::test] + async fn exec_command_with_env_vars() { + let data = MockSshRunner::new(); + let commands = data.commands.clone(); + data.queue_response("", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let mut env = HashMap::new(); + env.insert("FOO".to_string(), "bar".to_string()); + + sandbox + .exec_command("echo $FOO", 5000, None, Some(&env), None) + .await + .unwrap(); + + let recorded = commands.lock().unwrap(); + assert!( + recorded[0].command.contains("FOO='bar'"), + "expected env var, got: {}", + recorded[0].command, + ); + } + + #[tokio::test] + async fn exec_command_timeout() { + let data = MockSshRunner::new(); + // Use exit code -99 as the sentinel for timeout in our mock + data.queue_response_bytes(Vec::new(), "", -99); + let sandbox = sandbox_with_mock_data(data); + + let result = sandbox + .exec_command("sleep 999", 100, None, None, None) + .await + .unwrap(); + + assert!(result.timed_out); + assert_eq!(result.exit_code, -1); + } + + // ---- Step 3: read_file ---- + + #[tokio::test] + async fn read_file_returns_numbered_lines() { + let data = MockSshRunner::new(); + data.queue_response("line one\nline two\nline three\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let content = sandbox.read_file("test.txt", None, None).await.unwrap(); + assert!(content.contains("1 | line one")); + assert!(content.contains("2 | line two")); + assert!(content.contains("3 | line three")); + } + + #[tokio::test] + async fn read_file_with_offset_and_limit() { + let data = MockSshRunner::new(); + data.queue_response("a\nb\nc\nd\ne\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let content = sandbox.read_file("test.txt", Some(1), Some(2)).await.unwrap(); + assert!(content.contains("2 | b")); + assert!(content.contains("3 | c")); + assert!(!content.contains("1 | a")); + assert!(!content.contains("4 | d")); + } + + #[tokio::test] + async fn read_file_absolute_path() { + let data = MockSshRunner::new(); + let commands = data.commands.clone(); + data.queue_response("content\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + sandbox + .read_file("/etc/hosts", None, None) + .await + .unwrap(); + + let recorded = commands.lock().unwrap(); + assert!( + recorded[0].command.contains("/etc/hosts"), + "expected absolute path, got: {}", + recorded[0].command, + ); + assert!( + !recorded[0].command.contains("/home/user"), + "should not prepend working dir for absolute path", + ); + } + + // ---- Step 4: write_file ---- + + #[tokio::test] + async fn write_file_uploads_content() { + let data = MockSshRunner::new(); + let uploads = data.uploads.clone(); + // Response for mkdir -p + data.queue_response("", "", 0); + let sandbox = sandbox_with_mock_data(data); + + sandbox.write_file("src/main.rs", "fn main() {}").await.unwrap(); + + let recorded = uploads.lock().unwrap(); + assert_eq!(recorded[0].path, "/home/exedev/src/main.rs"); + assert_eq!(recorded[0].content, b"fn main() {}"); + } + + #[tokio::test] + async fn write_file_creates_parent_dirs() { + let data = MockSshRunner::new(); + let commands = data.commands.clone(); + // Response for mkdir -p + data.queue_response("", "", 0); + let sandbox = sandbox_with_mock_data(data); + + sandbox + .write_file("deep/nested/file.txt", "content") + .await + .unwrap(); + + let recorded = commands.lock().unwrap(); + assert!( + recorded[0].command.contains("mkdir -p"), + "expected mkdir -p, got: {}", + recorded[0].command, + ); + assert!( + recorded[0].command.contains("/home/exedev/deep/nested"), + "expected parent path, got: {}", + recorded[0].command, + ); + } + + // ---- Step 5: delete_file + file_exists ---- + + #[tokio::test] + async fn delete_file_runs_rm() { + let data = MockSshRunner::new(); + let commands = data.commands.clone(); + data.queue_response("", "", 0); + let sandbox = sandbox_with_mock_data(data); + + sandbox.delete_file("old.txt").await.unwrap(); + + let recorded = commands.lock().unwrap(); + assert!( + recorded[0].command.contains("rm -f"), + "expected rm -f, got: {}", + recorded[0].command, + ); + } + + #[tokio::test] + async fn file_exists_true() { + let data = MockSshRunner::new(); + data.queue_response("", "", 0); + let sandbox = sandbox_with_mock_data(data); + + assert!(sandbox.file_exists("exists.txt").await.unwrap()); + } + + #[tokio::test] + async fn file_exists_false() { + let data = MockSshRunner::new(); + data.queue_response("", "", 1); + let sandbox = sandbox_with_mock_data(data); + + assert!(!sandbox.file_exists("missing.txt").await.unwrap()); + } + + // ---- Step 6: list_directory ---- + + #[tokio::test] + async fn list_directory_parses_find_output() { + let data = MockSshRunner::new(); + // find output for list_directory (run via exec_command, so two responses: + // first for the rg_available check if it fires... but exec_command calls + // run_command_with_timeout directly, which will get the next response) + data.queue_response("f\t1024\tfile.txt\nd\t4096\tsrc\nf\t512\tREADME.md\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let entries = sandbox.list_directory(".", None).await.unwrap(); + assert_eq!(entries.len(), 3); + // Sorted alphabetically + assert_eq!(entries[0].name, "README.md"); + assert!(!entries[0].is_dir); + assert_eq!(entries[0].size, Some(512)); + assert_eq!(entries[1].name, "file.txt"); + assert_eq!(entries[2].name, "src"); + assert!(entries[2].is_dir); + assert!(entries[2].size.is_none()); + } + + // ---- Step 7: grep ---- + + #[tokio::test] + async fn grep_returns_matches() { + let data = MockSshRunner::new(); + // First call: rg --version check (cached) + data.queue_response("ripgrep 14.0.0", "", 0); + // Second call: the actual grep + data.queue_response("src/main.rs:1:fn main() {}\nsrc/lib.rs:5:fn helper() {}\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let results = sandbox + .grep("fn ", ".", &GrepOptions::default()) + .await + .unwrap(); + + assert_eq!(results.len(), 2); + assert!(results[0].contains("main.rs")); + } + + #[tokio::test] + async fn grep_no_matches_returns_empty() { + let data = MockSshRunner::new(); + // rg --version + data.queue_response("ripgrep 14.0.0", "", 0); + // grep with no matches (exit code 1) + data.queue_response("", "", 1); + let sandbox = sandbox_with_mock_data(data); + + let results = sandbox + .grep("nonexistent", ".", &GrepOptions::default()) + .await + .unwrap(); + + assert!(results.is_empty()); + } + + // ---- Step 8: glob ---- + + #[tokio::test] + async fn glob_finds_files() { + let data = MockSshRunner::new(); + data.queue_response("/home/user/src/main.rs\n/home/user/src/lib.rs\n", "", 0); + let sandbox = sandbox_with_mock_data(data); + + let results = sandbox.glob("*.rs", Some("src")).await.unwrap(); + + assert_eq!(results.len(), 2); + assert!(results[0].contains("main.rs")); + } + + // ---- Step 9: download_file_to_local ---- + + #[tokio::test] + async fn download_file_to_local_writes_bytes() { + let data = MockSshRunner::new(); + data.queue_download(b"binary content".to_vec()); + let sandbox = sandbox_with_mock_data(data); + + let tmp = tempfile::tempdir().unwrap(); + let local = tmp.path().join("downloaded.bin"); + sandbox + .download_file_to_local("artifact.bin", &local) + .await + .unwrap(); + + let bytes = tokio::fs::read(&local).await.unwrap(); + assert_eq!(bytes, b"binary content"); + } + + // ---- Step 10: initialize + cleanup (VM lifecycle) ---- + + #[tokio::test] + async fn initialize_creates_vm() { + let mgmt = MockSshRunner::new(); + mgmt.queue_response( + r#"{"vm_name": "my-vm", "ssh_dest": "my-vm.exe.xyz"}"#, + "", + 0, + ); + + let data_for_init = MockSshRunner::new(); + + let mut sandbox = ExeSandbox::new(Box::new(mgmt)); + // Override factory to return our mock data SSH + let data_box: Arc>>> = + Arc::new(Mutex::new(Some(Box::new(data_for_init)))); + sandbox.data_ssh_factory = Box::new(move |_host: &str| { + let data_box = Arc::clone(&data_box); + Box::pin(async move { + data_box + .lock() + .unwrap() + .take() + .ok_or_else(|| "mock data SSH already taken".to_string()) + }) + }); + + sandbox.initialize().await.unwrap(); + + assert_eq!(sandbox.sandbox_info(), "my-vm"); + } + + #[tokio::test] + async fn initialize_emits_events() { + let mgmt = MockSshRunner::new(); + mgmt.queue_response( + r#"{"vm_name": "ev-vm", "ssh_dest": "ev-vm.exe.xyz"}"#, + "", + 0, + ); + + let events: Arc>> = Arc::new(Mutex::new(Vec::new())); + let events_cb = Arc::clone(&events); + + let mut sandbox = ExeSandbox::new(Box::new(mgmt)); + sandbox.set_event_callback(Arc::new(move |event| { + events_cb.lock().unwrap().push(format!("{event:?}")); + })); + + let data_for_init = MockSshRunner::new(); + let data_box: Arc>>> = + Arc::new(Mutex::new(Some(Box::new(data_for_init)))); + sandbox.data_ssh_factory = Box::new(move |_host: &str| { + let data_box = Arc::clone(&data_box); + Box::pin(async move { + data_box + .lock() + .unwrap() + .take() + .ok_or_else(|| "mock data SSH already taken".to_string()) + }) + }); + + sandbox.initialize().await.unwrap(); + + let captured = events.lock().unwrap(); + assert!( + captured.iter().any(|e| e.contains("Initializing")), + "expected Initializing event, got: {captured:?}" + ); + assert!( + captured.iter().any(|e| e.contains("Ready")), + "expected Ready event, got: {captured:?}" + ); + } + + #[tokio::test] + async fn cleanup_destroys_vm() { + let mgmt = MockSshRunner::new(); + let mgmt_commands = mgmt.commands.clone(); + // Response for `rm ` + mgmt.queue_response("", "", 0); + + let sandbox = ExeSandbox::new(Box::new(mgmt)); + let _ = sandbox.vm_name.set("doomed-vm".to_string()); + + sandbox.cleanup().await.unwrap(); + + let recorded = mgmt_commands.lock().unwrap(); + assert_eq!(recorded[0].command, "rm doomed-vm"); + } + + #[tokio::test] + async fn cleanup_before_initialize_is_noop() { + let mgmt = MockSshRunner::new(); + let sandbox = ExeSandbox::new(Box::new(mgmt)); + // Should not error — no VM to destroy + sandbox.cleanup().await.unwrap(); + } +} diff --git a/crates/arc-exe/src/openssh_runner.rs b/crates/arc-exe/src/openssh_runner.rs new file mode 100644 index 000000000..e5e9fd72e --- /dev/null +++ b/crates/arc-exe/src/openssh_runner.rs @@ -0,0 +1,123 @@ +use async_trait::async_trait; +use openssh::{KnownHosts, Session}; + +use crate::{SshOutput, SshRunner}; + +/// Real SSH implementation using the `openssh` crate (multiplexed connections). +pub struct OpensshRunner { + session: Session, + /// When true, commands are sent via `raw_command` (no shell wrapping). + /// Used for the exe.dev management plane which has a custom SSH command handler. + raw_mode: bool, +} + +impl OpensshRunner { + /// Connect to a host via SSH, using the user's SSH agent for authentication. + /// Commands are executed through a shell (`sh -c`). + pub async fn connect(host: &str) -> Result { + let session = Session::connect(host, KnownHosts::Accept) + .await + .map_err(|e| format!("SSH connection to {host} failed: {e}"))?; + Ok(Self { + session, + raw_mode: false, + }) + } + + /// Connect to a host via SSH in raw mode (no shell wrapping). + /// Commands are sent directly as the SSH command string. + /// Used for the exe.dev management plane which has a custom command handler. + pub async fn connect_raw(host: &str) -> Result { + let session = Session::connect(host, KnownHosts::Accept) + .await + .map_err(|e| format!("SSH connection to {host} failed: {e}"))?; + Ok(Self { + session, + raw_mode: true, + }) + } + + fn build_command(&self, command: &str) -> openssh::OwningCommand<&Session> { + if self.raw_mode { + self.session.raw_command(command) + } else { + self.session.shell(command) + } + } +} + +#[async_trait] +impl SshRunner for OpensshRunner { + async fn run_command(&self, command: &str) -> Result { + let output = self + .build_command(command) + .output() + .await + .map_err(|e| format!("SSH command failed: {e}"))?; + + let exit_code = output.status.code().unwrap_or(-1); + Ok(SshOutput { + stdout: output.stdout, + stderr: output.stderr, + exit_code, + }) + } + + async fn run_command_with_timeout( + &self, + command: &str, + timeout: std::time::Duration, + ) -> Result { + let mut child = self.build_command(command); + let fut = child.output(); + + match tokio::time::timeout(timeout, fut).await { + Ok(Ok(output)) => { + let exit_code = output.status.code().unwrap_or(-1); + Ok(SshOutput { + stdout: output.stdout, + stderr: output.stderr, + exit_code, + }) + } + Ok(Err(e)) => Err(format!("SSH command failed: {e}")), + Err(_) => Err("Command timed out".to_string()), + } + } + + async fn upload_file(&self, path: &str, content: &[u8]) -> Result<(), String> { + use base64::Engine; + let encoded = base64::engine::general_purpose::STANDARD.encode(content); + let cmd = format!( + "echo '{}' | base64 -d > '{}'", + encoded, + path.replace('\'', "'\\''"), + ); + let output = self + .build_command(&cmd) + .output() + .await + .map_err(|e| format!("SSH upload failed: {e}"))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(format!("Upload to {path} failed: {stderr}")); + } + Ok(()) + } + + async fn download_file(&self, path: &str) -> Result, String> { + let cmd = format!("cat '{}'", path.replace('\'', "'\\''")); + let output = self + .build_command(&cmd) + .output() + .await + .map_err(|e| format!("SSH download failed: {e}"))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(format!("Download of {path} failed: {stderr}")); + } + Ok(output.stdout) + } +} diff --git a/crates/arc-exe/tests/integration.rs b/crates/arc-exe/tests/integration.rs new file mode 100644 index 000000000..545bf53f8 --- /dev/null +++ b/crates/arc-exe/tests/integration.rs @@ -0,0 +1,49 @@ +use arc_agent::sandbox::Sandbox; +use arc_exe::{ExeSandbox, OpensshRunner}; + +/// Full lifecycle test against a real exe.dev account. +/// Requires SSH agent with exe.dev credentials. +/// +/// Run with: cargo test -p arc-exe -- --ignored +#[tokio::test] +#[ignore] +async fn exe_sandbox_full_lifecycle() { + let mgmt_ssh = OpensshRunner::connect_raw("exe.dev") + .await + .expect("SSH to exe.dev failed — is your SSH agent running?"); + + let sandbox = ExeSandbox::new(Box::new(mgmt_ssh)); + + // Initialize (creates VM) + sandbox.initialize().await.unwrap(); + assert!(!sandbox.sandbox_info().is_empty()); + assert_eq!(sandbox.platform(), "linux"); + + // exec_command + let result = sandbox + .exec_command("echo hello", 10_000, None, None, None) + .await + .unwrap(); + assert_eq!(result.stdout.trim(), "hello"); + assert_eq!(result.exit_code, 0); + + // write_file + read_file + sandbox + .write_file("test.txt", "line1\nline2\nline3") + .await + .unwrap(); + let content = sandbox.read_file("test.txt", None, None).await.unwrap(); + assert!(content.contains("1 | line1")); + assert!(content.contains("2 | line2")); + + // file_exists + assert!(sandbox.file_exists("test.txt").await.unwrap()); + assert!(!sandbox.file_exists("nonexistent.txt").await.unwrap()); + + // delete_file + sandbox.delete_file("test.txt").await.unwrap(); + assert!(!sandbox.file_exists("test.txt").await.unwrap()); + + // Cleanup (destroys VM) + sandbox.cleanup().await.unwrap(); +} diff --git a/crates/arc-workflows/Cargo.toml b/crates/arc-workflows/Cargo.toml index 3403300f6..c162c9d74 100644 --- a/crates/arc-workflows/Cargo.toml +++ b/crates/arc-workflows/Cargo.toml @@ -17,6 +17,7 @@ clap.workspace = true anyhow.workspace = true dotenvy.workspace = true arc-agent = { path = "../arc-agent" } +arc-exe = { path = "../arc-exe" } arc-util = { path = "../arc-util" } arc-git-storage = { path = "../arc-git-storage" } arc-llm = { path = "../arc-llm" } diff --git a/crates/arc-workflows/src/cli/mod.rs b/crates/arc-workflows/src/cli/mod.rs index 35b15c3d2..837737e9a 100644 --- a/crates/arc-workflows/src/cli/mod.rs +++ b/crates/arc-workflows/src/cli/mod.rs @@ -28,6 +28,8 @@ pub enum SandboxProvider { Docker, /// Run tools inside a Daytona cloud sandbox Daytona, + /// Run tools inside an exe.dev VM + Exe, } impl fmt::Display for SandboxProvider { @@ -36,6 +38,7 @@ impl fmt::Display for SandboxProvider { Self::Local => write!(f, "local"), Self::Docker => write!(f, "docker"), Self::Daytona => write!(f, "daytona"), + Self::Exe => write!(f, "exe"), } } } @@ -48,6 +51,7 @@ impl FromStr for SandboxProvider { "local" => Ok(Self::Local), "docker" => Ok(Self::Docker), "daytona" => Ok(Self::Daytona), + "exe" => Ok(Self::Exe), other => Err(format!("unknown sandbox provider: {other}")), } } @@ -255,6 +259,14 @@ mod tests { "LOCAL".parse::().unwrap(), SandboxProvider::Local ); + assert_eq!( + "exe".parse::().unwrap(), + SandboxProvider::Exe + ); + assert_eq!( + "EXE".parse::().unwrap(), + SandboxProvider::Exe + ); assert!("invalid".parse::().is_err()); } @@ -263,6 +275,7 @@ mod tests { assert_eq!(SandboxProvider::Local.to_string(), "local"); assert_eq!(SandboxProvider::Docker.to_string(), "docker"); assert_eq!(SandboxProvider::Daytona.to_string(), "daytona"); + assert_eq!(SandboxProvider::Exe.to_string(), "exe"); } #[test] diff --git a/crates/arc-workflows/src/cli/run.rs b/crates/arc-workflows/src/cli/run.rs index bbc065927..94bb3f8d9 100644 --- a/crates/arc-workflows/src/cli/run.rs +++ b/crates/arc-workflows/src/cli/run.rs @@ -267,7 +267,7 @@ pub async fn run_command( SandboxProvider::Local | SandboxProvider::Docker => { crate::git::ensure_clean(&original_cwd).is_ok() } - SandboxProvider::Daytona => false, + SandboxProvider::Daytona | SandboxProvider::Exe => false, }; if args.preflight { @@ -457,6 +457,17 @@ pub async fn run_command( daytona_sandbox_ref = Some(Arc::clone(&daytona_arc)); daytona_arc } + SandboxProvider::Exe => { + let mgmt_ssh = arc_exe::OpensshRunner::connect_raw("exe.dev") + .await + .map_err(|e| anyhow::anyhow!("Failed to connect to exe.dev: {e}"))?; + let mut env = arc_exe::ExeSandbox::new(Box::new(mgmt_ssh)); + let emitter_cb = Arc::clone(&emitter); + env.set_event_callback(Arc::new(move |event| { + emitter_cb.emit(&crate::event::WorkflowRunEvent::Sandbox { event }); + })); + Arc::new(env) + } SandboxProvider::Local => { let mut env = LocalSandbox::new(cwd); let emitter_cb = Arc::clone(&emitter); @@ -667,6 +678,7 @@ pub async fn run_command( SandboxProvider::Daytona => daytona_base_sha .as_ref() .map(|_| GitCheckpointMode::Remote(original_cwd.clone())), + SandboxProvider::Exe => None, }, base_sha: worktree_base_sha.or(daytona_base_sha), run_branch: worktree_branch.or(daytona_branch), @@ -1202,6 +1214,13 @@ async fn run_preflight( } Err(e) => Err(format!("Daytona client creation failed: {e}")), }, + SandboxProvider::Exe => match arc_exe::OpensshRunner::connect_raw("exe.dev").await { + Ok(mgmt_ssh) => { + let env = arc_exe::ExeSandbox::new(Box::new(mgmt_ssh)); + Ok(Arc::new(env) as Arc) + } + Err(e) => Err(format!("exe.dev SSH connection failed: {e}")), + }, SandboxProvider::Local => { Ok(Arc::new(LocalSandbox::new(original_cwd.clone())) as Arc) } @@ -1653,6 +1672,7 @@ mod tests { provider: None, preserve: Some(false), daytona: None, + exe: None, }), vars: None, hooks: Vec::new(), @@ -1674,6 +1694,7 @@ mod tests { provider: None, preserve: Some(true), daytona: None, + exe: None, }), vars: None, hooks: Vec::new(), @@ -1683,6 +1704,7 @@ mod tests { provider: None, preserve: Some(false), daytona: None, + exe: None, }), ..RunDefaults::default() }; @@ -1696,6 +1718,7 @@ mod tests { provider: None, preserve: Some(true), daytona: None, + exe: None, }), ..RunDefaults::default() }; diff --git a/crates/arc-workflows/src/cli/run_config.rs b/crates/arc-workflows/src/cli/run_config.rs index f3ad01595..1ae9ede7d 100644 --- a/crates/arc-workflows/src/cli/run_config.rs +++ b/crates/arc-workflows/src/cli/run_config.rs @@ -42,6 +42,7 @@ pub struct SandboxConfig { pub provider: Option, pub preserve: Option, pub daytona: Option, + pub exe: Option, } /// Defaults for workflow runs, loaded from the server config. @@ -754,6 +755,7 @@ preserve = true provider: None, preserve: Some(false), daytona: None, + exe: None, }), ..RunDefaults::default() }; @@ -779,6 +781,7 @@ provider = "docker" provider: None, preserve: Some(true), daytona: None, + exe: None, }), ..RunDefaults::default() }; @@ -807,6 +810,7 @@ provider = "daytona" auto_stop_interval: Some(30), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() }; @@ -842,6 +846,7 @@ auto_stop_interval = 60 labels: Some(HashMap::from([("env".into(), "prod".into())])), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() }; @@ -876,6 +881,7 @@ env = "from_task" ])), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() }; @@ -914,6 +920,7 @@ cpu = 2 }), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() }; @@ -951,6 +958,7 @@ auto_stop_interval = 60 }), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() }; @@ -1183,6 +1191,7 @@ network = "block" network: Some(crate::daytona_sandbox::DaytonaNetwork::AllowAll), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() }; @@ -1216,6 +1225,7 @@ auto_stop_interval = 60 ])), ..DaytonaConfig::default() }), + exe: None, }), ..RunDefaults::default() };