From 67652ad4808aa49c58c751351a6fe3a0fe6bc3ae Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 2 Mar 2026 19:58:08 -0500 Subject: [PATCH] Add generic asset snapshot collection for workflow node outputs Automatically discovers and collects well-known output files (test reports, screenshots, trace files) from each node's execution into its log directory. Works across all three sandbox types (local, Docker, Daytona) via a new `download_file_to_local` trait method that handles binary files correctly. - New `asset_snapshot` module with pure functions for find command generation, output parsing, candidate matching, and budget enforcement - `Sandbox::download_file_to_local` implemented for Local (fs::copy), Docker (host bind-mount resolution), and Daytona (SDK download) - `AssetsCaptured` event variant with DEBUG-level tracing - Engine integration: baseline snapshot before handler, collection after (both success and error paths), non-fatal on errors - E2e tests for local sandbox (2 tests), Docker (#[ignore]), Daytona (#[ignore]) Co-Authored-By: Claude Opus 4.6 --- crates/arc-agent/src/docker_sandbox.rs | 36 + crates/arc-agent/src/local_sandbox.rs | 63 ++ crates/arc-agent/src/sandbox.rs | 16 + crates/arc-agent/src/test_support.rs | 43 + crates/arc-workflows/src/artifact.rs | 8 + crates/arc-workflows/src/asset_snapshot.rs | 808 ++++++++++++++++++ crates/arc-workflows/src/cli/mod.rs | 3 + crates/arc-workflows/src/daytona_sandbox.rs | 30 + crates/arc-workflows/src/engine.rs | 50 ++ crates/arc-workflows/src/event.rs | 20 + crates/arc-workflows/src/handler/parallel.rs | 7 + crates/arc-workflows/src/lib.rs | 1 + .../tests/daytona_integration.rs | 116 +++ crates/arc-workflows/tests/integration.rs | 336 ++++++++ 14 files changed, 1537 insertions(+) create mode 100644 crates/arc-workflows/src/asset_snapshot.rs diff --git a/crates/arc-agent/src/docker_sandbox.rs b/crates/arc-agent/src/docker_sandbox.rs index c94b0a995..d0335a0f5 100644 --- a/crates/arc-agent/src/docker_sandbox.rs +++ b/crates/arc-agent/src/docker_sandbox.rs @@ -270,6 +270,42 @@ impl DockerSandbox { #[async_trait] impl Sandbox for DockerSandbox { + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &std::path::Path, + ) -> Result<(), String> { + // Docker bind-mounts host_working_directory -> container_mount_point. + // Resolve the container path to the corresponding host path. + let container_path = self.resolve_container_path(remote_path); + let host_path = if container_path.starts_with(&self.config.container_mount_point) { + let relative = &container_path[self.config.container_mount_point.len()..]; + let relative = relative.strip_prefix('/').unwrap_or(relative); + std::path::PathBuf::from(&self.config.host_working_directory).join(relative) + } else { + return Err(format!( + "Path {container_path} is outside the bind-mounted directory {}", + self.config.container_mount_point + )); + }; + + 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::copy(&host_path, local_path) + .await + .map_err(|e| { + format!( + "Failed to copy {} to {}: {e}", + host_path.display(), + local_path.display() + ) + })?; + Ok(()) + } + async fn initialize(&self) -> Result<(), String> { self.emit(SandboxEvent::Initializing { provider: "docker".into(), diff --git a/crates/arc-agent/src/local_sandbox.rs b/crates/arc-agent/src/local_sandbox.rs index caa47dd4e..36c6e957c 100644 --- a/crates/arc-agent/src/local_sandbox.rs +++ b/crates/arc-agent/src/local_sandbox.rs @@ -345,6 +345,23 @@ impl Sandbox for LocalSandbox { Ok(results) } + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &Path, + ) -> Result<(), String> { + let full_path = self.resolve_path(remote_path); + 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::copy(&full_path, local_path) + .await + .map_err(|e| format!("Failed to copy {} to {}: {e}", full_path.display(), local_path.display()))?; + Ok(()) + } + async fn initialize(&self) -> Result<(), String> { self.emit(SandboxEvent::Initializing { provider: "local".into(), @@ -798,4 +815,50 @@ mod tests { assert_eq!(results.len(), 2); std::fs::remove_dir_all(&dir).unwrap(); } + + #[tokio::test] + async fn local_sandbox_download_file_to_local() { + let dir = temp_dir(); + std::fs::write(dir.join("source.txt"), "hello download").unwrap(); + + let env = LocalSandbox::new(dir.clone()); + let dest = dir.join("output/downloaded.txt"); + env.download_file_to_local("source.txt", &dest) + .await + .unwrap(); + + assert_eq!(std::fs::read_to_string(&dest).unwrap(), "hello download"); + std::fs::remove_dir_all(&dir).unwrap(); + } + + #[tokio::test] + async fn local_sandbox_download_file_to_local_creates_parent_dirs() { + let dir = temp_dir(); + std::fs::write(dir.join("data.bin"), "binary-ish").unwrap(); + + let env = LocalSandbox::new(dir.clone()); + let dest = dir.join("deep/nested/dir/data.bin"); + env.download_file_to_local("data.bin", &dest) + .await + .unwrap(); + + assert_eq!(std::fs::read_to_string(&dest).unwrap(), "binary-ish"); + std::fs::remove_dir_all(&dir).unwrap(); + } + + #[tokio::test] + async fn local_sandbox_download_file_to_local_binary() { + let dir = temp_dir(); + let binary_data: Vec = (0u8..=255).collect(); + std::fs::write(dir.join("binary.bin"), &binary_data).unwrap(); + + let env = LocalSandbox::new(dir.clone()); + let dest = dir.join("out/binary.bin"); + env.download_file_to_local("binary.bin", &dest) + .await + .unwrap(); + + assert_eq!(std::fs::read(&dest).unwrap(), binary_data); + std::fs::remove_dir_all(&dir).unwrap(); + } } diff --git a/crates/arc-agent/src/sandbox.rs b/crates/arc-agent/src/sandbox.rs index 67c0b5adb..7b8bd77e0 100644 --- a/crates/arc-agent/src/sandbox.rs +++ b/crates/arc-agent/src/sandbox.rs @@ -1,6 +1,7 @@ use async_trait::async_trait; use serde::{Deserialize, Serialize}; use std::fmt::Write; +use std::path::Path; use std::sync::Arc; use tokio_util::sync::CancellationToken; @@ -60,6 +61,14 @@ macro_rules! delegate_sandbox { self.$field.glob(pattern, path).await } + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &std::path::Path, + ) -> Result<(), String> { + self.$field.download_file_to_local(remote_path, local_path).await + } + async fn initialize(&self) -> Result<(), String> { self.$field.initialize().await } @@ -290,6 +299,13 @@ pub trait Sandbox: Send + Sync { options: &GrepOptions, ) -> Result, String>; async fn glob(&self, pattern: &str, path: Option<&str>) -> Result, String>; + /// Copy a file from the sandbox to a local filesystem path. + /// Handles binary files correctly across all sandbox types. + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &Path, + ) -> Result<(), String>; async fn initialize(&self) -> Result<(), String>; async fn cleanup(&self) -> Result<(), String>; fn working_directory(&self) -> &str; diff --git a/crates/arc-agent/src/test_support.rs b/crates/arc-agent/src/test_support.rs index 6c24d661c..1a5b03f55 100644 --- a/crates/arc-agent/src/test_support.rs +++ b/crates/arc-agent/src/test_support.rs @@ -162,6 +162,26 @@ impl Sandbox for MockSandbox { Ok(self.glob_results.clone()) } + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &std::path::Path, + ) -> Result<(), String> { + let content = self + .files + .get(remote_path) + .ok_or_else(|| format!("File not found: {remote_path}"))?; + 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, content.as_bytes()) + .await + .map_err(|e| format!("Failed to write {}: {e}", local_path.display()))?; + Ok(()) + } + async fn initialize(&self) -> Result<(), String> { self.emit(crate::sandbox::SandboxEvent::Initializing { provider: "mock".into(), @@ -288,6 +308,29 @@ impl Sandbox for MutableMockSandbox { Ok(vec![]) } + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &std::path::Path, + ) -> Result<(), String> { + let content = self + .files + .lock() + .expect("files lock poisoned") + .get(remote_path) + .cloned() + .ok_or_else(|| format!("File not found: {remote_path}"))?; + 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, content.as_bytes()) + .await + .map_err(|e| format!("Failed to write {}: {e}", local_path.display()))?; + Ok(()) + } + async fn initialize(&self) -> Result<(), String> { Ok(()) } diff --git a/crates/arc-workflows/src/artifact.rs b/crates/arc-workflows/src/artifact.rs index 1ec755c28..21509a87a 100644 --- a/crates/arc-workflows/src/artifact.rs +++ b/crates/arc-workflows/src/artifact.rs @@ -574,6 +574,14 @@ mod tests { Err("not implemented".to_string()) } + async fn download_file_to_local( + &self, + _remote_path: &str, + _local_path: &std::path::Path, + ) -> std::result::Result<(), String> { + Err("not implemented".to_string()) + } + async fn initialize(&self) -> std::result::Result<(), String> { Ok(()) } diff --git a/crates/arc-workflows/src/asset_snapshot.rs b/crates/arc-workflows/src/asset_snapshot.rs new file mode 100644 index 000000000..4107aa8a2 --- /dev/null +++ b/crates/arc-workflows/src/asset_snapshot.rs @@ -0,0 +1,808 @@ +use arc_agent::Sandbox; +use serde::Serialize; +use std::collections::HashMap; +use std::path::Path; +use tracing::{debug, warn}; + +/// Fingerprint of a file for change detection between snapshots. +#[derive(Debug, Clone, PartialEq)] +pub struct FileFingerprint { + pub size: u64, + pub mtime_epoch_secs: f64, +} + +/// A file discovered by the find command. +#[derive(Debug, Clone)] +pub struct DiscoveredFile { + pub relative_path: String, + pub size: u64, + pub mtime_epoch_secs: f64, +} + +/// Summary of an asset collection run. +#[derive(Debug, Clone, Serialize)] +pub struct AssetCollectionSummary { + pub files_copied: usize, + pub total_bytes: u64, + pub files_skipped: usize, + pub download_errors: usize, + pub copied_paths: Vec, +} + +/// Directory path segments that identify asset directories. +const DIRECTORY_SEGMENTS: &[&str] = &[ + "playwright-report", + "test-results", + "cypress/videos", + "cypress/screenshots", +]; + +/// Filename glob patterns for individual asset files. +const FILENAME_GLOBS: &[&str] = &["junit*.xml", "*.trace.zip"]; + +/// Directories to exclude from the find search. +const EXCLUDE_DIRS: &[&str] = &[ + ".git", + "node_modules", + ".pnpm-store", + ".npm", + "target", + ".next", + "__pycache__", +]; + +/// Path segments that indicate excluded tool cache directories. +const EXCLUDE_SEGMENTS: &[&str] = &[ + ".cache/ms-playwright", + "playwright/.cache", + ".yarn/cache", +]; + +/// Maximum size for a single file (10 MB). +const MAX_FILE_SIZE: u64 = 10 * 1024 * 1024; + +/// Maximum total size for all collected files (50 MB). +const MAX_TOTAL_SIZE: u64 = 50 * 1024 * 1024; + +/// Build a platform-aware find command to discover asset files. +pub fn build_find_command(root: &str, platform: &str) -> String { + let mut cmd = format!("find {root}"); + + // Prune excluded directories + let prune_parts: Vec = EXCLUDE_DIRS + .iter() + .map(|d| format!("-name '{d}'")) + .collect(); + cmd.push_str(" \\( "); + cmd.push_str(&prune_parts.join(" -o ")); + cmd.push_str(" \\) -prune -o"); + + // Match conditions: not a symlink, is a file, matches asset patterns + cmd.push_str(" -not -type l -type f \\("); + + let mut path_conditions: Vec = Vec::new(); + for segment in DIRECTORY_SEGMENTS { + path_conditions.push(format!(" -path '*/{segment}/*'")); + } + for glob in FILENAME_GLOBS { + path_conditions.push(format!(" -name '{glob}'")); + } + cmd.push_str(&path_conditions.join(" -o")); + cmd.push_str(" \\)"); + + // Platform-specific output format + match platform { + "darwin" => { + cmd.push_str(" -exec stat -f '%z %m' {} \\; -print"); + } + _ => { + // Linux: use -printf for size, mtime, and relative path + cmd.push_str(&format!(" -printf '%s\\t%T@\\t%P\\n'")); + } + } + + cmd +} + +/// Parse the output of the find command into discovered files. +pub fn parse_find_output(output: &str, platform: &str) -> Vec { + match platform { + "darwin" => parse_find_output_darwin(output), + _ => parse_find_output_linux(output), + } +} + +fn parse_find_output_linux(output: &str) -> Vec { + let mut files = Vec::new(); + for line in output.lines() { + let line = line.trim(); + if line.is_empty() { + continue; + } + let parts: Vec<&str> = line.splitn(3, '\t').collect(); + if parts.len() != 3 { + continue; + } + let size = match parts[0].parse::() { + Ok(s) => s, + Err(_) => continue, + }; + let mtime = match parts[1].parse::() { + Ok(m) => m, + Err(_) => continue, + }; + let path = parts[2].to_string(); + if path.is_empty() { + continue; + } + files.push(DiscoveredFile { + relative_path: path, + size, + mtime_epoch_secs: mtime, + }); + } + files +} + +fn parse_find_output_darwin(output: &str) -> Vec { + let mut files = Vec::new(); + let lines: Vec<&str> = output.lines().collect(); + // Darwin output comes in pairs: "size mtime" then "path" + let mut i = 0; + while i + 1 < lines.len() { + let stat_line = lines[i].trim(); + let path_line = lines[i + 1].trim(); + i += 2; + + if stat_line.is_empty() || path_line.is_empty() { + continue; + } + + let stat_parts: Vec<&str> = stat_line.splitn(2, ' ').collect(); + if stat_parts.len() != 2 { + continue; + } + + let size = match stat_parts[0].parse::() { + Ok(s) => s, + Err(_) => continue, + }; + let mtime = match stat_parts[1].parse::() { + Ok(m) => m, + Err(_) => continue, + }; + + files.push(DiscoveredFile { + relative_path: path_line.to_string(), + size, + mtime_epoch_secs: mtime, + }); + } + files +} + +/// Check whether a path matches known asset patterns. +pub fn is_asset_candidate(path: &str) -> bool { + // Check excluded segments first + for seg in EXCLUDE_SEGMENTS { + if path.contains(seg) { + return false; + } + } + + // Check directory segments — must appear as a complete path segment + for segment in DIRECTORY_SEGMENTS { + // segment may contain a slash (e.g., "cypress/videos"), so check + // that it appears bounded by / or start/end of string + if let Some(pos) = path.find(segment) { + let before_ok = pos == 0 || path.as_bytes()[pos - 1] == b'/'; + let after_pos = pos + segment.len(); + let after_ok = + after_pos >= path.len() || path.as_bytes()[after_pos] == b'/'; + if before_ok && after_ok { + return true; + } + } + } + + // Check filename globs against the last path component + if let Some(filename) = path.rsplit('/').next() { + for glob_pattern in FILENAME_GLOBS { + if matches_simple_glob(glob_pattern, filename) { + return true; + } + } + // Also check at root level (no slash in path) + if !path.contains('/') { + for glob_pattern in FILENAME_GLOBS { + if matches_simple_glob(glob_pattern, path) { + return true; + } + } + } + } + + false +} + +/// Simple glob matching supporting only `*` wildcard. +fn matches_simple_glob(pattern: &str, text: &str) -> bool { + let parts: Vec<&str> = pattern.split('*').collect(); + if parts.len() == 1 { + return pattern == text; + } + + // Check prefix + if !text.starts_with(parts[0]) { + return false; + } + // Check suffix + if !text.ends_with(parts[parts.len() - 1]) { + return false; + } + + // For patterns like "junit*.xml", verify the middle parts appear in order + let mut pos = parts[0].len(); + for part in &parts[1..parts.len() - 1] { + if let Some(found) = text[pos..].find(part) { + pos += found + part.len(); + } else { + return false; + } + } + + true +} + +/// Select which files should be collected based on fingerprint changes, timing, and size budgets. +pub fn select_files_to_collect( + discovered: &[DiscoveredFile], + baseline: &HashMap, + command_start_epoch: f64, +) -> Vec { + let mut candidates: Vec = discovered + .iter() + .filter(|f| { + // Skip files that haven't changed since baseline + if let Some(fp) = baseline.get(&f.relative_path) { + if fp.size == f.size && (fp.mtime_epoch_secs - f.mtime_epoch_secs).abs() < 0.01 { + return false; + } + } + + // Skip files older than command start + if f.mtime_epoch_secs < command_start_epoch { + return false; + } + + // Skip oversized files + if f.size > MAX_FILE_SIZE { + return false; + } + + true + }) + .cloned() + .collect(); + + // Sort by size ascending (smallest first) + candidates.sort_by_key(|f| f.size); + + // Enforce total budget + let mut total: u64 = 0; + let mut selected = Vec::new(); + for f in candidates { + if total + f.size > MAX_TOTAL_SIZE { + break; + } + total += f.size; + selected.push(f); + } + + selected +} + +/// Timeout for the find command (30 seconds). +const FIND_TIMEOUT_MS: u64 = 30_000; + +/// Normalize discovered file paths to be relative to the working directory. +/// On darwin, find outputs absolute paths; on linux, `-printf '%P'` gives relative paths. +/// This strips the working directory prefix and leading `./` to ensure consistent relative paths. +fn normalize_paths(discovered: Vec, root: &str) -> Vec { + let root_with_slash = if root.ends_with('/') { + root.to_string() + } else { + format!("{root}/") + }; + discovered + .into_iter() + .map(|mut f| { + if let Some(stripped) = f.relative_path.strip_prefix(&root_with_slash) { + f.relative_path = stripped.to_string(); + } else if let Some(stripped) = f.relative_path.strip_prefix(root) { + f.relative_path = stripped.strip_prefix('/').unwrap_or(stripped).to_string(); + } + if let Some(stripped) = f.relative_path.strip_prefix("./") { + f.relative_path = stripped.to_string(); + } + f + }) + .filter(|f| !f.relative_path.is_empty()) + .collect() +} + +/// Take a snapshot of current asset files in the sandbox. +/// Returns a fingerprint map of discovered files. +pub async fn snapshot( + sandbox: &dyn Sandbox, +) -> Result, String> { + let root = sandbox.working_directory(); + let platform = sandbox.platform(); + let cmd = build_find_command(root, platform); + + debug!("Taking asset snapshot"); + let result = sandbox + .exec_command(&cmd, FIND_TIMEOUT_MS, None, None, None) + .await?; + + // Ignore non-zero exit codes — find may return 1 if some dirs are unreadable + let discovered = parse_find_output(&result.stdout, platform); + let discovered = normalize_paths(discovered, root); + + let mut fingerprints = HashMap::new(); + for f in discovered { + if is_asset_candidate(&f.relative_path) { + fingerprints.insert( + f.relative_path, + FileFingerprint { + size: f.size, + mtime_epoch_secs: f.mtime_epoch_secs, + }, + ); + } + } + + Ok(fingerprints) +} + +/// Collect asset files that changed since the baseline snapshot. +pub async fn collect_assets( + sandbox: &dyn Sandbox, + stage_dir: &Path, + baseline: &HashMap, + command_start_epoch: f64, +) -> Result { + let root = sandbox.working_directory(); + let platform = sandbox.platform(); + let cmd = build_find_command(root, platform); + + let result = sandbox + .exec_command(&cmd, FIND_TIMEOUT_MS, None, None, None) + .await?; + + let discovered = parse_find_output(&result.stdout, platform); + let discovered = normalize_paths(discovered, root); + let candidates: Vec = discovered + .into_iter() + .filter(|f| is_asset_candidate(&f.relative_path)) + .collect(); + + let total_discovered = candidates.len(); + let to_collect = select_files_to_collect(&candidates, baseline, command_start_epoch); + let files_skipped = total_discovered - to_collect.len(); + + let mut files_copied: usize = 0; + let mut total_bytes: u64 = 0; + let mut download_errors: usize = 0; + let mut copied_paths: Vec = Vec::new(); + + for file in &to_collect { + let dest = stage_dir.join(&file.relative_path); + match sandbox + .download_file_to_local(&file.relative_path, &dest) + .await + { + Ok(()) => { + files_copied += 1; + total_bytes += file.size; + copied_paths.push(file.relative_path.clone()); + } + Err(e) => { + warn!( + path = file.relative_path.as_str(), + error = e.as_str(), + "Asset download failed" + ); + download_errors += 1; + } + } + } + + // Write manifest.json + let summary = AssetCollectionSummary { + files_copied, + total_bytes, + files_skipped, + download_errors, + copied_paths, + }; + + if files_copied > 0 { + if let Ok(json) = serde_json::to_string_pretty(&summary) { + let manifest_path = stage_dir.join("manifest.json"); + if let Some(parent) = manifest_path.parent() { + let _ = tokio::fs::create_dir_all(parent).await; + } + let _ = tokio::fs::write(&manifest_path, json).await; + } + } + + Ok(summary) +} + +#[cfg(test)] +mod tests { + use super::*; + use arc_agent::sandbox::ExecResult; + + /// Minimal mock sandbox for asset_snapshot tests. + struct AssetMockSandbox { + files: HashMap, + exec_result: ExecResult, + working_dir: &'static str, + platform_str: &'static str, + } + + impl AssetMockSandbox { + fn new( + files: HashMap, + exec_stdout: &str, + platform: &'static str, + ) -> Self { + Self { + files, + exec_result: ExecResult { + stdout: exec_stdout.to_string(), + stderr: String::new(), + exit_code: 0, + timed_out: false, + duration_ms: 10, + }, + working_dir: "/home/test", + platform_str: platform, + } + } + } + + #[async_trait::async_trait] + impl Sandbox for AssetMockSandbox { + async fn read_file(&self, _: &str, _: Option, _: Option) -> Result { + Err("not implemented".into()) + } + async fn write_file(&self, _: &str, _: &str) -> Result<(), String> { Ok(()) } + async fn delete_file(&self, _: &str) -> Result<(), String> { Ok(()) } + async fn file_exists(&self, _: &str) -> Result { Ok(false) } + async fn list_directory(&self, _: &str, _: Option) -> Result, String> { Ok(vec![]) } + async fn exec_command(&self, _: &str, _: u64, _: Option<&str>, _: Option<&std::collections::HashMap>, _: Option) -> Result { + Ok(self.exec_result.clone()) + } + async fn grep(&self, _: &str, _: &str, _: &arc_agent::sandbox::GrepOptions) -> Result, String> { Ok(vec![]) } + async fn glob(&self, _: &str, _: Option<&str>) -> Result, String> { Ok(vec![]) } + async fn download_file_to_local(&self, remote_path: &str, local_path: &std::path::Path) -> Result<(), String> { + let content = self.files.get(remote_path) + .ok_or_else(|| format!("File not found: {remote_path}"))?; + if let Some(parent) = local_path.parent() { + tokio::fs::create_dir_all(parent).await + .map_err(|e| format!("Failed to create dirs: {e}"))?; + } + tokio::fs::write(local_path, content.as_bytes()).await + .map_err(|e| format!("Failed to write: {e}"))?; + Ok(()) + } + async fn initialize(&self) -> Result<(), String> { Ok(()) } + async fn cleanup(&self) -> Result<(), String> { Ok(()) } + fn working_directory(&self) -> &str { self.working_dir } + fn platform(&self) -> &str { self.platform_str } + fn os_version(&self) -> String { "Linux 6.1.0".into() } + } + + #[test] + fn is_asset_candidate_matches_directory_segments() { + assert!(is_asset_candidate("playwright-report/index.html")); + } + + #[test] + fn is_asset_candidate_matches_nested_segments() { + assert!(is_asset_candidate("frontend/playwright-report/index.html")); + } + + #[test] + fn is_asset_candidate_rejects_partial_segments() { + assert!(!is_asset_candidate("playwright-reporter/index.html")); + } + + #[test] + fn is_asset_candidate_matches_filename_globs() { + assert!(is_asset_candidate("junit-report.xml")); + assert!(is_asset_candidate("some/dir/junit.xml")); + assert!(is_asset_candidate("output/results.trace.zip")); + } + + #[test] + fn is_asset_candidate_rejects_excluded_paths() { + assert!(!is_asset_candidate(".cache/ms-playwright/chromium/file.txt")); + assert!(!is_asset_candidate("playwright/.cache/some-file")); + assert!(!is_asset_candidate(".yarn/cache/something.zip")); + } + + #[test] + fn is_asset_candidate_matches_cypress_directories() { + assert!(is_asset_candidate("cypress/videos/test.mp4")); + assert!(is_asset_candidate("cypress/screenshots/fail.png")); + } + + #[test] + fn is_asset_candidate_rejects_unrelated_paths() { + assert!(!is_asset_candidate("src/main.rs")); + assert!(!is_asset_candidate("package.json")); + assert!(!is_asset_candidate("report.xml")); + } + + #[test] + fn parse_find_output_linux() { + let output = "1024\t1709312400.0\ttest-results/r.xml\n"; + let files = parse_find_output(output, "linux"); + assert_eq!(files.len(), 1); + assert_eq!(files[0].relative_path, "test-results/r.xml"); + assert_eq!(files[0].size, 1024); + assert!((files[0].mtime_epoch_secs - 1_709_312_400.0).abs() < 0.01); + } + + #[test] + fn parse_find_output_darwin() { + let output = "1024 1709312400\n/tmp/test/test-results/r.xml\n"; + let files = parse_find_output(output, "darwin"); + assert_eq!(files.len(), 1); + assert_eq!(files[0].relative_path, "/tmp/test/test-results/r.xml"); + assert_eq!(files[0].size, 1024); + assert!((files[0].mtime_epoch_secs - 1_709_312_400.0).abs() < 0.01); + } + + #[test] + fn parse_find_output_skips_malformed_lines() { + let output = "not-a-number\t1709312400.0\tfile.xml\n\ + 1024\t1709312400.0\ttest-results/good.xml\n\ + incomplete\n"; + let files = parse_find_output(output, "linux"); + assert_eq!(files.len(), 1); + assert_eq!(files[0].relative_path, "test-results/good.xml"); + } + + #[test] + fn select_files_skips_unchanged() { + let discovered = vec![DiscoveredFile { + relative_path: "test-results/r.xml".to_string(), + size: 1024, + mtime_epoch_secs: 1000.0, + }]; + let mut baseline = HashMap::new(); + baseline.insert( + "test-results/r.xml".to_string(), + FileFingerprint { + size: 1024, + mtime_epoch_secs: 1000.0, + }, + ); + let selected = select_files_to_collect(&discovered, &baseline, 500.0); + assert_eq!(selected.len(), 0); + } + + #[test] + fn select_files_skips_old_mtime() { + let discovered = vec![DiscoveredFile { + relative_path: "test-results/old.xml".to_string(), + size: 1024, + mtime_epoch_secs: 500.0, + }]; + let baseline = HashMap::new(); + let selected = select_files_to_collect(&discovered, &baseline, 1000.0); + assert_eq!(selected.len(), 0); + } + + #[test] + fn select_files_skips_oversized() { + let discovered = vec![DiscoveredFile { + relative_path: "test-results/huge.xml".to_string(), + size: MAX_FILE_SIZE + 1, + mtime_epoch_secs: 2000.0, + }]; + let baseline = HashMap::new(); + let selected = select_files_to_collect(&discovered, &baseline, 1000.0); + assert_eq!(selected.len(), 0); + } + + #[test] + fn select_files_sorts_smallest_first() { + let discovered = vec![ + DiscoveredFile { + relative_path: "a.xml".to_string(), + size: 3000, + mtime_epoch_secs: 2000.0, + }, + DiscoveredFile { + relative_path: "b.xml".to_string(), + size: 1000, + mtime_epoch_secs: 2000.0, + }, + DiscoveredFile { + relative_path: "c.xml".to_string(), + size: 2000, + mtime_epoch_secs: 2000.0, + }, + ]; + let baseline = HashMap::new(); + let selected = select_files_to_collect(&discovered, &baseline, 1000.0); + assert_eq!(selected.len(), 3); + assert_eq!(selected[0].size, 1000); + assert_eq!(selected[1].size, 2000); + assert_eq!(selected[2].size, 3000); + } + + #[test] + fn select_files_enforces_total_budget() { + let discovered: Vec = (0..6) + .map(|i| DiscoveredFile { + relative_path: format!("file{i}.xml"), + size: 9 * 1024 * 1024, // 9 MB each + mtime_epoch_secs: 2000.0, + }) + .collect(); + let baseline = HashMap::new(); + let selected = select_files_to_collect(&discovered, &baseline, 1000.0); + // 50 MB budget / 9 MB each = 5 fit (45 MB), 6th would be 54 MB + assert_eq!(selected.len(), 5); + } + + #[test] + fn build_find_command_linux() { + let cmd = build_find_command("/workspace", "linux"); + assert!(cmd.contains("-printf")); + assert!(cmd.contains("playwright-report")); + assert!(cmd.contains("junit*.xml")); + assert!(cmd.contains("-prune")); + assert!(cmd.contains("node_modules")); + } + + #[test] + fn build_find_command_darwin() { + let cmd = build_find_command("/workspace", "darwin"); + assert!(cmd.contains("-exec stat -f")); + assert!(cmd.contains("playwright-report")); + assert!(cmd.contains("junit*.xml")); + assert!(!cmd.contains("-printf")); + } + + #[test] + fn normalize_paths_strips_root_prefix() { + let files = vec![ + DiscoveredFile { + relative_path: "/workspace/test-results/r.xml".to_string(), + size: 100, + mtime_epoch_secs: 1000.0, + }, + DiscoveredFile { + relative_path: "./test-results/s.xml".to_string(), + size: 200, + mtime_epoch_secs: 1000.0, + }, + DiscoveredFile { + relative_path: "test-results/t.xml".to_string(), + size: 300, + mtime_epoch_secs: 1000.0, + }, + ]; + let normalized = normalize_paths(files, "/workspace"); + assert_eq!(normalized[0].relative_path, "test-results/r.xml"); + assert_eq!(normalized[1].relative_path, "test-results/s.xml"); + assert_eq!(normalized[2].relative_path, "test-results/t.xml"); + } + + #[tokio::test] + async fn snapshot_uses_exec_command_and_parses() { + let mock = AssetMockSandbox::new( + HashMap::new(), + "1024\t2000.0\ttest-results/r.xml\n512\t2000.0\tsrc/main.rs\n", + "linux", + ); + + let fingerprints = snapshot(&mock).await.unwrap(); + // Only test-results/r.xml is an asset candidate, src/main.rs is not + assert_eq!(fingerprints.len(), 1); + assert!(fingerprints.contains_key("test-results/r.xml")); + assert_eq!(fingerprints["test-results/r.xml"].size, 1024); + } + + #[tokio::test] + async fn collect_assets_downloads_and_writes_manifest() { + let stage_dir = tempfile::tempdir().unwrap(); + + let mut files = HashMap::new(); + files.insert("test-results/r.xml".to_string(), "".to_string()); + + let mock = AssetMockSandbox::new( + files, + "1024\t2000.0\ttest-results/r.xml\n", + "linux", + ); + + let baseline = HashMap::new(); + let summary = collect_assets(&mock, stage_dir.path(), &baseline, 1000.0) + .await + .unwrap(); + + assert_eq!(summary.files_copied, 1); + assert_eq!(summary.total_bytes, 1024); + assert_eq!(summary.download_errors, 0); + assert_eq!(summary.copied_paths, vec!["test-results/r.xml"]); + + // Check that the file was written to the stage dir + let dest = stage_dir.path().join("test-results/r.xml"); + assert!(dest.exists()); + let content = std::fs::read_to_string(&dest).unwrap(); + assert_eq!(content, ""); + + // Check manifest + let manifest = stage_dir.path().join("manifest.json"); + assert!(manifest.exists()); + } + + #[tokio::test] + async fn collect_assets_skips_unchanged_files() { + let stage_dir = tempfile::tempdir().unwrap(); + + let mut files = HashMap::new(); + files.insert("test-results/r.xml".to_string(), "".to_string()); + + let mock = AssetMockSandbox::new( + files, + "1024\t2000.0\ttest-results/r.xml\n", + "linux", + ); + + // Provide a baseline with the same fingerprint + let mut baseline = HashMap::new(); + baseline.insert( + "test-results/r.xml".to_string(), + FileFingerprint { + size: 1024, + mtime_epoch_secs: 2000.0, + }, + ); + + let summary = collect_assets(&mock, stage_dir.path(), &baseline, 1000.0) + .await + .unwrap(); + + assert_eq!(summary.files_copied, 0); + } + + #[tokio::test] + async fn collect_assets_non_fatal_on_download_error() { + let stage_dir = tempfile::tempdir().unwrap(); + + // Don't add the file to the mock files map — download will fail + let mock = AssetMockSandbox::new( + HashMap::new(), + "100\t2000.0\ttest-results/missing.xml\n200\t2000.0\ttest-results/also-missing.xml\n", + "linux", + ); + + let baseline = HashMap::new(); + let summary = collect_assets(&mock, stage_dir.path(), &baseline, 1000.0) + .await + .unwrap(); + + assert_eq!(summary.files_copied, 0); + assert_eq!(summary.download_errors, 2); + } +} diff --git a/crates/arc-workflows/src/cli/mod.rs b/crates/arc-workflows/src/cli/mod.rs index 487d48fef..046f029a5 100644 --- a/crates/arc-workflows/src/cli/mod.rs +++ b/crates/arc-workflows/src/cli/mod.rs @@ -571,6 +571,9 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String WorkflowRunEvent::StallWatchdogTimeout { node, idle_seconds } => { format!("[STALL_WATCHDOG_TIMEOUT] node={node} idle_seconds={idle_seconds}") } + WorkflowRunEvent::AssetsCaptured { node_id, files_copied, total_bytes, files_skipped } => { + format!("[ASSETS_CAPTURED] node={node_id} files_copied={files_copied} total_bytes={total_bytes} files_skipped={files_skipped}") + } }; format!("{dim}{body}{reset}", dim = styles.dim, reset = styles.reset) } diff --git a/crates/arc-workflows/src/daytona_sandbox.rs b/crates/arc-workflows/src/daytona_sandbox.rs index eef3683d3..ed72746c0 100644 --- a/crates/arc-workflows/src/daytona_sandbox.rs +++ b/crates/arc-workflows/src/daytona_sandbox.rs @@ -263,6 +263,36 @@ pub fn get_gh_token() -> Result { #[async_trait] impl Sandbox for DaytonaSandbox { + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &Path, + ) -> Result<(), String> { + let sandbox = self.sandbox()?; + let resolved = self.resolve_path(remote_path); + + let fs_svc = sandbox + .fs() + .await + .map_err(|e| format!("Failed to get fs service: {e}"))?; + + let bytes = fs_svc + .download_file(&resolved) + .await + .map_err(|e| format!("Failed to download file {resolved}: {e}"))?; + + 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(()) + } + async fn initialize(&self) -> Result<(), String> { self.emit(SandboxEvent::Initializing { provider: "daytona".into(), diff --git a/crates/arc-workflows/src/engine.rs b/crates/arc-workflows/src/engine.rs index ddbe782bd..2603258eb 100644 --- a/crates/arc-workflows/src/engine.rs +++ b/crates/arc-workflows/src/engine.rs @@ -14,6 +14,7 @@ use tokio_util::sync::CancellationToken; use arc_git_storage::trailerlink::{self, Trailer}; use crate::artifact::{offload_large_values, sync_artifacts_to_env, ArtifactStore}; +use crate::asset_snapshot; use crate::checkpoint::Checkpoint; use crate::condition::evaluate_condition; use crate::context::Context; @@ -844,6 +845,21 @@ impl WorkflowRunEngine { let node_timeout = node.timeout(); for attempt in 1..=policy.max_attempts { + // Take baseline asset snapshot before handler execution + let baseline = match asset_snapshot::snapshot(self.services.sandbox.as_ref()).await { + Ok(fp) => fp, + Err(e) => { + tracing::warn!(node = %node.id, error = %e, "Asset baseline snapshot failed"); + std::collections::HashMap::new() + } + }; + // Floor to integer seconds: macOS stat reports mtime as integer seconds, + // so a fractional epoch would reject files created in the same second. + let command_start_epoch = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs() as f64) + .unwrap_or(0.0); + // Gap #11: Panic safety -- catch panics from handler execution let result = { let future = handler.execute(node, context, graph, logs_root, &self.services); @@ -878,6 +894,40 @@ impl WorkflowRunEngine { } }; + // Collect assets after handler completes (both success and error) + { + let assets_dir = node_dir(logs_root, &node.id, visit) + .join("assets") + .join(format!("attempt_{attempt}")); + match asset_snapshot::collect_assets( + self.services.sandbox.as_ref(), + &assets_dir, + &baseline, + command_start_epoch, + ) + .await + { + Ok(summary) if summary.files_copied > 0 => { + self.services + .emitter + .emit(&WorkflowRunEvent::AssetsCaptured { + node_id: node.id.clone(), + files_copied: summary.files_copied, + total_bytes: summary.total_bytes, + files_skipped: summary.files_skipped, + }); + } + Ok(_) => {} + Err(e) => { + tracing::warn!( + node = %node.id, + error = %e, + "Asset collection failed" + ); + } + } + } + let outcome = match result { Ok(o) => o, Err(e) => { diff --git a/crates/arc-workflows/src/event.rs b/crates/arc-workflows/src/event.rs index 923858191..f7bf5eb81 100644 --- a/crates/arc-workflows/src/event.rs +++ b/crates/arc-workflows/src/event.rs @@ -181,6 +181,12 @@ pub enum WorkflowRunEvent { node: String, idle_seconds: u64, }, + AssetsCaptured { + node_id: String, + files_copied: usize, + total_bytes: u64, + files_skipped: usize, + }, } impl WorkflowRunEvent { @@ -447,6 +453,20 @@ impl WorkflowRunEvent { } => { warn!(node, idle_seconds, "Stall watchdog timeout"); } + Self::AssetsCaptured { + node_id, + files_copied, + total_bytes, + files_skipped, + } => { + debug!( + node_id, + files_copied, + total_bytes, + files_skipped, + "Assets captured" + ); + } } } } diff --git a/crates/arc-workflows/src/handler/parallel.rs b/crates/arc-workflows/src/handler/parallel.rs index c7a61f07e..d9c595b81 100644 --- a/crates/arc-workflows/src/handler/parallel.rs +++ b/crates/arc-workflows/src/handler/parallel.rs @@ -77,6 +77,13 @@ impl Sandbox for WorktreeSandbox { async fn glob(&self, pattern: &str, path: Option<&str>) -> Result, String> { self.inner.glob(pattern, path).await } + async fn download_file_to_local( + &self, + remote_path: &str, + local_path: &std::path::Path, + ) -> Result<(), String> { + self.inner.download_file_to_local(remote_path, local_path).await + } async fn initialize(&self) -> Result<(), String> { self.inner.initialize().await } diff --git a/crates/arc-workflows/src/lib.rs b/crates/arc-workflows/src/lib.rs index 1adabd793..45763f1b8 100644 --- a/crates/arc-workflows/src/lib.rs +++ b/crates/arc-workflows/src/lib.rs @@ -1,4 +1,5 @@ pub mod artifact; +pub mod asset_snapshot; pub mod checkpoint; pub mod cli; pub mod condition; diff --git a/crates/arc-workflows/tests/daytona_integration.rs b/crates/arc-workflows/tests/daytona_integration.rs index 79e4dc6ca..57f51f3fd 100644 --- a/crates/arc-workflows/tests/daytona_integration.rs +++ b/crates/arc-workflows/tests/daytona_integration.rs @@ -1025,3 +1025,119 @@ async fn daytona_git_checkpoint_with_shadow_branch() { env.cleanup().await.unwrap(); } + +// --------------------------------------------------------------------------- +// Asset collection e2e — Daytona sandbox +// --------------------------------------------------------------------------- + +/// Handler that creates asset files via exec_command on the sandbox. +struct AssetCreatorHandler; + +#[async_trait::async_trait] +impl Handler for AssetCreatorHandler { + async fn execute( + &self, + _node: &Node, + _context: &Context, + _graph: &Graph, + _logs_root: &Path, + services: &arc_workflows::handler::EngineServices, + ) -> Result { + let script = concat!( + "mkdir -p test-results && ", + "echo '' > test-results/report.xml && ", + "echo 'test output' > test-results/output.txt" + ); + services + .sandbox + .exec_command(script, 30_000, None, None, None) + .await + .map_err(|e| ArcError::Handler(format!("exec failed: {e}")))?; + Ok(Outcome::success()) + } +} + +/// Daytona sandbox: asset collection discovers files on the remote sandbox and +/// downloads them to the local logs directory. +#[tokio::test] +#[ignore] +async fn daytona_asset_collection() { + let env = create_env().await; + env.initialize().await.unwrap(); + let env: Arc = Arc::new(env); + + let dir = tempfile::tempdir().unwrap(); + + let mut registry = HandlerRegistry::new(Box::new(AssetCreatorHandler)); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + + let engine = WorkflowRunEngine::new(registry, Arc::new(EventEmitter::new()), env.clone()); + + let mut graph = Graph::new("DaytonaAssetTest"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test asset collection on Daytona".to_string()), + ); + + let mut start = Node::new("start"); + start + .attrs + .insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut create_assets = Node::new("create_assets"); + create_assets + .attrs + .insert("label".to_string(), AttrValue::String("Create Assets".to_string())); + graph + .nodes + .insert("create_assets".to_string(), create_assets); + + let mut exit = Node::new("exit"); + exit.attrs + .insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + + graph.edges.push(Edge::new("start", "create_assets")); + graph.edges.push(Edge::new("create_assets", "exit")); + + let config = RunConfig { + logs_root: dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: "asset-test-daytona".into(), + git_checkpoint: None, + base_sha: None, + run_branch: None, + meta_branch: None, + labels: std::collections::HashMap::new(), + }; + + let outcome = engine + .run(&graph, &config) + .await + .expect("pipeline should succeed"); + assert_eq!(outcome.status, StageStatus::Success); + + let assets_dir = dir + .path() + .join("nodes") + .join("create_assets") + .join("assets") + .join("attempt_1"); + + let report_path = assets_dir.join("test-results/report.xml"); + assert!( + report_path.exists(), + "report.xml should be collected from Daytona sandbox at {}", + report_path.display() + ); + let content = std::fs::read_to_string(&report_path).unwrap(); + assert!(content.contains("testsuites")); + + let manifest_path = assets_dir.join("manifest.json"); + assert!(manifest_path.exists(), "manifest.json should exist"); + + env.cleanup().await.unwrap(); +} diff --git a/crates/arc-workflows/tests/integration.rs b/crates/arc-workflows/tests/integration.rs index b771699dd..699038d38 100644 --- a/crates/arc-workflows/tests/integration.rs +++ b/crates/arc-workflows/tests/integration.rs @@ -7698,6 +7698,10 @@ impl arc_agent::Sandbox for RemoteMockEnv { Ok(()) } + async fn download_file_to_local(&self, _: &str, _: &std::path::Path) -> std::result::Result<(), String> { + Err("not implemented".to_string()) + } + fn working_directory(&self) -> &str { &self.working_dir } @@ -8061,6 +8065,10 @@ impl arc_agent::Sandbox for CliTestEnv { Ok(()) } + async fn download_file_to_local(&self, _: &str, _: &std::path::Path) -> Result<(), String> { + Err("not implemented".to_string()) + } + fn working_directory(&self) -> &str { "/tmp/test" } @@ -8299,6 +8307,9 @@ async fn cli_backend_run_fails_on_nonzero_exit() { async fn cleanup(&self) -> Result<(), String> { Ok(()) } + async fn download_file_to_local(&self, _: &str, _: &std::path::Path) -> Result<(), String> { + Err("not implemented".to_string()) + } fn working_directory(&self) -> &str { "/tmp" } @@ -11153,3 +11164,328 @@ async fn e2e_stall_watchdog_with_explicit_timeout_override() { } // Daytona parallel git branching test is in daytona_integration.rs + +// --------------------------------------------------------------------------- +// Asset collection e2e tests +// --------------------------------------------------------------------------- + +/// Handler that creates asset files in the sandbox working directory via exec_command. +struct AssetCreatorHandler { + should_fail: bool, +} + +impl AssetCreatorHandler { + fn success() -> Self { + Self { should_fail: false } + } + + fn failing() -> Self { + Self { should_fail: true } + } +} + +#[async_trait::async_trait] +impl Handler for AssetCreatorHandler { + async fn execute( + &self, + _node: &Node, + _context: &Context, + _graph: &Graph, + _logs_root: &Path, + services: &arc_workflows::handler::EngineServices, + ) -> Result { + // Create asset files via the sandbox's exec_command + let script = concat!( + "mkdir -p test-results && ", + "echo '' > test-results/report.xml && ", + "echo 'test output' > test-results/output.txt" + ); + services + .sandbox + .exec_command(script, 30_000, None, None, None) + .await + .map_err(|e| ArcError::Handler(format!("exec failed: {e}")))?; + + if self.should_fail { + Ok(Outcome::fail("intentional failure")) + } else { + Ok(Outcome::success()) + } + } +} + +/// Local sandbox: asset collection discovers and downloads files created by a handler. +#[tokio::test] +async fn asset_collection_local_sandbox_success() { + let work_dir = tempfile::tempdir().unwrap(); + let logs_dir = tempfile::tempdir().unwrap(); + + let sandbox: Arc = + Arc::new(arc_agent::LocalSandbox::new(work_dir.path().to_path_buf())); + sandbox.initialize().await.unwrap(); + + let mut registry = HandlerRegistry::new(Box::new(AssetCreatorHandler::success())); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + + let mut emitter = EventEmitter::new(); + let events = collect_events(&mut emitter); + + let engine = WorkflowRunEngine::new(registry, Arc::new(emitter), sandbox.clone()); + + let mut graph = Graph::new("AssetCollectionTest"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test asset collection".to_string()), + ); + + let mut start = Node::new("start"); + start + .attrs + .insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut create_assets = Node::new("create_assets"); + create_assets + .attrs + .insert("label".to_string(), AttrValue::String("Create Assets".to_string())); + graph + .nodes + .insert("create_assets".to_string(), create_assets); + + let mut exit = Node::new("exit"); + exit.attrs + .insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + + graph.edges.push(Edge::new("start", "create_assets")); + graph.edges.push(Edge::new("create_assets", "exit")); + + let config = RunConfig { + logs_root: logs_dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: "asset-test-local".into(), + git_checkpoint: None, + base_sha: None, + run_branch: None, + meta_branch: None, + labels: std::collections::HashMap::new(), + }; + + let outcome = engine.run(&graph, &config).await.expect("run should succeed"); + assert_eq!(outcome.status, StageStatus::Success); + + // Check that asset files were collected into the stage directory + let assets_dir = logs_dir + .path() + .join("nodes") + .join("create_assets") + .join("assets") + .join("attempt_1"); + + let report_path = assets_dir.join("test-results/report.xml"); + assert!( + report_path.exists(), + "report.xml should be collected at {}", + report_path.display() + ); + let report_content = std::fs::read_to_string(&report_path).unwrap(); + assert!(report_content.contains("testsuites")); + + // Check manifest.json was written + let manifest_path = assets_dir.join("manifest.json"); + assert!( + manifest_path.exists(), + "manifest.json should exist at {}", + manifest_path.display() + ); + let manifest: serde_json::Value = + serde_json::from_str(&std::fs::read_to_string(&manifest_path).unwrap()).unwrap(); + assert!(manifest["files_copied"].as_u64().unwrap() >= 1); + + // Check that AssetsCaptured event was emitted + let captured_events = events.lock().unwrap(); + let assets_events: Vec<&WorkflowRunEvent> = captured_events + .iter() + .filter(|e| matches!(e, WorkflowRunEvent::AssetsCaptured { .. })) + .collect(); + assert!( + !assets_events.is_empty(), + "should emit at least one AssetsCaptured event" + ); +} + +/// Local sandbox: assets are still collected even when the handler fails. +#[tokio::test] +async fn asset_collection_local_sandbox_on_failure() { + let work_dir = tempfile::tempdir().unwrap(); + let logs_dir = tempfile::tempdir().unwrap(); + + let sandbox: Arc = + Arc::new(arc_agent::LocalSandbox::new(work_dir.path().to_path_buf())); + sandbox.initialize().await.unwrap(); + + let mut registry = HandlerRegistry::new(Box::new(AssetCreatorHandler::failing())); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + + let engine = WorkflowRunEngine::new( + registry, + Arc::new(EventEmitter::new()), + sandbox.clone(), + ); + + let mut graph = Graph::new("AssetCollectionFailTest"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test asset collection on failure".to_string()), + ); + + let mut start = Node::new("start"); + start + .attrs + .insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut create_assets = Node::new("create_assets"); + create_assets + .attrs + .insert("label".to_string(), AttrValue::String("Create Assets".to_string())); + graph + .nodes + .insert("create_assets".to_string(), create_assets); + + let mut exit = Node::new("exit"); + exit.attrs + .insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + + graph.edges.push(Edge::new("start", "create_assets")); + graph.edges.push(Edge::new("create_assets", "exit")); + + let config = RunConfig { + logs_root: logs_dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: "asset-test-fail".into(), + git_checkpoint: None, + base_sha: None, + run_branch: None, + meta_branch: None, + labels: std::collections::HashMap::new(), + }; + + let outcome = engine.run(&graph, &config).await.expect("run should succeed"); + // The pipeline completes (handler returned Fail, not an error), but assets should still be collected + assert_eq!(outcome.status, StageStatus::Fail); + + let assets_dir = logs_dir + .path() + .join("nodes") + .join("create_assets") + .join("assets") + .join("attempt_1"); + + let report_path = assets_dir.join("test-results/report.xml"); + assert!( + report_path.exists(), + "report.xml should still be collected after handler failure, at {}", + report_path.display() + ); +} + +/// Docker sandbox: asset collection works across the bind-mount boundary. +/// Requires Docker with `arc-agent:latest` image available locally. +#[tokio::test] +#[ignore] +async fn asset_collection_docker_sandbox() { + let host_dir = tempfile::tempdir().unwrap(); + let logs_dir = tempfile::tempdir().unwrap(); + + let config = arc_agent::DockerSandboxConfig { + host_working_directory: host_dir.path().to_str().unwrap().to_string(), + auto_pull: false, + ..Default::default() + }; + let sandbox: Arc = + Arc::new(arc_agent::DockerSandbox::new(config).expect("Docker not available")); + sandbox.initialize().await.expect("Docker init failed"); + + let mut registry = HandlerRegistry::new(Box::new(AssetCreatorHandler::success())); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + + let engine = WorkflowRunEngine::new( + registry, + Arc::new(EventEmitter::new()), + sandbox.clone(), + ); + + let mut graph = Graph::new("DockerAssetTest"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test asset collection in Docker".to_string()), + ); + + let mut start = Node::new("start"); + start + .attrs + .insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut create_assets = Node::new("create_assets"); + create_assets + .attrs + .insert("label".to_string(), AttrValue::String("Create Assets".to_string())); + graph + .nodes + .insert("create_assets".to_string(), create_assets); + + let mut exit = Node::new("exit"); + exit.attrs + .insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + + graph.edges.push(Edge::new("start", "create_assets")); + graph.edges.push(Edge::new("create_assets", "exit")); + + let run_config = RunConfig { + logs_root: logs_dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: "asset-test-docker".into(), + git_checkpoint: None, + base_sha: None, + run_branch: None, + meta_branch: None, + labels: std::collections::HashMap::new(), + }; + + let outcome = engine + .run(&graph, &run_config) + .await + .expect("pipeline should succeed"); + assert_eq!(outcome.status, StageStatus::Success); + + let assets_dir = logs_dir + .path() + .join("nodes") + .join("create_assets") + .join("assets") + .join("attempt_1"); + + let report_path = assets_dir.join("test-results/report.xml"); + assert!( + report_path.exists(), + "report.xml should be collected from Docker container at {}", + report_path.display() + ); + let content = std::fs::read_to_string(&report_path).unwrap(); + assert!(content.contains("testsuites")); + + let manifest_path = assets_dir.join("manifest.json"); + assert!(manifest_path.exists(), "manifest.json should exist"); + + sandbox.cleanup().await.unwrap(); +}