mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
fix(glob): harden artifact traversal
This commit is contained in:
parent
5980627bc8
commit
dec67ec92e
17 changed files with 443 additions and 486 deletions
|
|
@ -119,7 +119,7 @@ Finds files matching a glob pattern.
|
|||
|
||||
Returns matching file paths, one per line, sorted lexicographically by their path relative to the search root.
|
||||
|
||||
Patterns are case-sensitive and relative to `path`: `*` and `?` stay within one path segment, bracket expressions such as `[abc]` match one character, and `**` crosses directories when used as a complete segment. Leading dots are matched normally. For example, `*.rs` searches only the root of `path`, and `**/*.rs` searches recursively. Patterns must be relative and cannot contain a `..` segment.
|
||||
Patterns are case-sensitive and relative to `path`: `*` and `?` stay within one path segment, bracket expressions such as `[abc]` match one character, and `**` crosses directories when used as a complete segment. Leading dots are matched normally. For example, `*.rs` searches only the root of `path`, and `**/*.rs` searches recursively. Patterns must use `/`, be relative, and cannot contain a backslash or `..` segment.
|
||||
|
||||
<Note>
|
||||
Unlike `grep`, `glob` does **not** mark files as read for the read-before-write guardrail. To modify a file found via glob, the agent must read it first.
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ date: "2026-07-25"
|
|||
|
||||
- Replace `*.trace.zip` with `**/*.trace.zip` when traces may be nested.
|
||||
- Replace `test-results/**` with `**/test-results/**` only when `test-results` may appear below any directory instead of at the workspace root.
|
||||
- Absolute patterns and patterns containing a `..` segment are now rejected.
|
||||
- Absolute patterns, backslashes, and patterns containing a `..` segment are now rejected.
|
||||
- Agent `glob` results are now sorted lexicographically; local results are no longer ordered by modification time.
|
||||
</Warning>
|
||||
|
||||
|
|
|
|||
|
|
@ -431,7 +431,7 @@ Artifact globs use `/` as the separator and have the same semantics in every san
|
|||
- Bracket expressions such as `[abc]` and `[!abc]` match one character.
|
||||
- `**` matches across directories when used as a complete segment.
|
||||
- Leading dots are matched normally.
|
||||
- Patterns are relative to the sandbox working directory and case-sensitive.
|
||||
- Patterns use `/`, are relative to the sandbox working directory, and are case-sensitive. Backslashes are invalid.
|
||||
- Absolute patterns and patterns containing a `..` segment are invalid.
|
||||
|
||||
For example, `.ai/reports/*.md` matches direct Markdown children of `.ai/reports`, while `.ai/reports/**/*.md` also matches nested reports. `*.trace.zip` matches only the working-directory root; use `**/*.trace.zip` to match at any depth. To collect date-named implementation plans, use `.ai/plans/????-??-??-*.md`.
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ use fabro_types::{
|
|||
SystemActorKind, WorkflowSettings, parse_blob_ref,
|
||||
};
|
||||
use fabro_util::version::FABRO_VERSION;
|
||||
use fabro_util::workspace_glob::{WorkspaceGlob, WorkspaceGlobError};
|
||||
use fabro_workflow::command_log::{command_log_path, read_json_string_blob, read_log_slice};
|
||||
use fabro_workflow::run_status::RunStatus;
|
||||
use fabro_workflow::{Error as WorkflowError, operations};
|
||||
|
|
@ -902,13 +903,31 @@ async fn snapshot_run_variables(
|
|||
state.stores.variables.value_map().await
|
||||
}
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
enum RunVariableSubstitutionError {
|
||||
#[error(transparent)]
|
||||
Interpolation(#[from] ResolveError),
|
||||
|
||||
#[error("run.artifacts.include[{index}]: {source}")]
|
||||
ArtifactGlob {
|
||||
index: usize,
|
||||
#[source]
|
||||
source: WorkspaceGlobError,
|
||||
},
|
||||
}
|
||||
|
||||
fn substitute_run_variables(
|
||||
variables: &HashMap<String, String>,
|
||||
settings: &mut WorkflowSettings,
|
||||
) -> Result<(), ResolveError> {
|
||||
) -> Result<(), RunVariableSubstitutionError> {
|
||||
settings
|
||||
.run
|
||||
.substitute_variables(|name| variables.get(name).cloned())
|
||||
.substitute_variables(|name| variables.get(name).cloned())?;
|
||||
for (index, pattern) in settings.run.artifacts.include.iter().enumerate() {
|
||||
WorkspaceGlob::try_new(pattern)
|
||||
.map_err(|source| RunVariableSubstitutionError::ArtifactGlob { index, source })?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn get_run_status(
|
||||
|
|
@ -1213,3 +1232,26 @@ fn build_command_log_response(
|
|||
})
|
||||
.into_response()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn rejects_artifact_glob_that_becomes_unsafe_after_interpolation() {
|
||||
let variables = HashMap::from([("PATTERN".to_string(), "../outside/**".to_string())]);
|
||||
let mut settings = WorkflowSettings::default();
|
||||
settings.run.artifacts.include = vec!["{{ vars.PATTERN }}".to_string()];
|
||||
|
||||
let error = substitute_run_variables(&variables, &mut settings)
|
||||
.expect_err("interpolated parent traversal should be rejected");
|
||||
|
||||
assert!(matches!(
|
||||
error,
|
||||
RunVariableSubstitutionError::ArtifactGlob {
|
||||
index: 0,
|
||||
source: WorkspaceGlobError::ParentTraversal { .. },
|
||||
}
|
||||
));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -31,21 +31,23 @@ exact match and unique unless replace_all is true. If the edit fails with 'old_s
|
|||
re-read the file and take the exact text from the fresh output rather than guessing again. \
|
||||
Preserve existing indentation.";
|
||||
|
||||
const GLOB_DESCRIPTION: &str = "Find files by name using a glob pattern, most recently modified \
|
||||
first.
|
||||
const GLOB_DESCRIPTION: &str = "Find files by search-root-relative path using a glob pattern. \
|
||||
Results are sorted lexicographically by relative path.
|
||||
|
||||
Use this instead of `find` or recursive `ls` through Bash. Prefer patterns with a literal anchor \
|
||||
— an extension or a subdirectory — over bare wildcards.
|
||||
|
||||
Good patterns:
|
||||
- `*.rs` — an extension at any depth below the search root
|
||||
- `*.rs` — direct children of the search root
|
||||
- `**/*.rs` — files at any depth below the search root
|
||||
- `src/*.rs` — directly inside `src/`, not recursive
|
||||
- `src/**/*.rs` — recursive walk under a subdirectory
|
||||
- `{src,tests}/**/*.rs` — brace expansion works
|
||||
- `src/[lm]ib.rs` — a bracket expression matches one character
|
||||
|
||||
Avoid recursing into dependency or build output (`node_modules/**`, `target/**`): those produce \
|
||||
thousands of matches and waste context. Narrow to a specific subpath instead. Results are files, \
|
||||
so to locate a directory, glob for something inside it.";
|
||||
so to locate a directory, glob for something inside it. Patterns must use `/`, be relative, and \
|
||||
cannot contain a `..` segment.";
|
||||
|
||||
pub struct KimiProfile {
|
||||
base: BaseProfile,
|
||||
|
|
@ -342,7 +344,8 @@ mod tests {
|
|||
// Grep must not promise ripgrep syntax: fabro falls back to POSIX grep.
|
||||
let grep = describe("Grep");
|
||||
assert!(grep.contains("POSIX"), "{grep}");
|
||||
assert!(describe("Glob").contains("most recently modified"));
|
||||
assert!(describe("Glob").contains("sorted lexicographically"));
|
||||
assert!(describe("Glob").contains("`*.rs` — direct children"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
|
|
@ -1,3 +1,5 @@
|
|||
use crate::sandbox;
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub(crate) enum CloneDecision {
|
||||
EmptyWorkspace {
|
||||
|
|
@ -32,9 +34,9 @@ pub(crate) fn github_repo_layout(
|
|||
})?;
|
||||
let workspace_root = trim_root(workspace_root);
|
||||
let repos_root = trim_root(repos_root);
|
||||
let repos_owner_path = join_remote_path(repos_root, &owner);
|
||||
let primary_repo_path = join_remote_path(&repos_owner_path, &repo);
|
||||
let primary_repo_link = join_remote_path(workspace_root, &repo);
|
||||
let repos_owner_path = sandbox::join_sandbox_path(repos_root, &owner);
|
||||
let primary_repo_path = sandbox::join_sandbox_path(&repos_owner_path, &repo);
|
||||
let primary_repo_link = sandbox::join_sandbox_path(workspace_root, &repo);
|
||||
|
||||
Ok(GitHubRepoLayout {
|
||||
owner,
|
||||
|
|
@ -51,14 +53,6 @@ fn trim_root(root: &str) -> &str {
|
|||
if trimmed.is_empty() { "/" } else { trimmed }
|
||||
}
|
||||
|
||||
fn join_remote_path(root: &str, name: &str) -> String {
|
||||
if root == "/" {
|
||||
format!("/{name}")
|
||||
} else {
|
||||
format!("{root}/{name}")
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub(crate) enum EmptyWorkspaceReason {
|
||||
SkipClone,
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
use std::collections::{HashMap, HashSet};
|
||||
use std::collections::HashMap;
|
||||
use std::fmt::Write;
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
|
|
@ -27,8 +27,8 @@ use tokio_util::sync::CancellationToken;
|
|||
use crate::clone_source::{self, CloneDecision, EmptyWorkspaceReason};
|
||||
use crate::redact::redact_auth_url;
|
||||
use crate::sandbox::{
|
||||
BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH, RefreshOutcome,
|
||||
join_sandbox_path, optional_timeout, resolve_path, validate_bash_probe,
|
||||
self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH,
|
||||
REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, optional_timeout, resolve_path, validate_bash_probe,
|
||||
};
|
||||
use crate::{
|
||||
CommandOutputCallback, DirEntry, ExecResult, ExecStreamingResult, GrepOptions, Sandbox,
|
||||
|
|
@ -561,84 +561,6 @@ impl DaytonaSandbox {
|
|||
})
|
||||
}
|
||||
|
||||
async fn walk_files_recursive(
|
||||
&self,
|
||||
base: &str,
|
||||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
if relative_start.split('/').any(|segment| {
|
||||
options
|
||||
.excluded_directory_names
|
||||
.iter()
|
||||
.any(|name| name == segment)
|
||||
}) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let sandbox = self.sandbox()?;
|
||||
let fs_svc = sandbox
|
||||
.fs()
|
||||
.await
|
||||
.map_err(|e| crate::Error::context("Failed to get Daytona fs service", e))?;
|
||||
let Some(root) = resolve_daytona_walk_root(&fs_svc, base, relative_start).await? else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let mut files = Vec::new();
|
||||
let mut stack = vec![(root, relative_start.to_string())];
|
||||
let mut visited_dirs = HashSet::new();
|
||||
|
||||
while let Some((dir, relative_dir)) = stack.pop() {
|
||||
if !visited_dirs.insert(dir.clone()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
let entries = match fs_svc.list_files(&dir).await {
|
||||
Ok(entries) => entries,
|
||||
Err(daytona_sdk::DaytonaError::NotFound { .. }) => continue,
|
||||
Err(err) => {
|
||||
return Err(crate::Error::context(
|
||||
format!("Failed to list Daytona directory {dir}"),
|
||||
err,
|
||||
));
|
||||
}
|
||||
};
|
||||
|
||||
for entry in entries {
|
||||
if entry.name.is_empty() || entry.name == "." || entry.name == ".." {
|
||||
continue;
|
||||
}
|
||||
|
||||
let child_path = join_sandbox_path(&dir, &entry.name);
|
||||
let relative_path = join_sandbox_path(&relative_dir, &entry.name);
|
||||
if entry.is_dir {
|
||||
if !options
|
||||
.excluded_directory_names
|
||||
.iter()
|
||||
.any(|name| name == &entry.name)
|
||||
{
|
||||
stack.push((child_path, relative_path));
|
||||
}
|
||||
} else {
|
||||
let size = u64::try_from(entry.size).map_err(|error| {
|
||||
crate::Error::context(
|
||||
format!("Invalid Daytona file size for {child_path}"),
|
||||
error,
|
||||
)
|
||||
})?;
|
||||
files.push(SandboxFile {
|
||||
path: child_path,
|
||||
relative_path,
|
||||
size,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
files.sort_by(|left, right| left.relative_path.cmp(&right.relative_path));
|
||||
Ok(files)
|
||||
}
|
||||
|
||||
/// Read-only access to the SDK sandbox once initialized. Returns `None`
|
||||
/// before `initialize()` or `reconnect()` has populated the cell.
|
||||
pub fn sandbox_handle(&self) -> Option<&daytona_sdk::Sandbox> {
|
||||
|
|
@ -2002,41 +1924,21 @@ impl Sandbox for DaytonaSandbox {
|
|||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
let base = self.resolve_path(base);
|
||||
self.walk_files_recursive(&base, relative_start, options)
|
||||
.await
|
||||
}
|
||||
}
|
||||
|
||||
async fn resolve_daytona_walk_root(
|
||||
fs_svc: &daytona_sdk::FileSystemService,
|
||||
base: &str,
|
||||
relative_start: &str,
|
||||
) -> crate::Result<Option<String>> {
|
||||
let mut root = base.to_string();
|
||||
for segment in relative_start
|
||||
.split('/')
|
||||
.filter(|segment| !segment.is_empty())
|
||||
{
|
||||
let entries = match fs_svc.list_files(&root).await {
|
||||
Ok(entries) => entries,
|
||||
Err(daytona_sdk::DaytonaError::NotFound { .. }) => return Ok(None),
|
||||
Err(error) => {
|
||||
return Err(crate::Error::context(
|
||||
format!("Failed to inspect Daytona traversal root {root}"),
|
||||
error,
|
||||
));
|
||||
}
|
||||
};
|
||||
let Some(entry) = entries.into_iter().find(|entry| entry.name == segment) else {
|
||||
return Ok(None);
|
||||
};
|
||||
if !entry.is_dir {
|
||||
return Ok(None);
|
||||
if options.excludes_relative_path(relative_start) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
root = join_sandbox_path(&root, segment);
|
||||
|
||||
let base = self.resolve_path(base);
|
||||
let command = sandbox::build_remote_walk_command(&base, relative_start, options);
|
||||
let result = self
|
||||
.exec_command(&command, REMOTE_WALK_TIMEOUT_MS, None, None, None)
|
||||
.await?;
|
||||
if !result.is_success() {
|
||||
return Err(crate::Error::exec("recursive file traversal", result));
|
||||
}
|
||||
|
||||
sandbox::parse_remote_walk_output(&base, relative_start, &result.stdout)
|
||||
}
|
||||
Ok(Some(root))
|
||||
}
|
||||
|
||||
fn daytona_callback_error(err: &crate::Error) -> DaytonaError {
|
||||
|
|
|
|||
|
|
@ -29,8 +29,9 @@ use crate::clone_source::{self, CloneDecision, EmptyWorkspaceReason};
|
|||
use crate::managed_labels::{self, MANAGED_LABEL, RUN_ID_LABEL};
|
||||
use crate::redact::redact_auth_url;
|
||||
use crate::sandbox::{
|
||||
BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH, RefreshOutcome,
|
||||
StdioProcessControl, join_sandbox_path, optional_timeout, resolve_path, validate_bash_probe,
|
||||
self, BASH_ENV_VAR, BASH_PROBE_SCRIPT, BASH_PROBE_TIMEOUT_MS, REMOTE_BASH,
|
||||
REMOTE_WALK_TIMEOUT_MS, RefreshOutcome, StdioProcessControl, optional_timeout, resolve_path,
|
||||
validate_bash_probe,
|
||||
};
|
||||
use crate::{
|
||||
CommandOutputCallback, DEFAULT_EXEC_OUTPUT_TAIL_BYTES, DirEntry, ExecResult,
|
||||
|
|
@ -1887,25 +1888,20 @@ impl Sandbox for DockerSandbox {
|
|||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
if relative_start.split('/').any(|segment| {
|
||||
options
|
||||
.excluded_directory_names
|
||||
.iter()
|
||||
.any(|name| name == segment)
|
||||
}) {
|
||||
if options.excludes_relative_path(relative_start) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
let base = self.resolve_container_path(base);
|
||||
let command = docker_walk_command(&base, relative_start, options);
|
||||
let command = sandbox::build_remote_walk_command(&base, relative_start, options);
|
||||
let result = self
|
||||
.docker_exec_shell(&command, 30_000, None, None, None)
|
||||
.docker_exec_shell(&command, REMOTE_WALK_TIMEOUT_MS, None, None, None)
|
||||
.await?;
|
||||
if !result.is_success() {
|
||||
return Err(crate::Error::exec("recursive file traversal", result));
|
||||
}
|
||||
|
||||
parse_docker_walk_output(&base, relative_start, &result.stdout)
|
||||
sandbox::parse_remote_walk_output(&base, relative_start, &result.stdout)
|
||||
}
|
||||
|
||||
fn working_directory(&self) -> &str {
|
||||
|
|
@ -2017,72 +2013,6 @@ impl Sandbox for DockerSandbox {
|
|||
}
|
||||
}
|
||||
|
||||
fn docker_walk_command(base: &str, relative_start: &str, options: &WalkOptions) -> String {
|
||||
let traversal_root = join_sandbox_path(base, relative_start);
|
||||
let quoted_root = shell_quote(&traversal_root);
|
||||
let mut command = format!("if [ -e {quoted_root} ]");
|
||||
let mut component_path = base.to_string();
|
||||
for segment in relative_start
|
||||
.split('/')
|
||||
.filter(|segment| !segment.is_empty())
|
||||
{
|
||||
component_path = join_sandbox_path(&component_path, segment);
|
||||
let _ = write!(command, " && [ ! -L {} ]", shell_quote(&component_path));
|
||||
}
|
||||
let _ = write!(command, "; then find -H {quoted_root}");
|
||||
|
||||
if !options.excluded_directory_names.is_empty() {
|
||||
command.push_str(" \\( -type d \\(");
|
||||
for (index, directory_name) in options.excluded_directory_names.iter().enumerate() {
|
||||
if index > 0 {
|
||||
command.push_str(" -o");
|
||||
}
|
||||
let _ = write!(command, " -name {}", shell_quote(directory_name));
|
||||
}
|
||||
command.push_str(" \\) -prune \\) -o");
|
||||
}
|
||||
|
||||
command.push_str(" -not -type l -type f -printf '%s\\0%P\\0'; fi");
|
||||
command
|
||||
}
|
||||
|
||||
fn parse_docker_walk_output(
|
||||
base: &str,
|
||||
relative_start: &str,
|
||||
output: &str,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
let mut fields = output.split('\0');
|
||||
let mut files = Vec::new();
|
||||
|
||||
while let Some(size) = fields.next() {
|
||||
if size.is_empty() {
|
||||
break;
|
||||
}
|
||||
let relative_to_start = fields.next().ok_or_else(|| {
|
||||
crate::Error::message("Malformed Docker file traversal output: missing path")
|
||||
})?;
|
||||
let size = size.parse::<u64>().map_err(|error| {
|
||||
crate::Error::context(
|
||||
format!("Malformed Docker file traversal size {size:?}"),
|
||||
error,
|
||||
)
|
||||
})?;
|
||||
let relative_path = if relative_to_start.is_empty() {
|
||||
relative_start.to_string()
|
||||
} else {
|
||||
join_sandbox_path(relative_start, relative_to_start)
|
||||
};
|
||||
files.push(SandboxFile {
|
||||
path: join_sandbox_path(base, &relative_path),
|
||||
relative_path,
|
||||
size,
|
||||
});
|
||||
}
|
||||
|
||||
files.sort_by(|left, right| left.relative_path.cmp(&right.relative_path));
|
||||
Ok(files)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
#[expect(
|
||||
|
|
@ -2100,8 +2030,8 @@ mod tests {
|
|||
use crate::sandbox::{BASH_PROBE_MARKER, bash_probe_passed};
|
||||
|
||||
#[test]
|
||||
fn docker_walk_command_only_uses_find_for_traversal() {
|
||||
let command = docker_walk_command("/workspace", ".ai", &WalkOptions {
|
||||
fn remote_walk_command_only_uses_find_for_traversal() {
|
||||
let command = sandbox::build_remote_walk_command("/workspace", ".ai", &WalkOptions {
|
||||
excluded_directory_names: vec!["target".to_string(), "node_modules".to_string()],
|
||||
});
|
||||
|
||||
|
|
@ -2114,9 +2044,10 @@ mod tests {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn docker_walk_output_is_relative_to_the_declared_base() {
|
||||
fn remote_walk_output_is_relative_to_the_declared_base() {
|
||||
let files =
|
||||
parse_docker_walk_output("/workspace", ".ai/reports", "12\0result.md\0").unwrap();
|
||||
sandbox::parse_remote_walk_output("/workspace", ".ai/reports", "12\0result.md\0")
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(files, vec![SandboxFile {
|
||||
path: "/workspace/.ai/reports/result.md".to_string(),
|
||||
|
|
|
|||
|
|
@ -978,12 +978,7 @@ async fn walk_local_files(
|
|||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
if relative_start.split('/').any(|segment| {
|
||||
options
|
||||
.excluded_directory_names
|
||||
.iter()
|
||||
.any(|name| name == segment)
|
||||
}) {
|
||||
if options.excludes_relative_path(relative_start) {
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
|
|
@ -1029,12 +1024,10 @@ async fn walk_local_files(
|
|||
});
|
||||
} else if file_type.is_dir() {
|
||||
if path != base
|
||||
&& path.file_name().is_some_and(|file_name| {
|
||||
options
|
||||
.excluded_directory_names
|
||||
.iter()
|
||||
.any(|name| file_name == name.as_str())
|
||||
})
|
||||
&& path
|
||||
.file_name()
|
||||
.and_then(std::ffi::OsStr::to_str)
|
||||
.is_some_and(|file_name| options.excludes_name(file_name))
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
|
@ -1059,7 +1052,6 @@ async fn walk_local_files(
|
|||
}
|
||||
}
|
||||
|
||||
files.sort_by(|left, right| left.relative_path.cmp(&right.relative_path));
|
||||
Ok(files)
|
||||
}
|
||||
|
||||
|
|
@ -1897,8 +1889,7 @@ mod tests {
|
|||
std::fs::write(dir.join("src/nested/readme.md"), "").unwrap();
|
||||
|
||||
let env = LocalSandbox::new(dir.clone());
|
||||
let mut results = env.glob("**/*.rs", None).await.unwrap();
|
||||
results.sort();
|
||||
let results = env.glob("**/*.rs", None).await.unwrap();
|
||||
|
||||
assert_eq!(results, vec![
|
||||
dir.join("a.rs").to_string_lossy().into_owned(),
|
||||
|
|
@ -1986,13 +1977,15 @@ mod tests {
|
|||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
files
|
||||
.iter()
|
||||
.map(|file| (file.relative_path.as_str(), file.size))
|
||||
.collect::<Vec<_>>(),
|
||||
vec![(".ai/reports/empty.md", 0), (".ai/reports/result.md", 6),]
|
||||
);
|
||||
let mut file_metadata = files
|
||||
.iter()
|
||||
.map(|file| (file.relative_path.as_str(), file.size))
|
||||
.collect::<Vec<_>>();
|
||||
file_metadata.sort_unstable();
|
||||
assert_eq!(file_metadata, vec![
|
||||
(".ai/reports/empty.md", 0),
|
||||
(".ai/reports/result.md", 6),
|
||||
]);
|
||||
let root = dir.to_string_lossy();
|
||||
assert!(
|
||||
files
|
||||
|
|
|
|||
|
|
@ -29,6 +29,10 @@ pub(crate) const BASH_PROBE_TIMEOUT_MS: u64 = 10_000;
|
|||
#[cfg(any(feature = "docker", feature = "daytona"))]
|
||||
pub(crate) const REMOTE_BASH: &str = "/bin/bash";
|
||||
|
||||
/// Timeout for provider-neutral remote file traversal.
|
||||
#[cfg(any(feature = "docker", feature = "daytona"))]
|
||||
pub(crate) const REMOTE_WALK_TIMEOUT_MS: u64 = 30_000;
|
||||
|
||||
/// Environment variable Bash consults for non-interactive startup source.
|
||||
///
|
||||
/// Sandbox providers must remove or blank this before invoking `bash -c`;
|
||||
|
|
@ -925,6 +929,22 @@ pub struct WalkOptions {
|
|||
pub excluded_directory_names: Vec<String>,
|
||||
}
|
||||
|
||||
impl WalkOptions {
|
||||
#[must_use]
|
||||
pub fn excludes_name(&self, name: &str) -> bool {
|
||||
self.excluded_directory_names
|
||||
.iter()
|
||||
.any(|excluded| excluded == name)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn excludes_relative_path(&self, relative_path: &str) -> bool {
|
||||
relative_path
|
||||
.split('/')
|
||||
.any(|segment| self.excludes_name(segment))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct GrepOptions {
|
||||
pub glob_filter: Option<String>,
|
||||
|
|
@ -1206,7 +1226,7 @@ pub(crate) fn resolve_path(path: &str, working_dir: &str) -> String {
|
|||
if std::path::Path::new(path).is_absolute() {
|
||||
path.to_string()
|
||||
} else {
|
||||
format!("{working_dir}/{path}")
|
||||
join_sandbox_path(working_dir, path)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -1224,6 +1244,77 @@ pub(crate) fn join_sandbox_path(base: &str, relative_path: &str) -> String {
|
|||
format!("{}/{relative_path}", base.trim_end_matches('/'))
|
||||
}
|
||||
|
||||
#[cfg(any(feature = "docker", feature = "daytona"))]
|
||||
pub(crate) fn build_remote_walk_command(
|
||||
base: &str,
|
||||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> String {
|
||||
let traversal_root = join_sandbox_path(base, relative_start);
|
||||
let quoted_root = shell_quote(&traversal_root);
|
||||
let mut command = format!("if [ -e {quoted_root} ]");
|
||||
let mut component_path = base.to_string();
|
||||
for segment in relative_start
|
||||
.split('/')
|
||||
.filter(|segment| !segment.is_empty())
|
||||
{
|
||||
component_path = join_sandbox_path(&component_path, segment);
|
||||
let _ = write!(command, " && [ ! -L {} ]", shell_quote(&component_path));
|
||||
}
|
||||
let _ = write!(command, "; then find -H {quoted_root}");
|
||||
|
||||
if !options.excluded_directory_names.is_empty() {
|
||||
command.push_str(" \\( -type d \\(");
|
||||
for (index, directory_name) in options.excluded_directory_names.iter().enumerate() {
|
||||
if index > 0 {
|
||||
command.push_str(" -o");
|
||||
}
|
||||
let _ = write!(command, " -name {}", shell_quote(directory_name));
|
||||
}
|
||||
command.push_str(" \\) -prune \\) -o");
|
||||
}
|
||||
|
||||
command.push_str(" -not -type l -type f -printf '%s\\0%P\\0'; fi");
|
||||
command
|
||||
}
|
||||
|
||||
#[cfg(any(feature = "docker", feature = "daytona"))]
|
||||
pub(crate) fn parse_remote_walk_output(
|
||||
base: &str,
|
||||
relative_start: &str,
|
||||
output: &str,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
let mut fields = output.split('\0');
|
||||
let mut files = Vec::new();
|
||||
|
||||
while let Some(size) = fields.next() {
|
||||
if size.is_empty() {
|
||||
break;
|
||||
}
|
||||
let relative_to_start = fields.next().ok_or_else(|| {
|
||||
crate::Error::message("Malformed recursive file traversal output: missing path")
|
||||
})?;
|
||||
let size = size.parse::<u64>().map_err(|error| {
|
||||
crate::Error::context(
|
||||
format!("Malformed recursive file traversal size {size:?}"),
|
||||
error,
|
||||
)
|
||||
})?;
|
||||
let relative_path = if relative_to_start.is_empty() {
|
||||
relative_start.to_string()
|
||||
} else {
|
||||
join_sandbox_path(relative_start, relative_to_start)
|
||||
};
|
||||
files.push(SandboxFile {
|
||||
path: join_sandbox_path(base, &relative_path),
|
||||
relative_path,
|
||||
size,
|
||||
});
|
||||
}
|
||||
|
||||
Ok(files)
|
||||
}
|
||||
|
||||
/// Shell-quote a string using `shlex::try_quote`, with a fallback for edge
|
||||
/// cases. Re-exported from [`fabro_util::shell::shell_quote`] so sandbox code
|
||||
/// and the config resolve layer share one audited implementation.
|
||||
|
|
|
|||
|
|
@ -12,8 +12,8 @@ use tokio_util::sync::CancellationToken;
|
|||
use crate::sandbox::{self, StdioProcessControl};
|
||||
use crate::{
|
||||
DEFAULT_EXEC_OUTPUT_TAIL_BYTES, DirEntry, ExecResult, GrepOptions, Sandbox, SandboxEvent,
|
||||
SandboxEventCallback, StderrCollector, StdioProcess, StdioProcessHandle,
|
||||
StdioProcessTermination,
|
||||
SandboxEventCallback, SandboxFile, StderrCollector, StdioProcess, StdioProcessHandle,
|
||||
StdioProcessTermination, WalkOptions,
|
||||
};
|
||||
|
||||
// --- MockSandbox ---
|
||||
|
|
@ -47,6 +47,10 @@ pub struct MockSandbox {
|
|||
/// Fails `exec_command` and `exec_command_streaming` before any process
|
||||
/// runs, so callers see a transport error rather than an `ExecResult`.
|
||||
pub exec_error: Option<String>,
|
||||
/// Files returned by `walk_files`, before traversal-root and exclusion
|
||||
/// filtering.
|
||||
pub walk_files: Vec<SandboxFile>,
|
||||
pub walk_files_error: Option<String>,
|
||||
/// Reported by `exec_command_streaming`. Set to `false` to model a
|
||||
/// provider that cannot separate stdout from stderr.
|
||||
pub streams_separated: bool,
|
||||
|
|
@ -83,6 +87,18 @@ impl MockSandbox {
|
|||
.lock()
|
||||
.expect("stdio_process lock poisoned") = Some(process);
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_walk_files(mut self, files: Vec<SandboxFile>) -> Self {
|
||||
self.walk_files = files;
|
||||
self
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn with_walk_files_error(mut self, error: impl Into<String>) -> Self {
|
||||
self.walk_files_error = Some(error.into());
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
impl MockSandbox {
|
||||
|
|
@ -123,6 +139,8 @@ impl Default for MockSandbox {
|
|||
stdio_process_error: None,
|
||||
stdio_process: Mutex::new(None),
|
||||
exec_error: None,
|
||||
walk_files: Vec::new(),
|
||||
walk_files_error: None,
|
||||
streams_separated: true,
|
||||
}
|
||||
}
|
||||
|
|
@ -343,6 +361,38 @@ impl Sandbox for MockSandbox {
|
|||
Ok(self.glob_results.clone())
|
||||
}
|
||||
|
||||
async fn walk_files(
|
||||
&self,
|
||||
_base: &str,
|
||||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> crate::Result<Vec<SandboxFile>> {
|
||||
if let Some(error) = &self.walk_files_error {
|
||||
return Err(crate::Error::message(error.clone()));
|
||||
}
|
||||
|
||||
Ok(self
|
||||
.walk_files
|
||||
.iter()
|
||||
.filter(|file| {
|
||||
relative_start.is_empty()
|
||||
|| file.relative_path == relative_start
|
||||
|| file
|
||||
.relative_path
|
||||
.strip_prefix(relative_start)
|
||||
.is_some_and(|suffix| suffix.starts_with('/'))
|
||||
})
|
||||
.filter(|file| {
|
||||
let parent = file
|
||||
.relative_path
|
||||
.rsplit_once('/')
|
||||
.map_or("", |(parent, _)| parent);
|
||||
!options.excludes_relative_path(parent)
|
||||
})
|
||||
.cloned()
|
||||
.collect())
|
||||
}
|
||||
|
||||
async fn download_file_to_local(
|
||||
&self,
|
||||
remote_path: &str,
|
||||
|
|
|
|||
|
|
@ -1,12 +1,13 @@
|
|||
use std::collections::BTreeMap;
|
||||
use std::path::Path;
|
||||
|
||||
use fabro_agent::Sandbox;
|
||||
use fabro_sandbox::{SandboxFile, WalkOptions};
|
||||
use fabro_types::ArtifactUpload;
|
||||
use fabro_util::workspace_glob::WorkspaceGlobSet;
|
||||
use futures::{StreamExt as _, TryStreamExt as _, stream};
|
||||
use sha2::{Digest, Sha256};
|
||||
use tokio::fs;
|
||||
use tokio::io::AsyncReadExt as _;
|
||||
use tracing::warn;
|
||||
|
||||
/// Summary of an artifact collection run.
|
||||
|
|
@ -47,12 +48,15 @@ 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;
|
||||
|
||||
/// Independent traversal roots may run concurrently, but remote providers
|
||||
/// should not receive an unbounded burst of file-walk operations.
|
||||
const MAX_CONCURRENT_ARTIFACT_WALKS: usize = 4;
|
||||
|
||||
/// Select which files should be collected based on size budgets.
|
||||
pub fn select_files_to_collect(discovered: &[SandboxFile]) -> Vec<SandboxFile> {
|
||||
pub fn select_files_to_collect(discovered: Vec<SandboxFile>) -> Vec<SandboxFile> {
|
||||
let mut candidates: Vec<SandboxFile> = discovered
|
||||
.iter()
|
||||
.into_iter()
|
||||
.filter(|file| file.size <= MAX_FILE_SIZE)
|
||||
.cloned()
|
||||
.collect();
|
||||
|
||||
candidates.sort_by(|left, right| {
|
||||
|
|
@ -77,23 +81,31 @@ pub fn select_files_to_collect(discovered: &[SandboxFile]) -> Vec<SandboxFile> {
|
|||
async fn compute_artifact_info(
|
||||
relative_path: &str,
|
||||
local_path: &Path,
|
||||
) -> std::result::Result<ArtifactUpload, String> {
|
||||
) -> std::result::Result<Option<ArtifactUpload>, String> {
|
||||
let mime = mime_guess::from_path(relative_path)
|
||||
.first_or_octet_stream()
|
||||
.to_string();
|
||||
let data = fs::read(local_path)
|
||||
let file = fs::File::open(local_path)
|
||||
.await
|
||||
.map_err(|error| format!("failed to open {}: {error}", local_path.display()))?;
|
||||
let mut data = Vec::new();
|
||||
file.take(MAX_FILE_SIZE + 1)
|
||||
.read_to_end(&mut data)
|
||||
.await
|
||||
.map_err(|error| format!("failed to read {}: {error}", local_path.display()))?;
|
||||
let bytes = u64::try_from(data.len()).unwrap_or(u64::MAX);
|
||||
if bytes > MAX_FILE_SIZE {
|
||||
return Ok(None);
|
||||
}
|
||||
let content_md5 = format!("{:x}", md5::compute(&data));
|
||||
let content_sha256 = hex::encode(Sha256::digest(&data));
|
||||
Ok(ArtifactUpload {
|
||||
Ok(Some(ArtifactUpload {
|
||||
path: relative_path.to_string(),
|
||||
mime,
|
||||
content_md5,
|
||||
content_sha256,
|
||||
bytes,
|
||||
})
|
||||
}))
|
||||
}
|
||||
|
||||
/// Collect artifact files matching the configured workspace globs.
|
||||
|
|
@ -108,49 +120,56 @@ pub async fn collect_artifacts(
|
|||
.map(|directory| (*directory).to_string())
|
||||
.collect(),
|
||||
};
|
||||
let mut discovered_by_path = BTreeMap::new();
|
||||
for traversal_root in globs.traversal_roots() {
|
||||
let files = sandbox
|
||||
.walk_files(sandbox.working_directory(), traversal_root, &walk_options)
|
||||
.await
|
||||
.map_err(|error| {
|
||||
format!(
|
||||
"artifact file traversal failed below {traversal_root:?}: {}",
|
||||
error.display_with_causes()
|
||||
)
|
||||
})?;
|
||||
for file in files {
|
||||
if globs.is_match(&file.relative_path) {
|
||||
discovered_by_path
|
||||
.entry(file.relative_path.clone())
|
||||
.or_insert(file);
|
||||
}
|
||||
}
|
||||
}
|
||||
let discovered = discovered_by_path.into_values().collect::<Vec<_>>();
|
||||
let walk_options = &walk_options;
|
||||
let traversal_roots = globs
|
||||
.traversal_roots()
|
||||
.into_iter()
|
||||
.map(str::to_string)
|
||||
.collect::<Vec<_>>();
|
||||
let walks = stream::iter(traversal_roots)
|
||||
.map(|traversal_root| async move {
|
||||
sandbox
|
||||
.walk_files(sandbox.working_directory(), &traversal_root, walk_options)
|
||||
.await
|
||||
.map_err(|error| {
|
||||
format!(
|
||||
"artifact file traversal failed below {traversal_root:?}: {}",
|
||||
error.display_with_causes()
|
||||
)
|
||||
})
|
||||
})
|
||||
.buffer_unordered(MAX_CONCURRENT_ARTIFACT_WALKS)
|
||||
.try_collect::<Vec<_>>()
|
||||
.await?;
|
||||
let discovered = walks
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.filter(|file| globs.is_match(&file.relative_path))
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let total_discovered = discovered.len();
|
||||
let to_collect = select_files_to_collect(&discovered);
|
||||
let files_skipped = total_discovered - to_collect.len();
|
||||
let to_collect = select_files_to_collect(discovered);
|
||||
let mut files_skipped = total_discovered - to_collect.len();
|
||||
|
||||
let mut files_copied = 0;
|
||||
let mut total_bytes = 0;
|
||||
let mut total_bytes: u64 = 0;
|
||||
let mut download_errors = 0;
|
||||
let mut hash_errors = 0;
|
||||
let mut captured_assets = Vec::new();
|
||||
|
||||
for file in &to_collect {
|
||||
let dest = artifact_capture_dir.join(&file.relative_path);
|
||||
match sandbox
|
||||
.download_file_to_local(&file.relative_path, &dest)
|
||||
.await
|
||||
{
|
||||
match sandbox.download_file_to_local(&file.path, &dest).await {
|
||||
Ok(()) => match compute_artifact_info(&file.relative_path, &dest).await {
|
||||
Ok(info) => {
|
||||
Ok(Some(info)) if total_bytes.saturating_add(info.bytes) <= MAX_TOTAL_SIZE => {
|
||||
files_copied += 1;
|
||||
total_bytes += info.bytes;
|
||||
captured_assets.push(info);
|
||||
}
|
||||
Ok(Some(_) | None) => {
|
||||
let _ = fs::remove_file(&dest).await;
|
||||
files_skipped += 1;
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(
|
||||
path = file.relative_path.as_str(),
|
||||
|
|
@ -188,172 +207,10 @@ pub async fn collect_artifacts(
|
|||
mod tests {
|
||||
use std::collections::{BTreeSet, HashMap};
|
||||
|
||||
use fabro_agent::sandbox::{DirEntry, ExecResult, GrepOptions};
|
||||
use fabro_types::CommandTermination;
|
||||
use fabro_sandbox::test_support::MockSandbox;
|
||||
|
||||
use super::*;
|
||||
|
||||
struct AssetMockSandbox {
|
||||
contents: HashMap<String, String>,
|
||||
discovered: Vec<SandboxFile>,
|
||||
walk_error: Option<String>,
|
||||
}
|
||||
|
||||
impl AssetMockSandbox {
|
||||
fn new(contents: HashMap<String, String>) -> Self {
|
||||
let discovered = contents
|
||||
.iter()
|
||||
.map(|(relative_path, content)| sandbox_file(relative_path, content.len() as u64))
|
||||
.collect();
|
||||
Self {
|
||||
contents,
|
||||
discovered,
|
||||
walk_error: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn with_discovered(mut self, discovered: Vec<SandboxFile>) -> Self {
|
||||
self.discovered = discovered;
|
||||
self
|
||||
}
|
||||
|
||||
fn with_walk_error(mut self, error: &str) -> Self {
|
||||
self.walk_error = Some(error.to_string());
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Sandbox for AssetMockSandbox {
|
||||
async fn read_file_bytes(&self, _: &str) -> fabro_sandbox::Result<Vec<u8>> {
|
||||
Err("not implemented".into())
|
||||
}
|
||||
|
||||
async fn write_file(&self, _: &str, _: &str) -> fabro_sandbox::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn delete_file(&self, _: &str) -> fabro_sandbox::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn file_exists(&self, _: &str) -> fabro_sandbox::Result<bool> {
|
||||
Ok(false)
|
||||
}
|
||||
|
||||
async fn list_directory(
|
||||
&self,
|
||||
_: &str,
|
||||
_: Option<usize>,
|
||||
) -> fabro_sandbox::Result<Vec<DirEntry>> {
|
||||
Ok(Vec::new())
|
||||
}
|
||||
|
||||
async fn exec_command(
|
||||
&self,
|
||||
_: &str,
|
||||
_: u64,
|
||||
_: Option<&str>,
|
||||
_: Option<&std::collections::HashMap<String, String>>,
|
||||
_: Option<tokio_util::sync::CancellationToken>,
|
||||
) -> fabro_sandbox::Result<ExecResult> {
|
||||
Ok(ExecResult {
|
||||
stdout: String::new(),
|
||||
stderr: String::new(),
|
||||
exit_code: Some(0),
|
||||
termination: CommandTermination::Exited,
|
||||
duration_ms: 0,
|
||||
})
|
||||
}
|
||||
|
||||
async fn grep(
|
||||
&self,
|
||||
_: &str,
|
||||
_: &str,
|
||||
_: &GrepOptions,
|
||||
) -> fabro_sandbox::Result<Vec<String>> {
|
||||
Ok(Vec::new())
|
||||
}
|
||||
|
||||
async fn walk_files(
|
||||
&self,
|
||||
_: &str,
|
||||
relative_start: &str,
|
||||
options: &WalkOptions,
|
||||
) -> fabro_sandbox::Result<Vec<SandboxFile>> {
|
||||
if let Some(error) = &self.walk_error {
|
||||
return Err(fabro_sandbox::Error::message(error.clone()));
|
||||
}
|
||||
|
||||
Ok(self
|
||||
.discovered
|
||||
.iter()
|
||||
.filter(|file| {
|
||||
relative_start.is_empty()
|
||||
|| file
|
||||
.relative_path
|
||||
.starts_with(&format!("{relative_start}/"))
|
||||
})
|
||||
.filter(|file| {
|
||||
let parent = file
|
||||
.relative_path
|
||||
.rsplit_once('/')
|
||||
.map_or("", |(parent, _)| parent);
|
||||
!parent.split('/').any(|segment| {
|
||||
options
|
||||
.excluded_directory_names
|
||||
.iter()
|
||||
.any(|name| name == segment)
|
||||
})
|
||||
})
|
||||
.cloned()
|
||||
.collect())
|
||||
}
|
||||
|
||||
async fn download_file_to_local(
|
||||
&self,
|
||||
remote_path: &str,
|
||||
local_path: &Path,
|
||||
) -> fabro_sandbox::Result<()> {
|
||||
let content = self.contents.get(remote_path).ok_or_else(|| {
|
||||
fabro_sandbox::Error::message(format!("File not found: {remote_path}"))
|
||||
})?;
|
||||
if let Some(parent) = local_path.parent() {
|
||||
fs::create_dir_all(parent).await.map_err(|error| {
|
||||
fabro_sandbox::Error::context("Failed to create dirs", error)
|
||||
})?;
|
||||
}
|
||||
fs::write(local_path, content.as_bytes())
|
||||
.await
|
||||
.map_err(|error| fabro_sandbox::Error::context("Failed to write", error))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn upload_file_from_local(&self, _: &Path, _: &str) -> fabro_sandbox::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn initialize(&self) -> fabro_sandbox::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn cleanup(&self) -> fabro_sandbox::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn working_directory(&self) -> &str {
|
||||
"/home/test"
|
||||
}
|
||||
|
||||
fn platform(&self) -> &str {
|
||||
"linux"
|
||||
}
|
||||
|
||||
fn os_version(&self) -> String {
|
||||
"Linux 6.1.0".to_string()
|
||||
}
|
||||
}
|
||||
|
||||
fn sandbox_file(relative_path: &str, size: u64) -> SandboxFile {
|
||||
SandboxFile {
|
||||
path: format!("/home/test/{relative_path}"),
|
||||
|
|
@ -362,13 +219,29 @@ mod tests {
|
|||
}
|
||||
}
|
||||
|
||||
fn asset_sandbox(contents: HashMap<String, String>) -> MockSandbox {
|
||||
let mut files = HashMap::new();
|
||||
let mut discovered = Vec::new();
|
||||
for (relative_path, content) in contents {
|
||||
let file = sandbox_file(&relative_path, content.len() as u64);
|
||||
files.insert(file.path.clone(), content);
|
||||
discovered.push(file);
|
||||
}
|
||||
|
||||
MockSandbox {
|
||||
files,
|
||||
..MockSandbox::linux()
|
||||
}
|
||||
.with_walk_files(discovered)
|
||||
}
|
||||
|
||||
fn workspace_globs(patterns: &[&str]) -> WorkspaceGlobSet {
|
||||
WorkspaceGlobSet::try_new(patterns).unwrap()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn select_files_skips_oversized_files() {
|
||||
let selected = select_files_to_collect(&[sandbox_file("huge.xml", MAX_FILE_SIZE + 1)]);
|
||||
let selected = select_files_to_collect(vec![sandbox_file("huge.xml", MAX_FILE_SIZE + 1)]);
|
||||
|
||||
assert!(selected.is_empty());
|
||||
}
|
||||
|
|
@ -381,7 +254,7 @@ mod tests {
|
|||
sandbox_file("c.xml", 2000),
|
||||
];
|
||||
|
||||
let selected = select_files_to_collect(&discovered);
|
||||
let selected = select_files_to_collect(discovered);
|
||||
|
||||
assert_eq!(
|
||||
selected
|
||||
|
|
@ -398,7 +271,7 @@ mod tests {
|
|||
.map(|index| sandbox_file(&format!("file{index}.xml"), 9 * 1024 * 1024))
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let selected = select_files_to_collect(&discovered);
|
||||
let selected = select_files_to_collect(discovered);
|
||||
|
||||
assert_eq!(selected.len(), 5);
|
||||
}
|
||||
|
|
@ -409,7 +282,7 @@ mod tests {
|
|||
.map(|index| sandbox_file(&format!("file{index}.txt"), 100))
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let selected = select_files_to_collect(&discovered);
|
||||
let selected = select_files_to_collect(discovered);
|
||||
|
||||
assert_eq!(selected.len(), MAX_FILE_COUNT);
|
||||
}
|
||||
|
|
@ -430,7 +303,7 @@ mod tests {
|
|||
(".ai/plans/DRAFTING.md".to_string(), "drafting".to_string()),
|
||||
("README.md".to_string(), "readme".to_string()),
|
||||
]);
|
||||
let sandbox = AssetMockSandbox::new(contents);
|
||||
let sandbox = asset_sandbox(contents);
|
||||
let globs = workspace_globs(&[".ai/reports/*.md", ".ai/plans/????-??-??-*.md"]);
|
||||
|
||||
let summary = collect_artifacts(&sandbox, stage_dir.path(), &globs)
|
||||
|
|
@ -452,7 +325,7 @@ mod tests {
|
|||
#[tokio::test]
|
||||
async fn collect_artifacts_preserves_content_metadata() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let sandbox = AssetMockSandbox::new(HashMap::from([(
|
||||
let sandbox = asset_sandbox(HashMap::from([(
|
||||
"test-results/r.xml".to_string(),
|
||||
"<test/>".to_string(),
|
||||
)]));
|
||||
|
|
@ -482,10 +355,55 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn collect_artifacts_downloads_provider_resolved_paths() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let file = SandboxFile {
|
||||
path: "provider-object:report-1".to_string(),
|
||||
relative_path: "test-results/r.xml".to_string(),
|
||||
size: 7,
|
||||
};
|
||||
let sandbox = MockSandbox {
|
||||
files: HashMap::from([(file.path.clone(), "<test/>".to_string())]),
|
||||
..MockSandbox::linux()
|
||||
}
|
||||
.with_walk_files(vec![file]);
|
||||
let globs = workspace_globs(&["test-results/**"]);
|
||||
|
||||
let summary = collect_artifacts(&sandbox, stage_dir.path(), &globs)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(summary.files_copied, 1);
|
||||
assert_eq!(summary.captured_assets[0].path, "test-results/r.xml");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn collect_artifacts_rechecks_downloaded_file_size() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let content = "x".repeat(usize::try_from(MAX_FILE_SIZE + 1).unwrap());
|
||||
let file = sandbox_file("test-results/grew.bin", 1);
|
||||
let sandbox = MockSandbox {
|
||||
files: HashMap::from([(file.path.clone(), content)]),
|
||||
..MockSandbox::linux()
|
||||
}
|
||||
.with_walk_files(vec![file]);
|
||||
let globs = workspace_globs(&["test-results/**"]);
|
||||
|
||||
let summary = collect_artifacts(&sandbox, stage_dir.path(), &globs)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(summary.files_copied, 0);
|
||||
assert_eq!(summary.files_skipped, 1);
|
||||
assert!(summary.captured_assets.is_empty());
|
||||
assert!(!stage_dir.path().join("test-results/grew.bin").exists());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn collect_artifacts_prunes_dependency_and_build_directories() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let sandbox = AssetMockSandbox::new(HashMap::from([
|
||||
let sandbox = asset_sandbox(HashMap::from([
|
||||
(".ai/reports/keep.md".to_string(), "keep".to_string()),
|
||||
("target/report.md".to_string(), "target".to_string()),
|
||||
(
|
||||
|
|
@ -506,7 +424,7 @@ mod tests {
|
|||
#[tokio::test]
|
||||
async fn collect_artifacts_deduplicates_overlapping_patterns() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let sandbox = AssetMockSandbox::new(HashMap::from([(
|
||||
let sandbox = asset_sandbox(HashMap::from([(
|
||||
".ai/reports/summary.md".to_string(),
|
||||
"summary".to_string(),
|
||||
)]));
|
||||
|
|
@ -523,7 +441,7 @@ mod tests {
|
|||
#[tokio::test]
|
||||
async fn collect_artifacts_reports_traversal_errors() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let sandbox = AssetMockSandbox::new(HashMap::new()).with_walk_error("permission denied");
|
||||
let sandbox = asset_sandbox(HashMap::new()).with_walk_files_error("permission denied");
|
||||
let globs = workspace_globs(&["test-results/**"]);
|
||||
|
||||
let error = collect_artifacts(&sandbox, stage_dir.path(), &globs)
|
||||
|
|
@ -537,7 +455,7 @@ mod tests {
|
|||
#[tokio::test]
|
||||
async fn collect_artifacts_keeps_download_errors_non_fatal() {
|
||||
let stage_dir = tempfile::tempdir().unwrap();
|
||||
let sandbox = AssetMockSandbox::new(HashMap::new()).with_discovered(vec![
|
||||
let sandbox = asset_sandbox(HashMap::new()).with_walk_files(vec![
|
||||
sandbox_file("test-results/missing.xml", 100),
|
||||
sandbox_file("test-results/also-missing.xml", 200),
|
||||
]);
|
||||
|
|
|
|||
|
|
@ -14,6 +14,7 @@ use fabro_types::{ArtifactUpload, EventBody, RunId, StageId};
|
|||
use fabro_util::error::collect_chain;
|
||||
use fabro_util::workspace_glob::{WorkspaceGlobError, WorkspaceGlobSet};
|
||||
use tokio::fs;
|
||||
use tokio::sync::OnceCell;
|
||||
use tokio::time::sleep;
|
||||
|
||||
use crate::artifact::{normalize_durable_updates, offload_large_values, sync_artifacts_to_env};
|
||||
|
|
@ -45,6 +46,7 @@ pub(crate) struct ArtifactLifecycle {
|
|||
artifact_globs: std::result::Result<WorkspaceGlobSet, WorkspaceGlobError>,
|
||||
pub artifact_sink: Option<ArtifactSink>,
|
||||
captured_artifacts: std::sync::Mutex<HashSet<ArtifactIdentity>>,
|
||||
ledger_initialized: OnceCell<()>,
|
||||
/// Run-scoped stage execution allocator shared with `RunServices`.
|
||||
stage_executions: StageExecutionTracker,
|
||||
}
|
||||
|
|
@ -55,7 +57,7 @@ impl ArtifactLifecycle {
|
|||
run_store: RunStoreHandle,
|
||||
emitter: Arc<Emitter>,
|
||||
run_id: RunId,
|
||||
artifact_globs: Vec<String>,
|
||||
artifact_globs: &[String],
|
||||
artifact_sink: Option<ArtifactSink>,
|
||||
stage_executions: StageExecutionTracker,
|
||||
) -> Self {
|
||||
|
|
@ -67,6 +69,7 @@ impl ArtifactLifecycle {
|
|||
artifact_globs: WorkspaceGlobSet::try_new(artifact_globs),
|
||||
artifact_sink,
|
||||
captured_artifacts: std::sync::Mutex::new(HashSet::new()),
|
||||
ledger_initialized: OnceCell::new(),
|
||||
stage_executions,
|
||||
}
|
||||
}
|
||||
|
|
@ -81,19 +84,27 @@ impl ArtifactLifecycle {
|
|||
#[async_trait]
|
||||
impl RunLifecycle<WorkflowGraph> for ArtifactLifecycle {
|
||||
async fn on_run_start(&self, _graph: &WorkflowGraph, _state: &WfRunState) -> CoreResult<()> {
|
||||
let _ = self.artifact_globs()?;
|
||||
let ledger = self
|
||||
.rebuild_captured_artifact_ledger()
|
||||
.await
|
||||
.map_err(|err| {
|
||||
let rendered = collect_chain(err.as_ref()).join(": ");
|
||||
CoreError::Other(format!(
|
||||
"failed to rebuild captured artifact ledger: {rendered}"
|
||||
))
|
||||
})?;
|
||||
*self.captured_artifacts.lock().expect(
|
||||
"artifact mutex should not be poisoned: no code panics while holding this lock",
|
||||
) = ledger;
|
||||
let artifact_globs = self.artifact_globs()?;
|
||||
if artifact_globs.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
self.ledger_initialized
|
||||
.get_or_try_init(|| async {
|
||||
let ledger = self
|
||||
.rebuild_captured_artifact_ledger()
|
||||
.await
|
||||
.map_err(|err| {
|
||||
let rendered = collect_chain(err.as_ref()).join(": ");
|
||||
CoreError::Other(format!(
|
||||
"failed to rebuild captured artifact ledger: {rendered}"
|
||||
))
|
||||
})?;
|
||||
*self.captured_artifacts.lock().expect(
|
||||
"artifact mutex should not be poisoned: no code panics while holding this lock",
|
||||
) = ledger;
|
||||
Ok::<(), CoreError>(())
|
||||
})
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -177,7 +177,7 @@ impl WorkflowLifecycle {
|
|||
run_store.clone(),
|
||||
Arc::clone(emitter),
|
||||
run_options.run_id,
|
||||
run_options.artifact_globs(),
|
||||
run_options.artifact_glob_patterns(),
|
||||
artifact_sink,
|
||||
stage_executions.clone(),
|
||||
);
|
||||
|
|
|
|||
|
|
@ -58,8 +58,8 @@ impl RunOptions {
|
|||
git_author_from_settings(&self.settings)
|
||||
}
|
||||
|
||||
pub fn artifact_globs(&self) -> Vec<String> {
|
||||
self.settings.run.artifacts.include.clone()
|
||||
pub fn artifact_glob_patterns(&self) -> &[String] {
|
||||
&self.settings.run.artifacts.include
|
||||
}
|
||||
|
||||
/// Run branch name from git checkpoint options, if set.
|
||||
|
|
|
|||
|
|
@ -72,16 +72,31 @@ include = ["/tmp/*.md", "../reports/*.md", "src/[abc"]
|
|||
)
|
||||
.expect_err("invalid artifact workspace globs should not resolve");
|
||||
|
||||
let message = error.to_string();
|
||||
assert!(
|
||||
message.contains("run.artifacts.include[0]")
|
||||
&& message.contains("must be relative")
|
||||
&& message.contains("run.artifacts.include[1]")
|
||||
&& message.contains("parent directory")
|
||||
&& message.contains("run.artifacts.include[2]")
|
||||
&& message.contains("invalid workspace glob"),
|
||||
"unexpected error: {message}"
|
||||
let errors = match error {
|
||||
crate::Error::Resolve { errors, .. } => errors,
|
||||
other => panic!("expected structured resolve errors, got {other:#}"),
|
||||
};
|
||||
let invalid = errors
|
||||
.into_iter()
|
||||
.map(|error| match error {
|
||||
crate::ResolveError::Invalid { path, reason } => (path, reason),
|
||||
other => panic!("expected invalid-value error, got {other}"),
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
assert_eq!(
|
||||
invalid
|
||||
.iter()
|
||||
.map(|(path, _)| path.as_str())
|
||||
.collect::<Vec<_>>(),
|
||||
vec![
|
||||
"run.artifacts.include[0]",
|
||||
"run.artifacts.include[1]",
|
||||
"run.artifacts.include[2]",
|
||||
]
|
||||
);
|
||||
assert!(invalid[0].1.contains("must be relative"));
|
||||
assert!(invalid[1].1.contains("parent directory"));
|
||||
assert!(invalid[2].1.contains("invalid workspace glob"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
|
|
@ -14,7 +14,6 @@ const MATCH_OPTIONS: glob::MatchOptions = glob::MatchOptions {
|
|||
/// path segment; `**` crosses directory boundaries.
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct WorkspaceGlob {
|
||||
source: String,
|
||||
pattern: glob::Pattern,
|
||||
traversal_root: String,
|
||||
}
|
||||
|
|
@ -25,6 +24,11 @@ impl WorkspaceGlob {
|
|||
if source.is_empty() {
|
||||
return Err(WorkspaceGlobError::Empty);
|
||||
}
|
||||
if source.contains('\\') {
|
||||
return Err(WorkspaceGlobError::BackslashSeparator {
|
||||
pattern: source.to_string(),
|
||||
});
|
||||
}
|
||||
if is_absolute(source) {
|
||||
return Err(WorkspaceGlobError::Absolute {
|
||||
pattern: source.to_string(),
|
||||
|
|
@ -43,17 +47,11 @@ impl WorkspaceGlob {
|
|||
})?;
|
||||
|
||||
Ok(Self {
|
||||
source: source.to_string(),
|
||||
pattern,
|
||||
traversal_root: literal_traversal_root(source),
|
||||
})
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn as_str(&self) -> &str {
|
||||
&self.source
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn is_match(&self, relative_path: &str) -> bool {
|
||||
let relative_path = normalize_candidate(relative_path);
|
||||
|
|
@ -137,6 +135,9 @@ pub enum WorkspaceGlobError {
|
|||
#[error("workspace glob cannot be empty")]
|
||||
Empty,
|
||||
|
||||
#[error("workspace glob must use '/' as its path separator: {pattern:?}")]
|
||||
BackslashSeparator { pattern: String },
|
||||
|
||||
#[error("workspace glob must be relative: {pattern:?}")]
|
||||
Absolute { pattern: String },
|
||||
|
||||
|
|
@ -161,7 +162,6 @@ fn strip_current_dir_prefix(mut path: &str) -> &str {
|
|||
fn is_absolute(path: &str) -> bool {
|
||||
let bytes = path.as_bytes();
|
||||
path.starts_with('/')
|
||||
|| path.starts_with("//")
|
||||
|| Path::new(path).is_absolute()
|
||||
|| matches!(bytes, [drive, b':', ..] if drive.is_ascii_alphabetic())
|
||||
}
|
||||
|
|
@ -257,7 +257,6 @@ mod tests {
|
|||
fn workspace_glob_normalizes_a_leading_current_directory() {
|
||||
let glob = WorkspaceGlob::try_new("./src/*.rs").unwrap();
|
||||
|
||||
assert_eq!(glob.as_str(), "src/*.rs");
|
||||
assert!(glob.is_match("./src/lib.rs"));
|
||||
assert_eq!(glob.traversal_root(), "src");
|
||||
}
|
||||
|
|
@ -276,6 +275,14 @@ mod tests {
|
|||
WorkspaceGlob::try_new("C:/tmp/*.md"),
|
||||
Err(WorkspaceGlobError::Absolute { .. })
|
||||
));
|
||||
assert!(matches!(
|
||||
WorkspaceGlob::try_new(r"dir\*.rs"),
|
||||
Err(WorkspaceGlobError::BackslashSeparator { .. })
|
||||
));
|
||||
assert!(matches!(
|
||||
WorkspaceGlob::try_new(r"\\server\share\*.md"),
|
||||
Err(WorkspaceGlobError::BackslashSeparator { .. })
|
||||
));
|
||||
assert!(matches!(
|
||||
WorkspaceGlob::try_new("../*.md"),
|
||||
Err(WorkspaceGlobError::ParentTraversal { .. })
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue