mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
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 <noreply@anthropic.com>
This commit is contained in:
parent
743a4fc677
commit
67652ad480
14 changed files with 1537 additions and 0 deletions
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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<u8> = (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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<Vec<String>, String>;
|
||||
async fn glob(&self, pattern: &str, path: Option<&str>) -> Result<Vec<String>, 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;
|
||||
|
|
|
|||
|
|
@ -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(())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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(())
|
||||
}
|
||||
|
|
|
|||
808
crates/arc-workflows/src/asset_snapshot.rs
Normal file
808
crates/arc-workflows/src/asset_snapshot.rs
Normal file
|
|
@ -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<String>,
|
||||
}
|
||||
|
||||
/// 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<String> = 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<String> = 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<DiscoveredFile> {
|
||||
match platform {
|
||||
"darwin" => parse_find_output_darwin(output),
|
||||
_ => parse_find_output_linux(output),
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_find_output_linux(output: &str) -> Vec<DiscoveredFile> {
|
||||
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::<u64>() {
|
||||
Ok(s) => s,
|
||||
Err(_) => continue,
|
||||
};
|
||||
let mtime = match parts[1].parse::<f64>() {
|
||||
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<DiscoveredFile> {
|
||||
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::<u64>() {
|
||||
Ok(s) => s,
|
||||
Err(_) => continue,
|
||||
};
|
||||
let mtime = match stat_parts[1].parse::<f64>() {
|
||||
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<String, FileFingerprint>,
|
||||
command_start_epoch: f64,
|
||||
) -> Vec<DiscoveredFile> {
|
||||
let mut candidates: Vec<DiscoveredFile> = 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<DiscoveredFile>, root: &str) -> Vec<DiscoveredFile> {
|
||||
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<HashMap<String, FileFingerprint>, 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<String, FileFingerprint>,
|
||||
command_start_epoch: f64,
|
||||
) -> Result<AssetCollectionSummary, String> {
|
||||
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<DiscoveredFile> = 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<String> = 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<String, String>,
|
||||
exec_result: ExecResult,
|
||||
working_dir: &'static str,
|
||||
platform_str: &'static str,
|
||||
}
|
||||
|
||||
impl AssetMockSandbox {
|
||||
fn new(
|
||||
files: HashMap<String, String>,
|
||||
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<usize>, _: Option<usize>) -> Result<String, String> {
|
||||
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<bool, String> { Ok(false) }
|
||||
async fn list_directory(&self, _: &str, _: Option<usize>) -> Result<Vec<arc_agent::sandbox::DirEntry>, String> { Ok(vec![]) }
|
||||
async fn exec_command(&self, _: &str, _: u64, _: Option<&str>, _: Option<&std::collections::HashMap<String, String>>, _: Option<tokio_util::sync::CancellationToken>) -> Result<ExecResult, String> {
|
||||
Ok(self.exec_result.clone())
|
||||
}
|
||||
async fn grep(&self, _: &str, _: &str, _: &arc_agent::sandbox::GrepOptions) -> Result<Vec<String>, String> { Ok(vec![]) }
|
||||
async fn glob(&self, _: &str, _: Option<&str>) -> Result<Vec<String>, 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<DiscoveredFile> = (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(), "<test/>".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, "<test/>");
|
||||
|
||||
// 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(), "<test/>".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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -263,6 +263,36 @@ pub fn get_gh_token() -> Result<String, String> {
|
|||
|
||||
#[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(),
|
||||
|
|
|
|||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -77,6 +77,13 @@ impl Sandbox for WorktreeSandbox {
|
|||
async fn glob(&self, pattern: &str, path: Option<&str>) -> Result<Vec<String>, 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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
pub mod artifact;
|
||||
pub mod asset_snapshot;
|
||||
pub mod checkpoint;
|
||||
pub mod cli;
|
||||
pub mod condition;
|
||||
|
|
|
|||
|
|
@ -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<Outcome, ArcError> {
|
||||
let script = concat!(
|
||||
"mkdir -p test-results && ",
|
||||
"echo '<testsuites><testsuite name=\"example\"/></testsuites>' > 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<dyn Sandbox> = 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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<Outcome, ArcError> {
|
||||
// Create asset files via the sandbox's exec_command
|
||||
let script = concat!(
|
||||
"mkdir -p test-results && ",
|
||||
"echo '<testsuites><testsuite name=\"example\"/></testsuites>' > 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<dyn arc_agent::Sandbox> =
|
||||
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<dyn arc_agent::Sandbox> =
|
||||
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<dyn arc_agent::Sandbox> =
|
||||
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();
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue