mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
Merge remote-tracking branch 'origin/main'
This commit is contained in:
commit
7feef1c6b2
34 changed files with 2180 additions and 489 deletions
File diff suppressed because it is too large
Load diff
|
|
@ -4,8 +4,10 @@
|
|||
)]
|
||||
|
||||
mod args;
|
||||
mod manifest_builder;
|
||||
|
||||
use clap::{Command, CommandFactory};
|
||||
pub use manifest_builder::{BuiltManifest, ManifestBuildInput, build_run_manifest};
|
||||
|
||||
pub fn command_for_reference() -> Command {
|
||||
args::Cli::command()
|
||||
|
|
|
|||
|
|
@ -10,6 +10,10 @@ mod gh;
|
|||
mod landing;
|
||||
mod local_server;
|
||||
mod logging;
|
||||
#[allow(
|
||||
unreachable_pub,
|
||||
reason = "The library exports manifest builder helpers for tests; the binary includes the same module privately."
|
||||
)]
|
||||
mod manifest_builder;
|
||||
mod server_client;
|
||||
mod server_runs;
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ use fabro_graphviz::graph::AttrValue;
|
|||
use fabro_graphviz::parser;
|
||||
use fabro_types::settings::run::{ResolvedGoalSource, ResolvedRunGoal};
|
||||
use fabro_types::{DirtyStatus, GitContext, PreRunPushOutcome, RunId, WorkflowSettings};
|
||||
use fabro_workflow::ManifestPath;
|
||||
use fabro_workflow::git::{
|
||||
GitSyncStatus, branch_needs_push, head_sha, push_branch_noninteractive, sync_status,
|
||||
};
|
||||
|
|
@ -22,7 +23,7 @@ use fabro_workflow::git::{
|
|||
use crate::args::{PreflightArgs, RunArgs};
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct ManifestBuildInput {
|
||||
pub struct ManifestBuildInput {
|
||||
pub workflow: PathBuf,
|
||||
pub cwd: PathBuf,
|
||||
pub run_overrides: Option<RunLayer>,
|
||||
|
|
@ -35,7 +36,7 @@ pub(crate) struct ManifestBuildInput {
|
|||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct BuiltManifest {
|
||||
pub struct BuiltManifest {
|
||||
pub manifest: types::RunManifest,
|
||||
pub target_path: PathBuf,
|
||||
}
|
||||
|
|
@ -49,11 +50,11 @@ struct CollectContext<'a> {
|
|||
#[derive(Clone)]
|
||||
struct WorkflowScanInput {
|
||||
absolute_dot_path: PathBuf,
|
||||
logical_dot_path: PathBuf,
|
||||
dot_path: ManifestPath,
|
||||
source: String,
|
||||
}
|
||||
|
||||
pub(crate) fn build_run_manifest(input: ManifestBuildInput) -> Result<BuiltManifest> {
|
||||
pub fn build_run_manifest(input: ManifestBuildInput) -> Result<BuiltManifest> {
|
||||
let root_resolution = resolve_workflow_path(&input.workflow, &input.cwd)?;
|
||||
if root_resolution.workflow_toml_path.is_none()
|
||||
&& !root_resolution.resolved_workflow_path.is_file()
|
||||
|
|
@ -92,8 +93,8 @@ pub(crate) fn build_run_manifest(input: ManifestBuildInput) -> Result<BuiltManif
|
|||
.build()
|
||||
.map_err(|errors| anyhow!("failed to resolve manifest settings: {errors}"))?;
|
||||
let target_path = root_resolution.dot_path.clone();
|
||||
let target_logical_path = to_logical_path(&target_path, &input.cwd)?;
|
||||
let target_logical_path_string = logical_path_string(&target_logical_path);
|
||||
let target_manifest_path = manifest_path_from_absolute(&target_path, &input.cwd)?;
|
||||
let target_key = target_manifest_path.to_string();
|
||||
|
||||
let mut context = CollectContext {
|
||||
cwd: &input.cwd,
|
||||
|
|
@ -104,7 +105,7 @@ pub(crate) fn build_run_manifest(input: ManifestBuildInput) -> Result<BuiltManif
|
|||
|
||||
let root_source = context
|
||||
.workflows
|
||||
.get(&target_logical_path_string)
|
||||
.get(&target_key)
|
||||
.map(|workflow| workflow.source.clone())
|
||||
.ok_or_else(|| anyhow!("root workflow missing from manifest bundle"))?;
|
||||
|
||||
|
|
@ -153,7 +154,7 @@ pub(crate) fn build_run_manifest(input: ManifestBuildInput) -> Result<BuiltManif
|
|||
run_id: input.run_id.map(|run_id| run_id.to_string()),
|
||||
target: types::ManifestTarget {
|
||||
identifier: input.workflow.display().to_string(),
|
||||
path: target_logical_path_string,
|
||||
path: target_key,
|
||||
},
|
||||
version: 1,
|
||||
workflows: context.workflows,
|
||||
|
|
@ -220,9 +221,9 @@ fn collect_workflow_entry(
|
|||
workflow.to_path_buf()
|
||||
};
|
||||
let resolution = resolve_workflow_path(&normalized_workflow, resolve_from)?;
|
||||
let logical_dot_path = to_logical_path(&resolution.dot_path, context.cwd)?;
|
||||
let logical_dot_key = logical_path_string(&logical_dot_path);
|
||||
if !context.visited_workflows.insert(logical_dot_key.clone()) {
|
||||
let dot_path = manifest_path_from_absolute(&resolution.dot_path, context.cwd)?;
|
||||
let dot_key = dot_path.to_string();
|
||||
if !context.visited_workflows.insert(dot_key.clone()) {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
|
|
@ -230,7 +231,7 @@ fn collect_workflow_entry(
|
|||
.with_context(|| format!("Failed to read {}", resolution.dot_path.display()))?;
|
||||
let config = if let Some(workflow_toml_path) = resolution.workflow_toml_path.as_ref() {
|
||||
Some(types::ManifestWorkflowConfig {
|
||||
path: logical_path_string(&to_logical_path(workflow_toml_path, context.cwd)?),
|
||||
path: manifest_path_from_absolute(workflow_toml_path, context.cwd)?.to_string(),
|
||||
source: std::fs::read_to_string(workflow_toml_path)
|
||||
.with_context(|| format!("Failed to read {}", workflow_toml_path.display()))?,
|
||||
})
|
||||
|
|
@ -240,7 +241,7 @@ fn collect_workflow_entry(
|
|||
|
||||
let scan = WorkflowScanInput {
|
||||
absolute_dot_path: resolution.dot_path,
|
||||
logical_dot_path,
|
||||
dot_path,
|
||||
source: source.clone(),
|
||||
};
|
||||
let mut files = HashMap::new();
|
||||
|
|
@ -250,13 +251,11 @@ fn collect_workflow_entry(
|
|||
}
|
||||
collect_workflow_files(context, &scan, &mut files, &mut visited_imports)?;
|
||||
|
||||
context
|
||||
.workflows
|
||||
.insert(logical_dot_key, types::ManifestWorkflow {
|
||||
config,
|
||||
files,
|
||||
source,
|
||||
});
|
||||
context.workflows.insert(dot_key, types::ManifestWorkflow {
|
||||
config,
|
||||
files,
|
||||
source,
|
||||
});
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -285,7 +284,7 @@ fn collect_workflow_files(
|
|||
context.cwd,
|
||||
goal_ref.trim_start_matches('@'),
|
||||
types::ManifestFileRefType::FileInline,
|
||||
Some(workflow.logical_dot_path.clone()),
|
||||
Some(workflow.dot_path.clone()),
|
||||
)?;
|
||||
}
|
||||
}
|
||||
|
|
@ -302,7 +301,7 @@ fn collect_workflow_files(
|
|||
context.cwd,
|
||||
prompt_ref.trim_start_matches('@'),
|
||||
types::ManifestFileRefType::FileInline,
|
||||
Some(workflow.logical_dot_path.clone()),
|
||||
Some(workflow.dot_path.clone()),
|
||||
)?;
|
||||
}
|
||||
}
|
||||
|
|
@ -317,9 +316,9 @@ fn collect_workflow_files(
|
|||
context.cwd,
|
||||
import_ref,
|
||||
types::ManifestFileRefType::Import,
|
||||
Some(workflow.logical_dot_path.clone()),
|
||||
Some(workflow.dot_path.clone()),
|
||||
)?;
|
||||
let import_key = logical_path_string(&imported.logical_path);
|
||||
let import_key = imported.path.to_string();
|
||||
if visited_imports.insert(import_key) {
|
||||
let imported_source = std::fs::read_to_string(&imported.absolute_path)
|
||||
.with_context(|| {
|
||||
|
|
@ -327,7 +326,7 @@ fn collect_workflow_files(
|
|||
})?;
|
||||
let imported_scan = WorkflowScanInput {
|
||||
absolute_dot_path: imported.absolute_path,
|
||||
logical_dot_path: imported.logical_path,
|
||||
dot_path: imported.path,
|
||||
source: imported_source,
|
||||
};
|
||||
collect_workflow_files(context, &imported_scan, files, visited_imports)?;
|
||||
|
|
@ -380,21 +379,25 @@ fn collect_workflow_config_files(
|
|||
return Ok(());
|
||||
};
|
||||
|
||||
let config_path = context.cwd.join(&config.path);
|
||||
let config_path = ManifestPath::from_wire(&config.path)
|
||||
.ok_or_else(|| anyhow!("invalid manifest workflow config path: {}", config.path))?;
|
||||
let absolute_config_path = context.cwd.join(config_path.as_path());
|
||||
collect_bundled_file(
|
||||
files,
|
||||
config_path.parent().unwrap_or_else(|| Path::new(".")),
|
||||
absolute_config_path
|
||||
.parent()
|
||||
.unwrap_or_else(|| Path::new(".")),
|
||||
context.cwd,
|
||||
path,
|
||||
types::ManifestFileRefType::Dockerfile,
|
||||
Some(PathBuf::from(&config.path)),
|
||||
Some(config_path),
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
struct BundledFile {
|
||||
absolute_path: PathBuf,
|
||||
logical_path: PathBuf,
|
||||
path: ManifestPath,
|
||||
}
|
||||
|
||||
fn collect_bundled_file(
|
||||
|
|
@ -403,19 +406,19 @@ fn collect_bundled_file(
|
|||
cwd: &Path,
|
||||
reference: &str,
|
||||
ref_type: types::ManifestFileRefType,
|
||||
from: Option<PathBuf>,
|
||||
from: Option<ManifestPath>,
|
||||
) -> Result<BundledFile> {
|
||||
let absolute_path = normalize_absolute_path(base_dir, reference)
|
||||
.ok_or_else(|| anyhow!("unsupported manifest reference: {reference}"))?;
|
||||
let logical_path = to_logical_path(&absolute_path, cwd)?;
|
||||
let key = logical_path_string(&logical_path);
|
||||
let path = manifest_path_from_absolute(&absolute_path, cwd)?;
|
||||
let key = path.to_string();
|
||||
if !files.contains_key(&key) {
|
||||
let content = std::fs::read_to_string(&absolute_path)
|
||||
.with_context(|| format!("Failed to read {}", absolute_path.display()))?;
|
||||
files.insert(key.clone(), types::ManifestFileEntry {
|
||||
content,
|
||||
ref_: types::ManifestFileRef {
|
||||
from: from.map(|value| logical_path_string(&value)),
|
||||
from: from.map(|value| value.to_string()),
|
||||
original: reference.to_string(),
|
||||
type_: ref_type,
|
||||
},
|
||||
|
|
@ -424,7 +427,7 @@ fn collect_bundled_file(
|
|||
|
||||
Ok(BundledFile {
|
||||
absolute_path,
|
||||
logical_path,
|
||||
path,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -625,49 +628,9 @@ fn normalize_absolute_path(base_dir: &Path, reference: &str) -> Option<PathBuf>
|
|||
Some(normalized)
|
||||
}
|
||||
|
||||
fn to_logical_path(path: &Path, cwd: &Path) -> Result<PathBuf> {
|
||||
if let Ok(stripped) = path.strip_prefix(cwd) {
|
||||
return Ok(stripped.to_path_buf());
|
||||
}
|
||||
|
||||
relative_path_from(path, cwd)
|
||||
.ok_or_else(|| anyhow!("Failed to compute logical path for {}", path.display()))
|
||||
}
|
||||
|
||||
fn relative_path_from(path: &Path, base: &Path) -> Option<PathBuf> {
|
||||
let path_components = path.components().collect::<Vec<_>>();
|
||||
let base_components = base.components().collect::<Vec<_>>();
|
||||
if path_components.is_empty() || base_components.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut common = 0;
|
||||
while common < path_components.len()
|
||||
&& common < base_components.len()
|
||||
&& path_components[common] == base_components[common]
|
||||
{
|
||||
common += 1;
|
||||
}
|
||||
|
||||
let mut relative = PathBuf::new();
|
||||
for component in &base_components[common..] {
|
||||
if matches!(component, Component::Normal(_)) {
|
||||
relative.push("..");
|
||||
}
|
||||
}
|
||||
for component in &path_components[common..] {
|
||||
match component {
|
||||
Component::Normal(part) => relative.push(part),
|
||||
Component::CurDir => {}
|
||||
Component::ParentDir => relative.push(".."),
|
||||
Component::RootDir | Component::Prefix(_) => return None,
|
||||
}
|
||||
}
|
||||
Some(relative)
|
||||
}
|
||||
|
||||
fn logical_path_string(path: &Path) -> String {
|
||||
path.to_string_lossy().to_string()
|
||||
fn manifest_path_from_absolute(path: &Path, cwd: &Path) -> Result<ManifestPath> {
|
||||
ManifestPath::from_absolute(path, cwd)
|
||||
.ok_or_else(|| anyhow!("Failed to compute manifest path for {}", path.display()))
|
||||
}
|
||||
|
||||
fn manifest_args_is_empty(args: &types::ManifestArgs) -> bool {
|
||||
|
|
|
|||
57
lib/crates/fabro-cli/tests/manifest_path_round_trip.rs
Normal file
57
lib/crates/fabro-cli/tests/manifest_path_round_trip.rs
Normal file
|
|
@ -0,0 +1,57 @@
|
|||
#![expect(
|
||||
clippy::disallowed_methods,
|
||||
reason = "Sync temp fixture writes keep this manifest round-trip test simple and isolated."
|
||||
)]
|
||||
|
||||
use std::path::PathBuf;
|
||||
|
||||
use fabro_cli::{ManifestBuildInput, build_run_manifest};
|
||||
use fabro_workflow::ManifestPath;
|
||||
|
||||
#[test]
|
||||
fn cli_built_manifest_resolves_user_global_at_path() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let workflow_dir = temp.path().join(".fabro/workflows/demo");
|
||||
let project = temp.path().join("project");
|
||||
std::fs::create_dir_all(workflow_dir.join("prompts")).unwrap();
|
||||
std::fs::create_dir_all(&project).unwrap();
|
||||
std::fs::write(
|
||||
workflow_dir.join("workflow.fabro"),
|
||||
r#"digraph Demo {
|
||||
graph [goal="Demo"]
|
||||
start [shape=Mdiamond]
|
||||
prompt [prompt="@prompts/hello.md"]
|
||||
exit [shape=Msquare]
|
||||
start -> prompt -> exit
|
||||
}"#,
|
||||
)
|
||||
.unwrap();
|
||||
std::fs::write(workflow_dir.join("prompts/hello.md"), "hello from bundle").unwrap();
|
||||
|
||||
let built = build_run_manifest(ManifestBuildInput {
|
||||
workflow: workflow_dir.join("workflow.fabro"),
|
||||
cwd: project,
|
||||
run_overrides: None,
|
||||
cli_overrides: None,
|
||||
args: None,
|
||||
run_id: None,
|
||||
user_settings_path: None,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
let bundle = fabro_server::workflow_bundle_from_manifest(&built.manifest.workflows).unwrap();
|
||||
let target_path = ManifestPath::from_wire(&built.manifest.target.path).unwrap();
|
||||
let workflow = bundle
|
||||
.workflow(&target_path)
|
||||
.expect("root workflow should be present");
|
||||
let resolved = workflow
|
||||
.file_resolver()
|
||||
.resolve(&workflow.current_dir(), "prompts/hello.md")
|
||||
.expect("prompt should resolve from bundle");
|
||||
|
||||
assert_eq!(resolved.content, "hello from bundle");
|
||||
assert_eq!(
|
||||
resolved.path,
|
||||
PathBuf::from("../.fabro/workflows/demo/prompts/hello.md")
|
||||
);
|
||||
}
|
||||
|
|
@ -13,6 +13,7 @@ use tokio::{fs, time};
|
|||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
use crate::clone_source::{self, CloneDecision, EmptyWorkspaceReason};
|
||||
use crate::redact::redact_auth_url;
|
||||
use crate::sandbox::resolve_path;
|
||||
use crate::{
|
||||
DirEntry, ExecResult, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback,
|
||||
|
|
@ -54,6 +55,23 @@ fn elapsed_ms(start: &Instant) -> u64 {
|
|||
u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX)
|
||||
}
|
||||
|
||||
fn command_kind(command: &str) -> &'static str {
|
||||
match command.split_whitespace().next().unwrap_or_default() {
|
||||
"git" => "git",
|
||||
"sh" | "/bin/sh" => "sh",
|
||||
"bash" | "/bin/bash" => "bash",
|
||||
"rg" => "rg",
|
||||
"grep" => "grep",
|
||||
"find" => "find",
|
||||
"cat" => "cat",
|
||||
"ls" => "ls",
|
||||
"mkdir" => "mkdir",
|
||||
"rm" => "rm",
|
||||
"printf" => "printf",
|
||||
_ => "other",
|
||||
}
|
||||
}
|
||||
|
||||
/// Sandbox that runs all operations inside a Daytona cloud sandbox.
|
||||
pub struct DaytonaSandbox {
|
||||
config: DaytonaConfig,
|
||||
|
|
@ -639,22 +657,25 @@ impl Sandbox for DaytonaSandbox {
|
|||
let wrapped = wrap_bash_command(&cmd);
|
||||
match ps.execute_command(&wrapped, opts).await {
|
||||
Ok(r) if r.exit_code != 0 => {
|
||||
let stderr = r.result.replace(
|
||||
&auth_url.raw_string(),
|
||||
&auth_url.redacted_string(),
|
||||
let err = crate::Error::exec(
|
||||
"git remote set-url origin (Daytona post-clone)",
|
||||
r.exit_code,
|
||||
false,
|
||||
0,
|
||||
redact_auth_url(&r.result, Some(&auth_url)),
|
||||
String::new(),
|
||||
);
|
||||
tracing::warn!(
|
||||
exit_code = r.exit_code,
|
||||
output = %stderr.trim(),
|
||||
error = %err,
|
||||
"Failed to set Daytona sandbox push credentials \
|
||||
on origin — subsequent git push from this \
|
||||
sandbox will fail"
|
||||
);
|
||||
}
|
||||
Ok(_) => {}
|
||||
Err(e) => {
|
||||
Err(_) => {
|
||||
tracing::warn!(
|
||||
error = %e,
|
||||
error_class = "daytona_set_url_exec_failed",
|
||||
"Daytona exec failed while setting push credentials \
|
||||
on origin — subsequent git push from this \
|
||||
sandbox will fail"
|
||||
|
|
@ -788,9 +809,9 @@ impl Sandbox for DaytonaSandbox {
|
|||
)]
|
||||
}
|
||||
|
||||
async fn git_push_ref(&self, refspec: &str) -> bool {
|
||||
async fn git_push_ref(&self, refspec: &str) -> crate::Result<()> {
|
||||
if !self.repo_cloned() {
|
||||
return false;
|
||||
return Ok(());
|
||||
}
|
||||
crate::git_push_via_exec(self, refspec).await
|
||||
}
|
||||
|
|
@ -857,7 +878,9 @@ impl Sandbox for DaytonaSandbox {
|
|||
origin_url,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| crate::Error::message(format!("Failed to refresh GitHub App token: {e}")))?;
|
||||
.map_err(|_| {
|
||||
crate::Error::message("Failed to refresh push credentials: token_mint_failed")
|
||||
})?;
|
||||
|
||||
let cmd = format!(
|
||||
"git -c maintenance.auto=0 remote set-url origin {}",
|
||||
|
|
@ -866,15 +889,14 @@ impl Sandbox for DaytonaSandbox {
|
|||
let result = self
|
||||
.exec_command(&cmd, 10_000, None, None, None)
|
||||
.await
|
||||
.map_err(|e| crate::Error::context("Failed to set refreshed push credentials", e))?;
|
||||
.map_err(|_| {
|
||||
crate::Error::message("Failed to refresh push credentials: set_url_exec_failed")
|
||||
})?;
|
||||
if result.exit_code != 0 {
|
||||
let stderr = result
|
||||
.stderr
|
||||
.replace(&auth_url.raw_string(), &auth_url.redacted_string());
|
||||
return Err(crate::Error::message(format!(
|
||||
"Failed to set refreshed push credentials (exit {}): {}",
|
||||
result.exit_code, stderr
|
||||
)));
|
||||
return Err(result.into_exec_error_with_redactor(
|
||||
"git remote set-url origin (refresh push credentials)",
|
||||
|s| redact_auth_url(s, Some(&auth_url)),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
|
|
@ -1021,7 +1043,12 @@ impl Sandbox for DaytonaSandbox {
|
|||
env_vars: Option<&HashMap<String, String>>,
|
||||
cancel_token: Option<CancellationToken>,
|
||||
) -> crate::Result<ExecResult> {
|
||||
tracing::info!(command, timeout_ms, "exec_command: entered");
|
||||
tracing::info!(
|
||||
timeout_ms,
|
||||
command_kind = command_kind(command),
|
||||
command_len = command.len(),
|
||||
"exec_command: entered"
|
||||
);
|
||||
|
||||
let sandbox = self.sandbox()?;
|
||||
let start = Instant::now();
|
||||
|
|
@ -1248,6 +1275,24 @@ mod tests {
|
|||
assert!(config.labels.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn command_kind_classifies_known_prefixes() {
|
||||
assert_eq!(command_kind(" git status"), "git");
|
||||
assert_eq!(command_kind("bash -lc 'echo ok'"), "bash");
|
||||
assert_eq!(command_kind("/bin/sh -c 'echo ok'"), "sh");
|
||||
assert_eq!(command_kind("rg --version"), "rg");
|
||||
assert_eq!(command_kind("find . -maxdepth 1"), "find");
|
||||
assert_eq!(command_kind(""), "other");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn command_kind_does_not_echo_auth_url_commands() {
|
||||
assert_eq!(
|
||||
command_kind("https://x-access-token:ghs_FAKE@github.com/owner/repo.git"),
|
||||
"other"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wrap_bash_uses_base64_encoding() {
|
||||
let wrapped = wrap_bash_command("echo hello");
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ use tokio::{fs, time};
|
|||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
use crate::clone_source::{self, CloneDecision, EmptyWorkspaceReason};
|
||||
use crate::redact::redact_auth_url;
|
||||
use crate::sandbox::resolve_path;
|
||||
use crate::{
|
||||
DirEntry, ExecResult, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback,
|
||||
|
|
@ -423,10 +424,12 @@ impl DockerSandbox {
|
|||
.docker_exec_shell(&command, 10_000, Some(WORKING_DIRECTORY), None, None)
|
||||
.await?;
|
||||
if result.exit_code != 0 {
|
||||
let stderr = redact_auth_url(&result.stderr, Some(auth_url));
|
||||
let err = result
|
||||
.into_exec_error_with_redactor("git remote set-url origin (post-clone)", |s| {
|
||||
redact_auth_url(s, Some(auth_url))
|
||||
});
|
||||
tracing::warn!(
|
||||
exit_code = result.exit_code,
|
||||
stderr = %stderr.trim(),
|
||||
error = %err,
|
||||
"Failed to set Docker sandbox push credentials on origin — \
|
||||
subsequent git push from this sandbox will fail"
|
||||
);
|
||||
|
|
@ -634,13 +637,6 @@ fn bash_remediation(error: &DockerError, image: &str) -> String {
|
|||
)
|
||||
}
|
||||
|
||||
fn redact_auth_url(text: &str, auth_url: Option<&fabro_redact::DisplaySafeUrl>) -> String {
|
||||
let Some(auth_url) = auth_url else {
|
||||
return text.to_string();
|
||||
};
|
||||
text.replace(&auth_url.raw_string(), &auth_url.redacted_string())
|
||||
}
|
||||
|
||||
fn build_single_file_tar(file_name: &str, bytes: &[u8]) -> crate::Result<Vec<u8>> {
|
||||
let mut tar_builder = tar::Builder::new(Vec::new());
|
||||
let mut header = tar::Header::new_gnu();
|
||||
|
|
@ -1208,9 +1204,9 @@ impl Sandbox for DockerSandbox {
|
|||
)]
|
||||
}
|
||||
|
||||
async fn git_push_ref(&self, refspec: &str) -> bool {
|
||||
async fn git_push_ref(&self, refspec: &str) -> crate::Result<()> {
|
||||
if !self.repo_cloned() {
|
||||
return false;
|
||||
return Ok(());
|
||||
}
|
||||
crate::git_push_via_exec(self, refspec).await
|
||||
}
|
||||
|
|
@ -1254,7 +1250,9 @@ impl Sandbox for DockerSandbox {
|
|||
origin_url,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| crate::Error::message(format!("Failed to refresh GitHub App token: {e}")))?;
|
||||
.map_err(|_| {
|
||||
crate::Error::message("Failed to refresh push credentials: token_mint_failed")
|
||||
})?;
|
||||
|
||||
let command = format!(
|
||||
"git -c maintenance.auto=0 remote set-url origin {}",
|
||||
|
|
@ -1264,11 +1262,10 @@ impl Sandbox for DockerSandbox {
|
|||
.docker_exec_shell(&command, 10_000, Some(WORKING_DIRECTORY), None, None)
|
||||
.await?;
|
||||
if result.exit_code != 0 {
|
||||
let stderr = redact_auth_url(&result.stderr, Some(&auth_url));
|
||||
return Err(crate::Error::message(format!(
|
||||
"Failed to set refreshed push credentials (exit {}): {}",
|
||||
result.exit_code, stderr
|
||||
)));
|
||||
return Err(result.into_exec_error_with_redactor(
|
||||
"git remote set-url origin (refresh push credentials)",
|
||||
|s| redact_auth_url(s, Some(&auth_url)),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
|
|
|
|||
|
|
@ -36,6 +36,21 @@ pub enum Error {
|
|||
#[source]
|
||||
source: BollardError,
|
||||
},
|
||||
|
||||
#[error(
|
||||
"{label} failed (exit {exit_code}, timed_out={timed_out}, duration_ms={duration_ms}) - hint: {hint}",
|
||||
hint = classify_exec_failure(stderr)
|
||||
.or_else(|| classify_exec_failure(stdout))
|
||||
.unwrap_or("unclassified")
|
||||
)]
|
||||
Exec {
|
||||
label: String,
|
||||
exit_code: i32,
|
||||
timed_out: bool,
|
||||
duration_ms: u64,
|
||||
stderr: String,
|
||||
stdout: String,
|
||||
},
|
||||
}
|
||||
|
||||
impl Error {
|
||||
|
|
@ -53,6 +68,24 @@ impl Error {
|
|||
}
|
||||
}
|
||||
|
||||
pub fn exec(
|
||||
label: impl Into<String>,
|
||||
exit_code: i32,
|
||||
timed_out: bool,
|
||||
duration_ms: u64,
|
||||
stderr: impl Into<String>,
|
||||
stdout: impl Into<String>,
|
||||
) -> Self {
|
||||
Self::Exec {
|
||||
label: label.into(),
|
||||
exit_code,
|
||||
timed_out,
|
||||
duration_ms,
|
||||
stderr: stderr.into(),
|
||||
stdout: stdout.into(),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "docker")]
|
||||
pub fn docker_connect(source: BollardError) -> Self {
|
||||
Self::DockerConnect { source }
|
||||
|
|
@ -95,4 +128,117 @@ impl From<&str> for Error {
|
|||
}
|
||||
}
|
||||
|
||||
pub(crate) fn classify_exec_failure(stderr: &str) -> Option<&'static str> {
|
||||
let lower = stderr.to_ascii_lowercase();
|
||||
if lower.contains("could not read username") || lower.contains("terminal prompts disabled") {
|
||||
Some(
|
||||
"no credentials in origin URL - check that the sandbox forwarded \
|
||||
GITHUB_APP_PRIVATE_KEY (or GITHUB_TOKEN) and that refresh_push_credentials succeeded",
|
||||
)
|
||||
} else if lower.contains("permission to") && lower.contains("denied") {
|
||||
Some(
|
||||
"github denied the push - installation token lacks contents:write \
|
||||
on this repo, or a branch protection / push ruleset is rejecting the ref",
|
||||
)
|
||||
} else if lower.contains("protected branch")
|
||||
|| lower.contains("ruleset")
|
||||
|| lower.contains("rejected")
|
||||
{
|
||||
Some("github rejected the ref - likely a branch protection rule or push ruleset")
|
||||
} else if lower.contains("authentication failed") || lower.contains("invalid username") {
|
||||
Some("github authentication failed - installation token may be expired or wrong scope")
|
||||
} else if lower.contains("could not resolve host") || lower.contains("network is unreachable") {
|
||||
Some("network failure inside sandbox - check DNS / egress from the run container")
|
||||
} else if lower.contains("repository not found") {
|
||||
Some("github 404 - the App installation may not include this repo")
|
||||
} else if lower.contains("no such remote") && lower.contains("origin") {
|
||||
Some("origin remote missing - push credentials could not be installed")
|
||||
} else if lower.contains("not a git repository")
|
||||
|| lower.contains("does not appear to be a git repository")
|
||||
{
|
||||
Some("git repository unavailable in sandbox working directory")
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
pub type Result<T> = std::result::Result<T, Error>;
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn exec_display_is_log_safe() {
|
||||
let stderr = "fatal: unable to access \
|
||||
'https://x-access-token:ghs_xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA@github.com/owner/repo/':\n\
|
||||
remote: Permission to owner/repo.git denied\n\
|
||||
identity ~/.ssh/id_rsa_work";
|
||||
let error = Error::exec(
|
||||
"git push origin refs/heads/run",
|
||||
128,
|
||||
false,
|
||||
210,
|
||||
stderr,
|
||||
"",
|
||||
);
|
||||
let rendered = error.to_string();
|
||||
|
||||
for forbidden in [
|
||||
"fatal:",
|
||||
"remote:",
|
||||
"x-access-token",
|
||||
"ghs_xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA",
|
||||
"~/.ssh",
|
||||
"id_rsa_work",
|
||||
] {
|
||||
assert!(
|
||||
!rendered.contains(forbidden),
|
||||
"Display leaked {forbidden:?}: {rendered}"
|
||||
);
|
||||
}
|
||||
assert!(rendered.contains("git push origin refs/heads/run"));
|
||||
assert!(rendered.contains("exit 128"));
|
||||
assert!(rendered.contains("timed_out=false"));
|
||||
assert!(rendered.contains("duration_ms=210"));
|
||||
assert!(rendered.contains("hint:"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_exec_failure_documents_known_branches() {
|
||||
let cases = [
|
||||
(
|
||||
"fatal: could not read Username for 'https://github.com'",
|
||||
"no credentials in origin URL",
|
||||
),
|
||||
(
|
||||
"remote: Permission to owner/repo.git denied to fabro-app[bot].",
|
||||
"github denied the push",
|
||||
),
|
||||
(
|
||||
"remote: error: GH013: Repository rule violations found due to ruleset",
|
||||
"github rejected the ref",
|
||||
),
|
||||
(
|
||||
"fatal: Authentication failed for 'https://github.com/owner/repo'",
|
||||
"github authentication failed",
|
||||
),
|
||||
(
|
||||
"fatal: could not resolve host: github.com",
|
||||
"network failure",
|
||||
),
|
||||
("remote: Repository not found.", "github 404"),
|
||||
("error: No such remote 'origin'", "origin remote missing"),
|
||||
("fatal: not a git repository", "git repository unavailable"),
|
||||
];
|
||||
|
||||
for (stderr, expected) in cases {
|
||||
let hint = classify_exec_failure(stderr).expect(stderr);
|
||||
assert!(
|
||||
hint.contains(expected),
|
||||
"expected {hint:?} to contain {expected:?}"
|
||||
);
|
||||
}
|
||||
assert_eq!(classify_exec_failure("weird new git error"), None);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -8,6 +8,9 @@ mod clone_source;
|
|||
|
||||
pub mod read_guard;
|
||||
|
||||
#[cfg(any(feature = "docker", feature = "daytona", test))]
|
||||
pub(crate) mod redact;
|
||||
|
||||
pub mod reconnect;
|
||||
|
||||
pub mod sandbox_provider;
|
||||
|
|
|
|||
|
|
@ -491,14 +491,17 @@ impl Sandbox for LocalSandbox {
|
|||
result
|
||||
}
|
||||
|
||||
async fn git_push_ref(&self, refspec: &str) -> bool {
|
||||
let has_origin = matches!(
|
||||
self.exec_command("git remote get-url origin", 10_000, None, None, None)
|
||||
.await,
|
||||
Ok(result) if result.exit_code == 0
|
||||
);
|
||||
async fn git_push_ref(&self, refspec: &str) -> crate::Result<()> {
|
||||
let has_origin = match self
|
||||
.exec_command("git remote get-url origin", 10_000, None, None, None)
|
||||
.await
|
||||
{
|
||||
Ok(result) if result.exit_code == 0 => true,
|
||||
Ok(_) => false,
|
||||
Err(err) => return Err(crate::Error::context("git remote get-url origin", err)),
|
||||
};
|
||||
if !has_origin {
|
||||
return true;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
crate::git_push_via_exec(self, refspec).await
|
||||
|
|
|
|||
9
lib/crates/fabro-sandbox/src/redact.rs
Normal file
9
lib/crates/fabro-sandbox/src/redact.rs
Normal file
|
|
@ -0,0 +1,9 @@
|
|||
pub(crate) fn redact_auth_url(
|
||||
text: &str,
|
||||
auth_url: Option<&fabro_redact::DisplaySafeUrl>,
|
||||
) -> String {
|
||||
let Some(auth_url) = auth_url else {
|
||||
return text.to_string();
|
||||
};
|
||||
text.replace(&auth_url.raw_string(), &auth_url.redacted_string())
|
||||
}
|
||||
|
|
@ -145,7 +145,7 @@ macro_rules! delegate_sandbox {
|
|||
self.$field.resume_setup_commands(run_branch)
|
||||
}
|
||||
|
||||
async fn git_push_ref(&self, refspec: &str) -> bool {
|
||||
async fn git_push_ref(&self, refspec: &str) -> $crate::Result<()> {
|
||||
self.$field.git_push_ref(refspec).await
|
||||
}
|
||||
|
||||
|
|
@ -380,6 +380,48 @@ pub struct ExecResult {
|
|||
pub duration_ms: u64,
|
||||
}
|
||||
|
||||
impl ExecResult {
|
||||
pub fn is_success(&self) -> bool {
|
||||
self.exit_code == 0 && !self.timed_out
|
||||
}
|
||||
|
||||
pub fn into_exec_error(self, label: impl Into<String>) -> crate::Error {
|
||||
crate::Error::exec(
|
||||
label,
|
||||
self.exit_code,
|
||||
self.timed_out,
|
||||
self.duration_ms,
|
||||
self.stderr,
|
||||
self.stdout,
|
||||
)
|
||||
}
|
||||
|
||||
pub fn into_exec_error_with_redactor(
|
||||
self,
|
||||
label: impl Into<String>,
|
||||
redactor: impl Fn(&str) -> String,
|
||||
) -> crate::Error {
|
||||
let stderr = redactor(&self.stderr);
|
||||
let stdout = redactor(&self.stdout);
|
||||
crate::Error::exec(
|
||||
label,
|
||||
self.exit_code,
|
||||
self.timed_out,
|
||||
self.duration_ms,
|
||||
stderr,
|
||||
stdout,
|
||||
)
|
||||
}
|
||||
|
||||
pub fn into_result(self, label: impl Into<String>) -> crate::Result<Self> {
|
||||
if self.is_success() {
|
||||
Ok(self)
|
||||
} else {
|
||||
Err(self.into_exec_error(label))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct DirEntry {
|
||||
pub name: String,
|
||||
|
|
@ -478,8 +520,10 @@ pub trait Sandbox: Send + Sync {
|
|||
}
|
||||
|
||||
/// Push a full refspec to origin from inside the sandbox.
|
||||
async fn git_push_ref(&self, _refspec: &str) -> bool {
|
||||
false
|
||||
async fn git_push_ref(&self, _refspec: &str) -> crate::Result<()> {
|
||||
Err(crate::Error::message(
|
||||
"git_push_ref not implemented for this sandbox",
|
||||
))
|
||||
}
|
||||
|
||||
/// Compute the filesystem path for a parallel branch worktree.
|
||||
|
|
@ -577,13 +621,8 @@ pub async fn setup_git_via_exec(
|
|||
let sha_result = sandbox
|
||||
.exec_command("git rev-parse HEAD", 10_000, None, None, None)
|
||||
.await
|
||||
.map_err(|e| crate::Error::message(format!("git rev-parse HEAD failed: {e}")))?;
|
||||
if sha_result.exit_code != 0 {
|
||||
return Err(crate::Error::message(format!(
|
||||
"git rev-parse HEAD failed (exit {}): {}",
|
||||
sha_result.exit_code, sha_result.stderr
|
||||
)));
|
||||
}
|
||||
.map_err(|e| crate::Error::context("git rev-parse HEAD", e))?
|
||||
.into_result("git rev-parse HEAD")?;
|
||||
(
|
||||
sha_result.stdout.trim().to_string(),
|
||||
format!("fabro/run/{run_id}"),
|
||||
|
|
@ -604,16 +643,11 @@ pub async fn setup_git_via_exec(
|
|||
shell_quote(&branch_name),
|
||||
shell_quote(&base_sha)
|
||||
);
|
||||
let checkout_result = sandbox
|
||||
sandbox
|
||||
.exec_command(&checkout_cmd, 10_000, None, None, None)
|
||||
.await
|
||||
.map_err(|e| crate::Error::message(format!("git checkout failed: {e}")))?;
|
||||
if checkout_result.exit_code != 0 {
|
||||
return Err(crate::Error::message(format!(
|
||||
"git checkout -B failed (exit {}): {}",
|
||||
checkout_result.exit_code, checkout_result.stderr
|
||||
)));
|
||||
}
|
||||
.map_err(|e| crate::Error::context("git checkout -B", e))?
|
||||
.into_result("git checkout -B")?;
|
||||
|
||||
Ok(GitRunInfo {
|
||||
base_sha,
|
||||
|
|
@ -646,11 +680,9 @@ pub(crate) async fn fetch_source_run_ref(
|
|||
.exec_command(&fetch_cmd, 30_000, None, None, None)
|
||||
.await?;
|
||||
if fetch.exit_code != 0 {
|
||||
last_error = format!(
|
||||
"git fetch source run ref failed (exit {}): {}",
|
||||
fetch.exit_code,
|
||||
fetch.stderr.trim()
|
||||
);
|
||||
last_error = fetch
|
||||
.into_exec_error("git fetch source run ref")
|
||||
.to_string();
|
||||
} else {
|
||||
let check = sandbox
|
||||
.exec_command(&check_cmd, 10_000, None, None, None)
|
||||
|
|
@ -658,11 +690,11 @@ pub(crate) async fn fetch_source_run_ref(
|
|||
if check.exit_code == 0 {
|
||||
return Ok(());
|
||||
}
|
||||
last_error = format!(
|
||||
"checkpoint {checkpoint_sha} is not reachable from {remote_ref} (exit {}): {}",
|
||||
check.exit_code,
|
||||
check.stderr.trim()
|
||||
);
|
||||
last_error = check
|
||||
.into_exec_error(format!(
|
||||
"checkpoint {checkpoint_sha} is not reachable from {remote_ref}"
|
||||
))
|
||||
.to_string();
|
||||
}
|
||||
time::sleep(Duration::from_millis(500)).await;
|
||||
}
|
||||
|
|
@ -672,93 +704,23 @@ pub(crate) async fn fetch_source_run_ref(
|
|||
|
||||
/// Helper for sandbox implementations that manage git internally.
|
||||
/// Pushes a refspec to origin via exec_command inside the sandbox.
|
||||
pub async fn git_push_via_exec(sandbox: &dyn Sandbox, refspec: &str) -> bool {
|
||||
pub async fn git_push_via_exec(sandbox: &dyn Sandbox, refspec: &str) -> crate::Result<()> {
|
||||
if let Err(e) = sandbox.refresh_push_credentials().await {
|
||||
tracing::warn!(
|
||||
refspec,
|
||||
error = %fabro_redact::redact_string(&e.to_string()),
|
||||
refspec = %refspec,
|
||||
error = %e,
|
||||
"Failed to refresh push credentials before git push"
|
||||
);
|
||||
}
|
||||
let cmd = format!("{GIT} push origin {}", shell_quote(refspec));
|
||||
match sandbox.exec_command(&cmd, 60_000, None, None, None).await {
|
||||
Ok(r) if r.exit_code == 0 => {
|
||||
tracing::info!(refspec, "Pushed git ref to origin");
|
||||
true
|
||||
}
|
||||
Ok(r) => {
|
||||
tracing::warn!(
|
||||
refspec,
|
||||
exit_code = r.exit_code,
|
||||
timed_out = r.timed_out,
|
||||
stderr = %trim_for_log(&r.stderr, GIT_LOG_TAIL_BYTES),
|
||||
stdout = %trim_for_log(&r.stdout, GIT_LOG_TAIL_BYTES),
|
||||
hint = classify_git_push_failure(&r.stderr).unwrap_or(""),
|
||||
"Failed to push git ref"
|
||||
);
|
||||
false
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
refspec,
|
||||
error = %fabro_redact::redact_string(&e.to_string()),
|
||||
"Failed to invoke git push in sandbox"
|
||||
);
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Maximum bytes of git stdout/stderr to include in a single log line.
|
||||
/// Long enough to capture the typical 1-3 line `fatal:` / `remote:` output
|
||||
/// without flooding the log when git emits a large progress dump.
|
||||
const GIT_LOG_TAIL_BYTES: usize = 2048;
|
||||
|
||||
/// Redact secrets from `text`, then keep at most the trailing `limit` bytes.
|
||||
/// Trailing because git's relevant `fatal:` / `remote: rejected` lines are
|
||||
/// emitted at the end of the output.
|
||||
fn trim_for_log(text: &str, limit: usize) -> String {
|
||||
let redacted = fabro_redact::redact_string(text);
|
||||
let trimmed = redacted.trim_end();
|
||||
if trimmed.len() <= limit {
|
||||
return trimmed.to_string();
|
||||
}
|
||||
let start = trimmed.len() - limit;
|
||||
let safe_start = (start..=trimmed.len())
|
||||
.find(|i| trimmed.is_char_boundary(*i))
|
||||
.unwrap_or(trimmed.len());
|
||||
format!("…{}", &trimmed[safe_start..])
|
||||
}
|
||||
|
||||
/// Map a git stderr to a short hint pointing at the likely cause. Returns
|
||||
/// `None` when no known pattern matches; callers should still log the raw
|
||||
/// (redacted) stderr so unknown failures stay debuggable.
|
||||
fn classify_git_push_failure(stderr: &str) -> Option<&'static str> {
|
||||
let lower = stderr.to_ascii_lowercase();
|
||||
if lower.contains("could not read username") || lower.contains("terminal prompts disabled") {
|
||||
Some(
|
||||
"no credentials in origin URL — check that the sandbox forwarded \
|
||||
GITHUB_APP_PRIVATE_KEY (or GITHUB_TOKEN) and that refresh_push_credentials succeeded",
|
||||
)
|
||||
} else if lower.contains("permission to") && lower.contains("denied") {
|
||||
Some(
|
||||
"github denied the push — installation token lacks contents:write \
|
||||
on this repo, or a branch protection / push ruleset is rejecting the ref",
|
||||
)
|
||||
} else if lower.contains("protected branch")
|
||||
|| lower.contains("ruleset")
|
||||
|| lower.contains("rejected")
|
||||
{
|
||||
Some("github rejected the ref — likely a branch protection rule or push ruleset")
|
||||
} else if lower.contains("authentication failed") || lower.contains("invalid username") {
|
||||
Some("github authentication failed — installation token may be expired or wrong scope")
|
||||
} else if lower.contains("could not resolve host") || lower.contains("network is unreachable") {
|
||||
Some("network failure inside sandbox — check DNS / egress from the run container")
|
||||
} else if lower.contains("repository not found") {
|
||||
Some("github 404 — the App installation may not include this repo")
|
||||
} else {
|
||||
None
|
||||
}
|
||||
let label = format!("git push origin {refspec}");
|
||||
sandbox
|
||||
.exec_command(&cmd, 60_000, None, None, None)
|
||||
.await
|
||||
.map_err(|e| crate::Error::context(label.clone(), e))?
|
||||
.into_result(&label)?;
|
||||
tracing::info!(refspec = %refspec, "Pushed git ref to origin");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
|
@ -780,57 +742,74 @@ mod tests {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn trim_for_log_keeps_short_output_intact() {
|
||||
assert_eq!(trim_for_log("fatal: nope\n", 2048), "fatal: nope");
|
||||
fn exec_result_helpers_convert_failure_to_exec_error() {
|
||||
let result = ExecResult {
|
||||
stdout: "out".into(),
|
||||
stderr: "fatal: could not read Username".into(),
|
||||
exit_code: 128,
|
||||
timed_out: false,
|
||||
duration_ms: 42,
|
||||
};
|
||||
let error = result.into_result("git push").unwrap_err();
|
||||
let crate::Error::Exec {
|
||||
label, exit_code, ..
|
||||
} = &error
|
||||
else {
|
||||
panic!("expected Error::Exec, got {error:?}");
|
||||
};
|
||||
assert_eq!(label, "git push");
|
||||
assert_eq!(*exit_code, 128);
|
||||
assert!(error.to_string().contains("no credentials in origin URL"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn trim_for_log_keeps_trailing_bytes_when_oversized() {
|
||||
let prefix = "x".repeat(3000);
|
||||
let suffix = "fatal: rejected";
|
||||
let trimmed = trim_for_log(&format!("{prefix}{suffix}"), 64);
|
||||
assert!(trimmed.starts_with('…'));
|
||||
assert!(trimmed.ends_with(suffix));
|
||||
assert!(trimmed.chars().count() <= 64 + 1);
|
||||
fn exec_result_success_honors_timeouts() {
|
||||
let success = ExecResult {
|
||||
stdout: String::new(),
|
||||
stderr: String::new(),
|
||||
exit_code: 0,
|
||||
timed_out: false,
|
||||
duration_ms: 1,
|
||||
};
|
||||
assert!(success.is_success());
|
||||
|
||||
let timeout = ExecResult {
|
||||
timed_out: true,
|
||||
..success
|
||||
};
|
||||
assert!(!timeout.is_success());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn trim_for_log_redacts_high_entropy_secrets() {
|
||||
let stderr = "fatal: unable to access \
|
||||
'https://x-access-token:ghs_xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA@github.com/owner/repo/'";
|
||||
let trimmed = trim_for_log(stderr, 2048);
|
||||
assert!(!trimmed.contains("ghs_xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA"));
|
||||
assert!(trimmed.contains("REDACTED"));
|
||||
fn exec_result_redactor_applies_to_stderr_and_stdout() {
|
||||
let result = ExecResult {
|
||||
stdout: "stdout https://token@example.com".into(),
|
||||
stderr: "stderr https://token@example.com".into(),
|
||||
exit_code: 1,
|
||||
timed_out: false,
|
||||
duration_ms: 1,
|
||||
};
|
||||
let error = result.into_exec_error_with_redactor("git set-url", |s| {
|
||||
s.replace("https://token@example.com", "https://****@example.com")
|
||||
});
|
||||
|
||||
let crate::Error::Exec { stderr, stdout, .. } = &error else {
|
||||
panic!("expected Error::Exec, got {error:?}");
|
||||
};
|
||||
assert_eq!(stderr, "stderr https://****@example.com");
|
||||
assert_eq!(stdout, "stdout https://****@example.com");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_git_push_failure_recognises_missing_credentials() {
|
||||
let hint = classify_git_push_failure(
|
||||
"fatal: could not read Username for 'https://github.com': No such device or address",
|
||||
fn sandbox_tracing_events_do_not_log_raw_command_fields() {
|
||||
let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("src");
|
||||
let mut failures = Vec::new();
|
||||
scan_for_command_tracing(&root, &mut failures);
|
||||
assert!(
|
||||
failures.is_empty(),
|
||||
"raw command/cmd tracing fields found:\n{}",
|
||||
failures.join("\n")
|
||||
);
|
||||
assert!(hint.unwrap().contains("no credentials in origin URL"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_git_push_failure_recognises_permission_denied() {
|
||||
let hint = classify_git_push_failure(
|
||||
"remote: Permission to owner/repo.git denied to fabro-app[bot].",
|
||||
);
|
||||
assert!(hint.unwrap().contains("github denied the push"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_git_push_failure_recognises_branch_protection() {
|
||||
let hint = classify_git_push_failure(
|
||||
"remote: error: GH013: Repository rule violations found for refs/heads/main\n\
|
||||
remote: - Cannot create ref 'refs/heads/fabro/run/X' due to ruleset",
|
||||
);
|
||||
assert!(hint.unwrap().contains("ruleset"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_git_push_failure_returns_none_for_unknown() {
|
||||
assert!(classify_git_push_failure("fatal: weird new git error message").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -960,4 +939,77 @@ mod tests {
|
|||
assert_eq!(shell_quote("hello"), "hello");
|
||||
assert_eq!(shell_quote("hello world"), "'hello world'");
|
||||
}
|
||||
|
||||
#[expect(
|
||||
clippy::disallowed_methods,
|
||||
reason = "unit test performs a small synchronous source scan of local Rust files"
|
||||
)]
|
||||
fn scan_for_command_tracing(path: &std::path::Path, failures: &mut Vec<String>) {
|
||||
for entry in std::fs::read_dir(path).unwrap() {
|
||||
let entry = entry.unwrap();
|
||||
let path = entry.path();
|
||||
if path.is_dir() {
|
||||
scan_for_command_tracing(&path, failures);
|
||||
continue;
|
||||
}
|
||||
if path.extension().and_then(|ext| ext.to_str()) != Some("rs") {
|
||||
continue;
|
||||
}
|
||||
let source = std::fs::read_to_string(&path).unwrap();
|
||||
for macro_name in [
|
||||
"tracing::trace!",
|
||||
"tracing::debug!",
|
||||
"tracing::info!",
|
||||
"tracing::warn!",
|
||||
"tracing::error!",
|
||||
"trace!",
|
||||
"debug!",
|
||||
"info!",
|
||||
"warn!",
|
||||
"error!",
|
||||
] {
|
||||
let mut rest = source.as_str();
|
||||
while let Some(idx) = rest.find(macro_name) {
|
||||
let start = source.len() - rest.len() + idx;
|
||||
if start > 0 && source.as_bytes()[start - 1] == b'"' {
|
||||
rest = &source[start + macro_name.len()..];
|
||||
continue;
|
||||
}
|
||||
let Some(call) = tracing_call(&source[start..]) else {
|
||||
break;
|
||||
};
|
||||
if call.contains("command,")
|
||||
|| call.contains("command =")
|
||||
|| call.contains("cmd,")
|
||||
|| call.contains("cmd =")
|
||||
{
|
||||
failures.push(format!(
|
||||
"{}: {}",
|
||||
path.display(),
|
||||
call.lines().next().unwrap_or(call)
|
||||
));
|
||||
}
|
||||
rest = &source[start + call.len()..];
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn tracing_call(source: &str) -> Option<&str> {
|
||||
let open = source.find('(')?;
|
||||
let mut depth = 0usize;
|
||||
for (idx, ch) in source.char_indices().skip(open) {
|
||||
match ch {
|
||||
'(' => depth += 1,
|
||||
')' => {
|
||||
depth = depth.saturating_sub(1);
|
||||
if depth == 0 {
|
||||
return Some(&source[..=idx]);
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -344,14 +344,17 @@ impl Sandbox for WorktreeSandbox {
|
|||
self.inner.resume_setup_commands(run_branch)
|
||||
}
|
||||
|
||||
async fn git_push_ref(&self, refspec: &str) -> bool {
|
||||
let has_origin = matches!(
|
||||
self.exec_command("git remote get-url origin", 10_000, None, None, None)
|
||||
.await,
|
||||
Ok(result) if result.exit_code == 0
|
||||
);
|
||||
async fn git_push_ref(&self, refspec: &str) -> crate::Result<()> {
|
||||
let has_origin = match self
|
||||
.exec_command("git remote get-url origin", 10_000, None, None, None)
|
||||
.await
|
||||
{
|
||||
Ok(result) if result.exit_code == 0 => true,
|
||||
Ok(_) => false,
|
||||
Err(err) => return Err(crate::Error::context("git remote get-url origin", err)),
|
||||
};
|
||||
if !has_origin {
|
||||
return true;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
crate::git_push_via_exec(self, refspec).await
|
||||
|
|
|
|||
|
|
@ -25,6 +25,7 @@ pub mod install;
|
|||
pub mod ip_allowlist;
|
||||
pub mod jwt_auth;
|
||||
pub mod manifest_validation;
|
||||
mod request_id;
|
||||
mod run_files;
|
||||
mod run_files_security;
|
||||
mod run_manifest;
|
||||
|
|
@ -40,5 +41,6 @@ pub mod web_auth;
|
|||
mod worker_token;
|
||||
|
||||
pub use error::{ApiError, Error, Result};
|
||||
pub use run_manifest::workflow_bundle_from_manifest;
|
||||
pub use server_secrets::process_env_snapshot;
|
||||
pub use startup::validate_startup;
|
||||
|
|
|
|||
418
lib/crates/fabro-server/src/request_id.rs
Normal file
418
lib/crates/fabro-server/src/request_id.rs
Normal file
|
|
@ -0,0 +1,418 @@
|
|||
use std::fmt;
|
||||
|
||||
use axum::body::{Body, HttpBody as _, to_bytes};
|
||||
use axum::extract::Request;
|
||||
use axum::http::header::HeaderName;
|
||||
use axum::http::{HeaderValue, header};
|
||||
use axum::middleware::Next;
|
||||
use axum::response::Response;
|
||||
use serde_json::{Value, json};
|
||||
use tracing::warn;
|
||||
use uuid::Uuid;
|
||||
|
||||
const BODY_REWRITE_LIMIT: usize = 1 << 20;
|
||||
const REQUEST_ID_HEADER: HeaderName = HeaderName::from_static("x-request-id");
|
||||
|
||||
#[derive(Clone, Copy, Debug)]
|
||||
pub(crate) struct RequestId(pub Uuid);
|
||||
|
||||
impl RequestId {
|
||||
pub(crate) fn render(self) -> String {
|
||||
self.0.to_string()
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for RequestId {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
self.0.fmt(f)
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn layer(mut req: Request, next: Next) -> Response {
|
||||
let request_id = RequestId(Uuid::new_v4());
|
||||
req.extensions_mut().insert(request_id);
|
||||
|
||||
let mut response = next.run(req).await;
|
||||
response.headers_mut().insert(
|
||||
REQUEST_ID_HEADER,
|
||||
HeaderValue::from_str(&request_id.render()).expect("uuid should be a valid header value"),
|
||||
);
|
||||
if response.status().as_u16() >= 400 && is_json_response(&response) {
|
||||
inject_request_id_into_body(response, request_id).await
|
||||
} else {
|
||||
response
|
||||
}
|
||||
}
|
||||
|
||||
fn is_json_response(response: &Response) -> bool {
|
||||
response
|
||||
.headers()
|
||||
.get(header::CONTENT_TYPE)
|
||||
.and_then(|value| value.to_str().ok())
|
||||
.is_some_and(|value| {
|
||||
value.split(';').next().is_some_and(|media_type| {
|
||||
media_type.trim().eq_ignore_ascii_case("application/json")
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
async fn inject_request_id_into_body(response: Response, request_id: RequestId) -> Response {
|
||||
match response.body().size_hint().upper() {
|
||||
Some(size) if size <= BODY_REWRITE_LIMIT as u64 => {}
|
||||
_ => return response,
|
||||
}
|
||||
|
||||
let (mut parts, body) = response.into_parts();
|
||||
let bytes = match to_bytes(body, BODY_REWRITE_LIMIT).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(err) => {
|
||||
warn!(
|
||||
request_id = %request_id,
|
||||
?err,
|
||||
"request_id middleware: failed to buffer response body for rewrite"
|
||||
);
|
||||
parts.headers.remove(header::CONTENT_LENGTH);
|
||||
return Response::from_parts(parts, Body::empty());
|
||||
}
|
||||
};
|
||||
|
||||
let mut value: Value = match serde_json::from_slice(&bytes) {
|
||||
Ok(value) => value,
|
||||
Err(_) => return Response::from_parts(parts, Body::from(bytes)),
|
||||
};
|
||||
|
||||
let Value::Object(object) = &mut value else {
|
||||
return Response::from_parts(parts, Body::from(bytes));
|
||||
};
|
||||
|
||||
let rendered = request_id.render();
|
||||
if let Some(errors) = object.get_mut("errors").and_then(Value::as_array_mut) {
|
||||
for error in errors {
|
||||
if let Value::Object(error) = error {
|
||||
error.insert("request_id".to_owned(), json!(rendered));
|
||||
}
|
||||
}
|
||||
}
|
||||
object.insert("request_id".to_owned(), json!(rendered));
|
||||
|
||||
let new_bytes = serde_json::to_vec(&value).unwrap_or_else(|_| bytes.to_vec());
|
||||
parts.headers.insert(
|
||||
header::CONTENT_LENGTH,
|
||||
HeaderValue::from_str(&new_bytes.len().to_string())
|
||||
.expect("content length should be a valid header value"),
|
||||
);
|
||||
Response::from_parts(parts, Body::from(new_bytes))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use axum::body::{Body, to_bytes};
|
||||
use axum::http::{Request, StatusCode, header};
|
||||
use axum::response::{IntoResponse, Response};
|
||||
use axum::routing::get;
|
||||
use axum::{Json, Router, middleware};
|
||||
use bytes::Bytes;
|
||||
use futures_util::stream;
|
||||
use serde_json::json;
|
||||
use tower::ServiceExt as _;
|
||||
use uuid::Uuid;
|
||||
|
||||
async fn ok_handler() -> impl IntoResponse {
|
||||
StatusCode::OK
|
||||
}
|
||||
|
||||
async fn api_error_handler() -> impl IntoResponse {
|
||||
(
|
||||
StatusCode::BAD_REQUEST,
|
||||
Json(json!({
|
||||
"errors": [
|
||||
{
|
||||
"status": "400",
|
||||
"title": "Bad Request",
|
||||
"detail": "invalid input"
|
||||
},
|
||||
{
|
||||
"status": "400",
|
||||
"title": "Bad Request",
|
||||
"detail": "missing field"
|
||||
}
|
||||
]
|
||||
})),
|
||||
)
|
||||
}
|
||||
|
||||
async fn legacy_error_handler() -> impl IntoResponse {
|
||||
(
|
||||
StatusCode::UNAUTHORIZED,
|
||||
Json(json!({ "error": "login required" })),
|
||||
)
|
||||
}
|
||||
|
||||
async fn stale_request_id_handler() -> impl IntoResponse {
|
||||
(
|
||||
StatusCode::BAD_REQUEST,
|
||||
Json(json!({
|
||||
"request_id": "stale-top-level",
|
||||
"errors": [
|
||||
{
|
||||
"status": "400",
|
||||
"title": "Bad Request",
|
||||
"detail": "invalid input",
|
||||
"request_id": "stale-entry"
|
||||
}
|
||||
]
|
||||
})),
|
||||
)
|
||||
}
|
||||
|
||||
async fn text_error_handler() -> impl IntoResponse {
|
||||
(
|
||||
StatusCode::BAD_REQUEST,
|
||||
[(header::CONTENT_TYPE, "text/plain")],
|
||||
"plain error",
|
||||
)
|
||||
}
|
||||
|
||||
async fn malformed_json_handler() -> impl IntoResponse {
|
||||
(
|
||||
StatusCode::BAD_REQUEST,
|
||||
[(header::CONTENT_TYPE, "application/json")],
|
||||
r#"{"error":"#,
|
||||
)
|
||||
}
|
||||
|
||||
async fn mixed_case_json_handler() -> impl IntoResponse {
|
||||
Response::builder()
|
||||
.status(StatusCode::BAD_REQUEST)
|
||||
.header(header::CONTENT_TYPE, "Application/JSON; charset=utf-8")
|
||||
.body(Body::from(r#"{"error":"mixed case"}"#))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn oversized_json_handler() -> impl IntoResponse {
|
||||
Response::builder()
|
||||
.status(StatusCode::BAD_REQUEST)
|
||||
.header(header::CONTENT_TYPE, "application/json")
|
||||
.body(Body::from(vec![b'a'; super::BODY_REWRITE_LIMIT + 1]))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn unknown_size_json_handler() -> impl IntoResponse {
|
||||
let stream = stream::once(async {
|
||||
Ok::<_, std::convert::Infallible>(Bytes::from_static(b"{\"error\":\"streamed\"}"))
|
||||
});
|
||||
Response::builder()
|
||||
.status(StatusCode::BAD_REQUEST)
|
||||
.header(header::CONTENT_TYPE, "application/json")
|
||||
.body(Body::from_stream(stream))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn send(request: Request<Body>) -> axum::response::Response {
|
||||
Router::new()
|
||||
.route("/", get(ok_handler))
|
||||
.route("/api-error", get(api_error_handler))
|
||||
.route("/legacy-error", get(legacy_error_handler))
|
||||
.route("/stale-request-id", get(stale_request_id_handler))
|
||||
.route("/text-error", get(text_error_handler))
|
||||
.route("/malformed-json", get(malformed_json_handler))
|
||||
.route("/mixed-case-json", get(mixed_case_json_handler))
|
||||
.route("/oversized-json", get(oversized_json_handler))
|
||||
.route("/unknown-size-json", get(unknown_size_json_handler))
|
||||
.layer(middleware::from_fn(super::layer))
|
||||
.oneshot(request)
|
||||
.await
|
||||
.expect("request should complete")
|
||||
}
|
||||
|
||||
fn request_id_header(response: &axum::response::Response) -> String {
|
||||
let request_id = response
|
||||
.headers()
|
||||
.get("x-request-id")
|
||||
.expect("request id header should be set")
|
||||
.to_str()
|
||||
.expect("request id should be ascii")
|
||||
.to_owned();
|
||||
Uuid::parse_str(&request_id).expect("request id should be a hyphenated uuid");
|
||||
request_id
|
||||
}
|
||||
|
||||
async fn response_bytes(response: axum::response::Response) -> Bytes {
|
||||
to_bytes(response.into_body(), usize::MAX)
|
||||
.await
|
||||
.expect("body should buffer")
|
||||
}
|
||||
|
||||
async fn response_json(response: axum::response::Response) -> serde_json::Value {
|
||||
let bytes = response_bytes(response).await;
|
||||
serde_json::from_slice(&bytes).expect("body should remain JSON")
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn sets_request_id_header_on_success_response() {
|
||||
let response = send(Request::builder().uri("/").body(Body::empty()).unwrap()).await;
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
request_id_header(&response);
|
||||
|
||||
let bytes = response_bytes(response).await;
|
||||
assert!(bytes.is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn overwrites_inbound_request_id_header() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/")
|
||||
.header("x-request-id", "GARBAGE-SHOULD-BE-IGNORED")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let request_id = request_id_header(&response);
|
||||
assert_ne!(request_id, "GARBAGE-SHOULD-BE-IGNORED");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn injects_request_id_into_api_error_body() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/api-error")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
|
||||
let request_id = request_id_header(&response);
|
||||
|
||||
let body = response_json(response).await;
|
||||
|
||||
assert_eq!(body["request_id"], request_id);
|
||||
let errors = body["errors"]
|
||||
.as_array()
|
||||
.expect("errors should be an array");
|
||||
assert_eq!(errors.len(), 2);
|
||||
assert!(errors.iter().all(|error| error["request_id"] == request_id));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn injects_request_id_into_legacy_json_error_body() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/legacy-error")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
|
||||
let request_id = request_id_header(&response);
|
||||
let body = response_json(response).await;
|
||||
|
||||
assert_eq!(body["error"], "login required");
|
||||
assert_eq!(body["request_id"], request_id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn overwrites_stale_request_ids_in_json_error_body() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/stale-request-id")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let request_id = request_id_header(&response);
|
||||
let body = response_json(response).await;
|
||||
|
||||
assert_eq!(body["request_id"], request_id);
|
||||
assert_eq!(body["errors"][0]["request_id"], request_id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn leaves_non_json_error_body_untouched() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/text-error")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
request_id_header(&response);
|
||||
assert_eq!(
|
||||
response_bytes(response).await,
|
||||
Bytes::from_static(b"plain error")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn leaves_malformed_json_error_body_untouched() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/malformed-json")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
request_id_header(&response);
|
||||
assert_eq!(
|
||||
response_bytes(response).await,
|
||||
Bytes::from_static(br#"{"error":"#)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn injects_request_id_into_mixed_case_json_error_body() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/mixed-case-json")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let request_id = request_id_header(&response);
|
||||
let body = response_json(response).await;
|
||||
|
||||
assert_eq!(body["error"], "mixed case");
|
||||
assert_eq!(body["request_id"], request_id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn leaves_oversized_json_error_body_untouched() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/oversized-json")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
request_id_header(&response);
|
||||
let bytes = response_bytes(response).await;
|
||||
assert_eq!(bytes.len(), super::BODY_REWRITE_LIMIT + 1);
|
||||
assert!(bytes.iter().all(|byte| *byte == b'a'));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn leaves_unknown_size_json_error_body_untouched() {
|
||||
let response = send(
|
||||
Request::builder()
|
||||
.uri("/unknown-size-json")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
|
||||
request_id_header(&response);
|
||||
assert_eq!(
|
||||
response_bytes(response).await,
|
||||
Bytes::from_static(b"{\"error\":\"streamed\"}")
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,5 +1,5 @@
|
|||
use std::collections::HashMap;
|
||||
use std::path::{Component, Path, PathBuf};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
|
||||
use anyhow::{Result, anyhow, bail};
|
||||
|
|
@ -30,11 +30,11 @@ use fabro_types::settings::run::{
|
|||
use fabro_types::{RunId, WorkflowSettings};
|
||||
use fabro_util::check_report::{CheckDetail, CheckReport, CheckResult, CheckSection, CheckStatus};
|
||||
use fabro_validate::Severity;
|
||||
use fabro_workflow::Error as WorkflowError;
|
||||
use fabro_workflow::operations::{CreateRunInput, ValidateInput, WorkflowInput, validate};
|
||||
use fabro_workflow::pipeline::Validated;
|
||||
use fabro_workflow::run_materialization::materialize_run;
|
||||
use fabro_workflow::workflow_bundle::{BundledWorkflow, WorkflowBundle};
|
||||
use fabro_workflow::workflow_bundle::{BundledWorkflow, ParsedWorkflowConfig, WorkflowBundle};
|
||||
use fabro_workflow::{Error as WorkflowError, ManifestPath};
|
||||
|
||||
use crate::server::AppState;
|
||||
|
||||
|
|
@ -45,7 +45,7 @@ pub(crate) struct PreparedManifest {
|
|||
pub root_source: String,
|
||||
pub run_id: Option<RunId>,
|
||||
pub settings: WorkflowSettings,
|
||||
pub target_path: PathBuf,
|
||||
pub target_path: ManifestPath,
|
||||
pub workflow_bundle: WorkflowBundle,
|
||||
pub workflow_input: BundledWorkflow,
|
||||
pub source_directory: PathBuf,
|
||||
|
|
@ -72,7 +72,8 @@ pub(crate) fn prepare_manifest(
|
|||
}
|
||||
|
||||
let cwd = PathBuf::from(&manifest.cwd);
|
||||
let target_path = PathBuf::from(&manifest.target.path);
|
||||
let target_path = ManifestPath::from_wire(&manifest.target.path)
|
||||
.ok_or_else(|| anyhow!("invalid manifest target path: {}", manifest.target.path))?;
|
||||
let workflow_bundle = workflow_bundle_from_manifest(&manifest.workflows)?;
|
||||
let workflow_input = workflow_bundle
|
||||
.workflow(&target_path)
|
||||
|
|
@ -81,7 +82,7 @@ pub(crate) fn prepare_manifest(
|
|||
let root_source = workflow_input.source.clone();
|
||||
|
||||
let args_overrides = manifest_args_overrides(manifest.args.as_ref());
|
||||
let workflow_run_layer = root_workflow_run_layer(manifest, &workflow_input)?;
|
||||
let workflow_run_layer = root_workflow_run_layer(&workflow_input)?;
|
||||
let mut workflow_settings_builder =
|
||||
WorkflowSettingsBuilder::new().server_run_defaults(manifest_run_defaults.clone());
|
||||
if let Some(run) = args_overrides.run {
|
||||
|
|
@ -181,7 +182,12 @@ pub(crate) async fn run_preflight(
|
|||
let (report, checks_ok) = build_preflight_report(state, prepared, validated).await?;
|
||||
let preflight_ok = !validated.has_errors() && checks_ok;
|
||||
Ok((
|
||||
preflight_response(validated, &prepared.target_path, &report, preflight_ok),
|
||||
preflight_response(
|
||||
validated,
|
||||
prepared.target_path.as_path(),
|
||||
&report,
|
||||
preflight_ok,
|
||||
),
|
||||
preflight_ok,
|
||||
))
|
||||
}
|
||||
|
|
@ -192,7 +198,7 @@ pub(crate) fn validate_response(
|
|||
) -> types::ValidateResponse {
|
||||
types::ValidateResponse {
|
||||
ok: !validated.has_errors(),
|
||||
workflow: workflow_summary(validated, &prepared.target_path),
|
||||
workflow: workflow_summary(validated, prepared.target_path.as_path()),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -203,35 +209,77 @@ pub(crate) fn graph_source(prepared: &PreparedManifest, direction: Option<&str>)
|
|||
)
|
||||
}
|
||||
|
||||
fn workflow_bundle_from_manifest(
|
||||
pub fn workflow_bundle_from_manifest(
|
||||
workflows: &HashMap<String, types::ManifestWorkflow>,
|
||||
) -> Result<WorkflowBundle> {
|
||||
let workflows = workflows
|
||||
.iter()
|
||||
.map(|(path, workflow)| {
|
||||
let files = workflow
|
||||
.files
|
||||
.iter()
|
||||
.map(|(key, entry)| (PathBuf::from(key), entry.content.clone()))
|
||||
.collect::<HashMap<_, _>>();
|
||||
Ok::<_, anyhow::Error>((PathBuf::from(path), BundledWorkflow {
|
||||
logical_path: PathBuf::from(path),
|
||||
source: workflow.source.clone(),
|
||||
files,
|
||||
}))
|
||||
})
|
||||
.collect::<Result<HashMap<_, _>>>()?;
|
||||
Ok(WorkflowBundle::new(workflows))
|
||||
let mut bundled = HashMap::new();
|
||||
let mut workflow_wire_keys = HashMap::new();
|
||||
|
||||
for (wire_key, workflow) in workflows {
|
||||
let path = ManifestPath::from_wire(wire_key)
|
||||
.ok_or_else(|| anyhow!("invalid manifest workflow key: {wire_key}"))?;
|
||||
if let Some(previous) = workflow_wire_keys.get(&path) {
|
||||
bail!(
|
||||
"duplicate canonical workflow key: {path} (from wire keys {previous:?} and \
|
||||
{wire_key:?})"
|
||||
);
|
||||
}
|
||||
workflow_wire_keys.insert(path.clone(), wire_key.clone());
|
||||
|
||||
let files = workflow_files_from_manifest(&workflow.files)?;
|
||||
let config = workflow
|
||||
.config
|
||||
.as_ref()
|
||||
.map(|config| {
|
||||
let path = ManifestPath::from_wire(&config.path).ok_or_else(|| {
|
||||
anyhow!("invalid manifest workflow config path: {}", config.path)
|
||||
})?;
|
||||
Ok::<_, anyhow::Error>(ParsedWorkflowConfig {
|
||||
path,
|
||||
source: config.source.clone(),
|
||||
})
|
||||
})
|
||||
.transpose()?;
|
||||
|
||||
bundled.insert(path.clone(), BundledWorkflow {
|
||||
path,
|
||||
source: workflow.source.clone(),
|
||||
config,
|
||||
files,
|
||||
});
|
||||
}
|
||||
|
||||
Ok(WorkflowBundle::new(bundled))
|
||||
}
|
||||
|
||||
fn root_workflow_run_layer(
|
||||
manifest: &types::RunManifest,
|
||||
workflow: &BundledWorkflow,
|
||||
) -> Result<RunLayer> {
|
||||
let Some(root) = manifest.workflows.get(&manifest.target.path) else {
|
||||
bail!("manifest target path is missing from workflows map");
|
||||
};
|
||||
let Some(config) = root.config.as_ref() else {
|
||||
fn workflow_files_from_manifest(
|
||||
files: &HashMap<String, types::ManifestFileEntry>,
|
||||
) -> Result<HashMap<ManifestPath, String>> {
|
||||
let mut bundled = HashMap::new();
|
||||
let mut file_wire_keys = HashMap::new();
|
||||
|
||||
for (wire_key, entry) in files {
|
||||
let path = ManifestPath::from_wire(wire_key)
|
||||
.ok_or_else(|| anyhow!("invalid manifest file key: {wire_key}"))?;
|
||||
if let Some(previous) = file_wire_keys.get(&path) {
|
||||
bail!(
|
||||
"duplicate canonical file key: {path} (from wire keys {previous:?} and \
|
||||
{wire_key:?})"
|
||||
);
|
||||
}
|
||||
if let Some(from) = entry.ref_.from.as_deref() {
|
||||
ManifestPath::from_wire(from)
|
||||
.ok_or_else(|| anyhow!("invalid manifest file ref from: {from}"))?;
|
||||
}
|
||||
file_wire_keys.insert(path.clone(), wire_key.clone());
|
||||
bundled.insert(path, entry.content.clone());
|
||||
}
|
||||
|
||||
Ok(bundled)
|
||||
}
|
||||
|
||||
fn root_workflow_run_layer(workflow: &BundledWorkflow) -> Result<RunLayer> {
|
||||
let Some(config) = workflow.config.as_ref() else {
|
||||
return Ok(RunLayer::default());
|
||||
};
|
||||
|
||||
|
|
@ -245,7 +293,7 @@ fn root_workflow_run_layer(
|
|||
.transpose()
|
||||
.map_err(|err| anyhow!("Failed to parse run config TOML: {err}"))?
|
||||
.unwrap_or_default();
|
||||
resolve_manifest_dockerfile(&mut run, Path::new(&config.path), &workflow.files)?;
|
||||
resolve_manifest_dockerfile(&mut run, &config.path, &workflow.files)?;
|
||||
Ok(run)
|
||||
}
|
||||
|
||||
|
|
@ -373,8 +421,8 @@ fn process_env_var(name: &str) -> Option<String> {
|
|||
|
||||
fn resolve_manifest_dockerfile(
|
||||
run: &mut RunLayer,
|
||||
config_path: &Path,
|
||||
files: &HashMap<PathBuf, String>,
|
||||
config_path: &ManifestPath,
|
||||
files: &HashMap<ManifestPath, String>,
|
||||
) -> Result<()> {
|
||||
let source = run
|
||||
.sandbox
|
||||
|
|
@ -389,39 +437,16 @@ fn resolve_manifest_dockerfile(
|
|||
return Ok(());
|
||||
};
|
||||
let path_owned = path.clone();
|
||||
let logical_path = normalize_logical_path(
|
||||
config_path.parent().unwrap_or_else(|| Path::new(".")),
|
||||
&path_owned,
|
||||
)
|
||||
.ok_or_else(|| anyhow!("unsupported dockerfile reference: {path_owned}"))?;
|
||||
let manifest_path = ManifestPath::from_reference(config_path.parent_or_dot(), &path_owned)
|
||||
.ok_or_else(|| anyhow!("unsupported dockerfile reference: {path_owned}"))?;
|
||||
let content = files
|
||||
.get(&logical_path)
|
||||
.get(&manifest_path)
|
||||
.cloned()
|
||||
.ok_or_else(|| anyhow!("missing bundled dockerfile: {}", logical_path.display()))?;
|
||||
.ok_or_else(|| anyhow!("missing bundled dockerfile: {manifest_path}"))?;
|
||||
*source = DaytonaDockerfileLayer::Inline(content);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn normalize_logical_path(current_dir: &Path, reference: &str) -> Option<PathBuf> {
|
||||
let path = Path::new(reference);
|
||||
if path.is_absolute() || reference.starts_with('~') {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut normalized = PathBuf::new();
|
||||
for component in current_dir.join(path).components() {
|
||||
match component {
|
||||
Component::CurDir => {}
|
||||
Component::Normal(part) => normalized.push(part),
|
||||
Component::ParentDir => {
|
||||
normalized.pop();
|
||||
}
|
||||
Component::RootDir | Component::Prefix(_) => return None,
|
||||
}
|
||||
}
|
||||
Some(normalized)
|
||||
}
|
||||
|
||||
async fn build_preflight_report(
|
||||
state: &AppState,
|
||||
prepared: &PreparedManifest,
|
||||
|
|
@ -1120,6 +1145,88 @@ mod tests {
|
|||
RunLayer::default()
|
||||
}
|
||||
|
||||
fn manifest_workflow() -> types::ManifestWorkflow {
|
||||
types::ManifestWorkflow {
|
||||
config: None,
|
||||
files: HashMap::new(),
|
||||
source: "digraph Demo { start [shape=Mdiamond] exit [shape=Msquare] start -> exit }"
|
||||
.to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
fn manifest_file(content: &str) -> types::ManifestFileEntry {
|
||||
types::ManifestFileEntry {
|
||||
content: content.to_string(),
|
||||
ref_: types::ManifestFileRef {
|
||||
from: Some("workflow.fabro".to_string()),
|
||||
original: "prompt.md".to_string(),
|
||||
type_: types::ManifestFileRefType::FileInline,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn workflow_bundle_rejects_duplicate_canonical_workflow_keys() {
|
||||
let workflows = HashMap::from([
|
||||
("bar.fabro".to_string(), manifest_workflow()),
|
||||
("./foo/../bar.fabro".to_string(), manifest_workflow()),
|
||||
]);
|
||||
|
||||
let error = workflow_bundle_from_manifest(&workflows).unwrap_err();
|
||||
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("duplicate canonical workflow key: bar.fabro")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn workflow_bundle_rejects_duplicate_canonical_file_keys() {
|
||||
let mut workflow = manifest_workflow();
|
||||
workflow.files = HashMap::from([
|
||||
("prompts/hello.md".to_string(), manifest_file("first")),
|
||||
("./prompts/./hello.md".to_string(), manifest_file("second")),
|
||||
]);
|
||||
let workflows = HashMap::from([("workflow.fabro".to_string(), workflow)]);
|
||||
|
||||
let error = workflow_bundle_from_manifest(&workflows).unwrap_err();
|
||||
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("duplicate canonical file key: prompts/hello.md")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn workflow_bundle_rejects_invalid_workflow_key() {
|
||||
let workflows = HashMap::from([("/abs/path.fabro".to_string(), manifest_workflow())]);
|
||||
|
||||
let error = workflow_bundle_from_manifest(&workflows).unwrap_err();
|
||||
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("invalid manifest workflow key: /abs/path.fabro")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn workflow_bundle_rejects_invalid_file_key() {
|
||||
let mut workflow = manifest_workflow();
|
||||
workflow.files = HashMap::from([("~/foo.md".to_string(), manifest_file("content"))]);
|
||||
let workflows = HashMap::from([("workflow.fabro".to_string(), workflow)]);
|
||||
|
||||
let error = workflow_bundle_from_manifest(&workflows).unwrap_err();
|
||||
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("invalid manifest file key: ~/foo.md")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn prepare_manifest_preserves_explicit_manifest_dry_run() {
|
||||
let server_settings = manifest_run_defaults(Some(&server_settings_fixture(
|
||||
|
|
|
|||
|
|
@ -120,6 +120,7 @@ use crate::github_webhooks::{
|
|||
};
|
||||
use crate::ip_allowlist::{IpAllowlistConfig, ip_allowlist_middleware};
|
||||
use crate::jwt_auth::{self, AuthMode, AuthenticatedService, AuthenticatedSubject};
|
||||
use crate::request_id::{self, RequestId};
|
||||
use crate::run_files::{FilesInFlight, list_run_files, new_files_in_flight};
|
||||
use crate::run_selector::{ResolveRunError, resolve_run_by_selector};
|
||||
use crate::server_secrets::{LlmClientResult, ServerSecrets};
|
||||
|
|
@ -1083,6 +1084,7 @@ pub fn build_router_with_options(
|
|||
))
|
||||
.layer(middleware::from_fn(security_headers::layer))
|
||||
.layer(middleware::from_fn(http_log_middleware))
|
||||
.layer(middleware::from_fn(request_id::layer))
|
||||
}
|
||||
|
||||
async fn http_log_middleware(req: axum_extract::Request, next: Next) -> Response {
|
||||
|
|
@ -1091,14 +1093,20 @@ async fn http_log_middleware(req: axum_extract::Request, next: Next) -> Response
|
|||
}
|
||||
let method = req.method().clone();
|
||||
let path = req.uri().path().to_string();
|
||||
let request_id = req
|
||||
.extensions()
|
||||
.get::<RequestId>()
|
||||
.copied()
|
||||
.map(RequestId::render)
|
||||
.unwrap_or_default();
|
||||
let start = std::time::Instant::now();
|
||||
let response = next.run(req).await;
|
||||
let status = response.status().as_u16();
|
||||
let latency_ms = start.elapsed().as_millis();
|
||||
if status >= 500 {
|
||||
error!(%method, %path, status, latency_ms, "HTTP response");
|
||||
error!(%method, %path, status, latency_ms, request_id = %request_id, "HTTP response");
|
||||
} else {
|
||||
info!(%method, %path, status, latency_ms, "HTTP response");
|
||||
info!(%method, %path, status, latency_ms, request_id = %request_id, "HTTP response");
|
||||
}
|
||||
response
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,7 +4,9 @@
|
|||
)]
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::path::{Component, Path, PathBuf};
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
use crate::ManifestPath;
|
||||
|
||||
pub trait FileResolver: Send + Sync {
|
||||
fn resolve(&self, current_dir: &Path, reference: &str) -> Option<ResolvedFile>;
|
||||
|
|
@ -12,52 +14,33 @@ pub trait FileResolver: Send + Sync {
|
|||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub struct ResolvedFile {
|
||||
pub logical_path: PathBuf,
|
||||
pub content: String,
|
||||
pub path: PathBuf,
|
||||
pub content: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default)]
|
||||
pub struct BundleFileResolver {
|
||||
files: HashMap<PathBuf, String>,
|
||||
files: HashMap<ManifestPath, String>,
|
||||
}
|
||||
|
||||
impl BundleFileResolver {
|
||||
#[must_use]
|
||||
pub fn new(files: HashMap<PathBuf, String>) -> Self {
|
||||
pub fn new(files: HashMap<ManifestPath, String>) -> Self {
|
||||
Self { files }
|
||||
}
|
||||
}
|
||||
|
||||
impl FileResolver for BundleFileResolver {
|
||||
fn resolve(&self, current_dir: &Path, reference: &str) -> Option<ResolvedFile> {
|
||||
let logical_path = normalize_logical_path(current_dir, reference)?;
|
||||
self.files.get(&logical_path).map(|content| ResolvedFile {
|
||||
logical_path,
|
||||
content: content.clone(),
|
||||
let path = ManifestPath::from_reference(current_dir, reference)?;
|
||||
let content = self.files.get(&path)?.clone();
|
||||
Some(ResolvedFile {
|
||||
path: path.into(),
|
||||
content,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn normalize_logical_path(current_dir: &Path, reference: &str) -> Option<PathBuf> {
|
||||
let path = Path::new(reference);
|
||||
if path.is_absolute() || reference.starts_with('~') {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut normalized = PathBuf::new();
|
||||
for component in current_dir.join(path).components() {
|
||||
match component {
|
||||
Component::CurDir => {}
|
||||
Component::Normal(part) => normalized.push(part),
|
||||
Component::ParentDir => {
|
||||
normalized.pop();
|
||||
}
|
||||
Component::RootDir | Component::Prefix(_) => return None,
|
||||
}
|
||||
}
|
||||
Some(normalized)
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default)]
|
||||
pub struct FilesystemFileResolver {
|
||||
fallback_dir: Option<PathBuf>,
|
||||
|
|
@ -97,7 +80,7 @@ impl FileResolver for FilesystemFileResolver {
|
|||
|
||||
match std::fs::read_to_string(&resolved_path) {
|
||||
Ok(content) => Some(ResolvedFile {
|
||||
logical_path: resolved_path,
|
||||
path: resolved_path,
|
||||
content,
|
||||
}),
|
||||
Err(error) => {
|
||||
|
|
@ -116,10 +99,14 @@ impl FileResolver for FilesystemFileResolver {
|
|||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn manifest_path(value: &str) -> ManifestPath {
|
||||
ManifestPath::from_wire(value).expect("path should parse")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bundle_resolver_returns_exact_match() {
|
||||
let resolver = BundleFileResolver::new(HashMap::from([(
|
||||
PathBuf::from("prompts/review.md"),
|
||||
manifest_path("prompts/review.md"),
|
||||
"check it".to_string(),
|
||||
)]));
|
||||
|
||||
|
|
@ -127,14 +114,14 @@ mod tests {
|
|||
.resolve(Path::new("."), "prompts/review.md")
|
||||
.expect("file should resolve");
|
||||
|
||||
assert_eq!(resolved.logical_path, PathBuf::from("prompts/review.md"));
|
||||
assert_eq!(resolved.path, PathBuf::from("prompts/review.md"));
|
||||
assert_eq!(resolved.content, "check it");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bundle_resolver_normalizes_relative_segments() {
|
||||
let resolver = BundleFileResolver::new(HashMap::from([(
|
||||
PathBuf::from("prompts/review.md"),
|
||||
manifest_path("prompts/review.md"),
|
||||
"check it".to_string(),
|
||||
)]));
|
||||
|
||||
|
|
@ -142,7 +129,7 @@ mod tests {
|
|||
.resolve(Path::new("subflows"), "../prompts/review.md")
|
||||
.expect("file should resolve");
|
||||
|
||||
assert_eq!(resolved.logical_path, PathBuf::from("prompts/review.md"));
|
||||
assert_eq!(resolved.path, PathBuf::from("prompts/review.md"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -150,4 +137,18 @@ mod tests {
|
|||
let resolver = BundleFileResolver::new(HashMap::new());
|
||||
assert!(resolver.resolve(Path::new("."), "missing.md").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bundle_resolver_resolves_outside_cwd_paths() {
|
||||
let resolver = BundleFileResolver::new(HashMap::from([(
|
||||
manifest_path("../.fabro/workflows/demo/prompts/hello.md"),
|
||||
"prompt content".to_string(),
|
||||
)]));
|
||||
|
||||
let resolved = resolver
|
||||
.resolve(Path::new("../.fabro/workflows/demo"), "prompts/hello.md")
|
||||
.expect("file should resolve for out-of-CWD workflow");
|
||||
|
||||
assert_eq!(resolved.content, "prompt content");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -19,10 +19,10 @@ use crate::context::{Context, WorkflowContext, keys};
|
|||
use crate::error::Error;
|
||||
use crate::operations::{ValidateInput, WorkflowInput, validate};
|
||||
use crate::outcome::{Outcome, OutcomeExt, StageStatus};
|
||||
use crate::pipeline;
|
||||
use crate::pipeline::types::Initialized;
|
||||
use crate::run_dir::visit_from_context;
|
||||
use crate::run_options::RunOptions;
|
||||
use crate::{ManifestPath, pipeline};
|
||||
|
||||
/// Orchestrates a child workflow engine, polling for completion or stop
|
||||
/// conditions.
|
||||
|
|
@ -30,7 +30,7 @@ pub struct SubWorkflowHandler;
|
|||
|
||||
struct ParsedChildWorkflow {
|
||||
graph: Graph,
|
||||
workflow_path: Option<PathBuf>,
|
||||
workflow_path: Option<ManifestPath>,
|
||||
}
|
||||
|
||||
/// Parse a duration string like "45s", "200ms", "5m" into a Duration.
|
||||
|
|
@ -109,9 +109,8 @@ fn parse_child_graph(node: &Node, services: &EngineServices) -> Result<ParsedChi
|
|||
(None, _) => WorkflowInput::Path(PathBuf::from(path)),
|
||||
};
|
||||
let workflow_path = match &workflow {
|
||||
WorkflowInput::Bundled(workflow) => Some(workflow.logical_path.clone()),
|
||||
WorkflowInput::Path(path) => Some(path.clone()),
|
||||
WorkflowInput::DotSource { .. } => None,
|
||||
WorkflowInput::Bundled(workflow) => Some(workflow.path.clone()),
|
||||
WorkflowInput::Path(_) | WorkflowInput::DotSource { .. } => None,
|
||||
};
|
||||
let validated = validate(ValidateInput {
|
||||
workflow,
|
||||
|
|
@ -592,13 +591,14 @@ mod tests {
|
|||
);
|
||||
|
||||
let mut services = make_services();
|
||||
services.workflow_path = Some(PathBuf::from("workflow.fabro"));
|
||||
services.workflow_path = Some(ManifestPath::from_wire("workflow.fabro").unwrap());
|
||||
services.workflow_bundle = Some(Arc::new(WorkflowBundle::new(HashMap::from([(
|
||||
PathBuf::from("children/review.fabro"),
|
||||
ManifestPath::from_wire("children/review.fabro").unwrap(),
|
||||
BundledWorkflow {
|
||||
logical_path: PathBuf::from("children/review.fabro"),
|
||||
source: child_dot_succeeds().to_string(),
|
||||
files: HashMap::new(),
|
||||
path: ManifestPath::from_wire("children/review.fabro").unwrap(),
|
||||
source: child_dot_succeeds().to_string(),
|
||||
config: None,
|
||||
files: HashMap::new(),
|
||||
},
|
||||
)]))));
|
||||
|
||||
|
|
@ -633,7 +633,7 @@ mod tests {
|
|||
);
|
||||
|
||||
let mut services = make_services();
|
||||
services.workflow_path = Some(PathBuf::from("workflow.fabro"));
|
||||
services.workflow_path = Some(ManifestPath::from_wire("workflow.fabro").unwrap());
|
||||
services.workflow_bundle = Some(Arc::new(WorkflowBundle::new(HashMap::new())));
|
||||
|
||||
let context = Context::new();
|
||||
|
|
|
|||
|
|
@ -132,6 +132,7 @@ mod hook_context;
|
|||
reason = "The lifecycle module remains crate-visible for tests and pending integrations."
|
||||
)]
|
||||
pub(crate) mod lifecycle;
|
||||
mod manifest_path;
|
||||
pub(crate) mod node_handler;
|
||||
pub mod operations;
|
||||
pub mod outcome;
|
||||
|
|
@ -145,6 +146,7 @@ pub mod run_dump;
|
|||
pub mod run_lookup;
|
||||
|
||||
pub use error::{Error, FailureCategory, FailureSignature, FailureSignatureExt, Result};
|
||||
pub use manifest_path::ManifestPath;
|
||||
pub mod run_materialization;
|
||||
pub mod run_options;
|
||||
pub mod run_status;
|
||||
|
|
|
|||
|
|
@ -175,7 +175,17 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
|
|||
.and_then(|g| g.run_branch.as_ref())
|
||||
{
|
||||
let refspec = format!("refs/heads/{branch}:refs/heads/{branch}");
|
||||
let push_ok = self.sandbox.git_push_ref(&refspec).await;
|
||||
let push_ok = match self.sandbox.git_push_ref(&refspec).await {
|
||||
Ok(()) => true,
|
||||
Err(err) => {
|
||||
tracing::warn!(
|
||||
refspec = %refspec,
|
||||
error = %err,
|
||||
"git push from run lifecycle failed"
|
||||
);
|
||||
false
|
||||
}
|
||||
};
|
||||
git_result.push_results.push((refspec, push_ok));
|
||||
}
|
||||
}
|
||||
|
|
@ -248,10 +258,10 @@ impl GitLifecycle {
|
|||
);
|
||||
match writer.write_snapshot(dump, message).await {
|
||||
Ok(snapshot) => {
|
||||
if !snapshot.pushed {
|
||||
if let Some(detail) = snapshot.push_error.as_deref() {
|
||||
self.emit_metadata_warning(
|
||||
"checkpoint_metadata_push_failed",
|
||||
format!("failed to push metadata ref refs/heads/{meta_branch}"),
|
||||
format!("failed to push metadata ref refs/heads/{meta_branch}: {detail}"),
|
||||
);
|
||||
}
|
||||
Some(snapshot.commit_sha)
|
||||
|
|
|
|||
378
lib/crates/fabro-workflow/src/manifest_path.rs
Normal file
378
lib/crates/fabro-workflow/src/manifest_path.rs
Normal file
|
|
@ -0,0 +1,378 @@
|
|||
use std::fmt;
|
||||
use std::path::{Component, Path, PathBuf};
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
/// A path used as a key inside a run manifest. It is anchored at the run's
|
||||
/// cwd, and may contain leading `..` segments for files outside that cwd.
|
||||
#[derive(Clone, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
#[serde(into = "String", try_from = "String")]
|
||||
pub struct ManifestPath(PathBuf);
|
||||
|
||||
impl ManifestPath {
|
||||
#[must_use]
|
||||
pub fn from_reference(current_dir: &Path, reference: &str) -> Option<Self> {
|
||||
let path = Path::new(reference);
|
||||
if path.is_absolute() || reference.starts_with('~') {
|
||||
return None;
|
||||
}
|
||||
|
||||
normalize_components(current_dir.join(path)).map(Self)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn from_absolute(absolute: &Path, cwd: &Path) -> Option<Self> {
|
||||
if let Ok(stripped) = absolute.strip_prefix(cwd) {
|
||||
return normalize_components(stripped).map(Self);
|
||||
}
|
||||
|
||||
relative_path_from(absolute, cwd).and_then(|path| normalize_components(path).map(Self))
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn from_wire(value: &str) -> Option<Self> {
|
||||
Self::from_reference(Path::new("."), value)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn as_path(&self) -> &Path {
|
||||
&self.0
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn parent(&self) -> Option<&Path> {
|
||||
self.0.parent()
|
||||
}
|
||||
|
||||
/// Directory that contains this path, falling back to `.` when the path
|
||||
/// has no parent component (e.g. a bare file name).
|
||||
#[must_use]
|
||||
pub fn parent_or_dot(&self) -> &Path {
|
||||
self.0.parent().unwrap_or_else(|| Path::new("."))
|
||||
}
|
||||
}
|
||||
|
||||
impl From<ManifestPath> for PathBuf {
|
||||
fn from(value: ManifestPath) -> Self {
|
||||
value.0
|
||||
}
|
||||
}
|
||||
|
||||
impl TryFrom<String> for ManifestPath {
|
||||
type Error = ManifestPathParseError;
|
||||
|
||||
fn try_from(value: String) -> Result<Self, Self::Error> {
|
||||
Self::from_wire(&value).ok_or(ManifestPathParseError(value))
|
||||
}
|
||||
}
|
||||
|
||||
impl From<ManifestPath> for String {
|
||||
fn from(value: ManifestPath) -> Self {
|
||||
value.to_string()
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for ManifestPath {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "{}", self.0.display())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct ManifestPathParseError(String);
|
||||
|
||||
impl fmt::Display for ManifestPathParseError {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "invalid ManifestPath: {}", self.0)
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for ManifestPathParseError {}
|
||||
|
||||
fn normalize_components(path: impl AsRef<Path>) -> Option<PathBuf> {
|
||||
let mut normalized = PathBuf::new();
|
||||
for component in path.as_ref().components() {
|
||||
match component {
|
||||
Component::CurDir => {}
|
||||
Component::Normal(part) => normalized.push(part),
|
||||
Component::ParentDir => {
|
||||
if normalized.file_name().is_some() {
|
||||
normalized.pop();
|
||||
} else {
|
||||
normalized.push("..");
|
||||
}
|
||||
}
|
||||
Component::RootDir | Component::Prefix(_) => return None,
|
||||
}
|
||||
}
|
||||
Some(normalized)
|
||||
}
|
||||
|
||||
fn relative_path_from(path: &Path, base: &Path) -> Option<PathBuf> {
|
||||
let path_components = path.components().collect::<Vec<_>>();
|
||||
let base_components = base.components().collect::<Vec<_>>();
|
||||
if path_components.is_empty() || base_components.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut common = 0;
|
||||
while common < path_components.len()
|
||||
&& common < base_components.len()
|
||||
&& path_components[common] == base_components[common]
|
||||
{
|
||||
common += 1;
|
||||
}
|
||||
|
||||
let mut relative = PathBuf::new();
|
||||
for component in &base_components[common..] {
|
||||
if matches!(component, Component::Normal(_)) {
|
||||
relative.push("..");
|
||||
}
|
||||
}
|
||||
for component in &path_components[common..] {
|
||||
match component {
|
||||
Component::Normal(part) => relative.push(part),
|
||||
Component::CurDir => {}
|
||||
Component::ParentDir => relative.push(".."),
|
||||
Component::RootDir | Component::Prefix(_) => return None,
|
||||
}
|
||||
}
|
||||
Some(relative)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::Path;
|
||||
|
||||
use super::*;
|
||||
|
||||
fn manifest_path(value: &str) -> ManifestPath {
|
||||
ManifestPath::from_wire(value).expect("path should parse")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_rejects_absolute_path() {
|
||||
assert!(ManifestPath::from_reference(Path::new("."), "/tmp/workflow.fabro").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_rejects_tilde_reference() {
|
||||
assert!(ManifestPath::from_reference(Path::new("."), "~/.fabro/workflow.fabro").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_simple_relative() {
|
||||
let path = ManifestPath::from_reference(Path::new("flows"), "workflow.fabro").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("flows/workflow.fabro"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_collapses_curdir_segments() {
|
||||
let path = ManifestPath::from_reference(Path::new("./flows"), "./workflow.fabro").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("flows/workflow.fabro"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_collapses_mid_path_parent_dir() {
|
||||
let path = ManifestPath::from_reference(Path::new("../foo"), "../bar/file.md").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("../bar/file.md"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_preserves_single_leading_parent_dir() {
|
||||
let path =
|
||||
ManifestPath::from_reference(Path::new("../.fabro/workflows/demo"), "prompts/hello.md")
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
path.as_path(),
|
||||
Path::new("../.fabro/workflows/demo/prompts/hello.md")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_preserves_multiple_leading_parent_dirs() {
|
||||
let path =
|
||||
ManifestPath::from_reference(Path::new("../../shared/workflows"), "file.md").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("../../shared/workflows/file.md"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_reference_collapses_then_escapes() {
|
||||
let path = ManifestPath::from_reference(Path::new("."), "foo/../..").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new(".."));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_absolute_file_inside_cwd_strips_prefix() {
|
||||
let path = ManifestPath::from_absolute(
|
||||
Path::new("/repo/project/workflow.fabro"),
|
||||
Path::new("/repo/project"),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("workflow.fabro"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_absolute_sibling_directory_uses_single_parent_dir() {
|
||||
let path = ManifestPath::from_absolute(
|
||||
Path::new("/repo/shared/workflow.fabro"),
|
||||
Path::new("/repo/project"),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("../shared/workflow.fabro"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_absolute_grandparent_uses_two_parent_dirs() {
|
||||
let path = ManifestPath::from_absolute(
|
||||
Path::new("/repo/shared/workflow.fabro"),
|
||||
Path::new("/repo/project/nested"),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("../../shared/workflow.fabro"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_absolute_user_global_workflow_from_unrelated_cwd() {
|
||||
let path = ManifestPath::from_absolute(
|
||||
Path::new("/tmp/.fabro/workflows/demo/workflow.fabro"),
|
||||
Path::new("/tmp/project"),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
path.as_path(),
|
||||
Path::new("../.fabro/workflows/demo/workflow.fabro")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_wire_accepts_canonical_relative() {
|
||||
let path = ManifestPath::from_wire("workflow.fabro").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("workflow.fabro"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_wire_rejects_absolute_path() {
|
||||
assert!(ManifestPath::from_wire("/tmp/workflow.fabro").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_wire_rejects_tilde_path() {
|
||||
assert!(ManifestPath::from_wire("~/.fabro/workflow.fabro").is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_wire_renormalizes_uncollapsed_curdir() {
|
||||
let path = ManifestPath::from_wire("./foo/./bar").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("foo/bar"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_wire_renormalizes_mid_path_parent_dir() {
|
||||
let path = ManifestPath::from_wire("foo/../bar").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("bar"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_wire_preserves_leading_parent_dir() {
|
||||
let path = ManifestPath::from_wire("../.fabro/workflows/demo/workflow.fabro").unwrap();
|
||||
|
||||
assert_eq!(
|
||||
path.as_path(),
|
||||
Path::new("../.fabro/workflows/demo/workflow.fabro")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deserialize_rejects_absolute_path() {
|
||||
assert!(serde_json::from_str::<ManifestPath>("\"/abs\"").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deserialize_rejects_tilde_path() {
|
||||
assert!(serde_json::from_str::<ManifestPath>("\"~/.fabro/x\"").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deserialize_renormalizes_non_canonical() {
|
||||
let path: ManifestPath = serde_json::from_str("\"foo/../bar\"").unwrap();
|
||||
|
||||
assert_eq!(path, manifest_path("bar"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn round_trip_file_inside_cwd() {
|
||||
let produced = ManifestPath::from_absolute(
|
||||
Path::new("/repo/project/prompts/hello.md"),
|
||||
Path::new("/repo/project"),
|
||||
)
|
||||
.unwrap();
|
||||
let consumed = ManifestPath::from_reference(Path::new("prompts"), "hello.md").unwrap();
|
||||
|
||||
assert_eq!(produced, consumed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn round_trip_user_global_workflow() {
|
||||
let produced = ManifestPath::from_absolute(
|
||||
Path::new("/tmp/.fabro/workflows/demo/prompts/hello.md"),
|
||||
Path::new("/tmp/project"),
|
||||
)
|
||||
.unwrap();
|
||||
let workflow = ManifestPath::from_absolute(
|
||||
Path::new("/tmp/.fabro/workflows/demo/workflow.fabro"),
|
||||
Path::new("/tmp/project"),
|
||||
)
|
||||
.unwrap();
|
||||
let consumed =
|
||||
ManifestPath::from_reference(workflow.parent().unwrap(), "prompts/hello.md").unwrap();
|
||||
|
||||
assert_eq!(produced, consumed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn round_trip_subworkflow_relative_reference() {
|
||||
let produced = ManifestPath::from_absolute(
|
||||
Path::new("/repo/.fabro/workflows/child/workflow.fabro"),
|
||||
Path::new("/repo"),
|
||||
)
|
||||
.unwrap();
|
||||
let root = ManifestPath::from_absolute(
|
||||
Path::new("/repo/.fabro/workflows/demo/workflow.fabro"),
|
||||
Path::new("/repo"),
|
||||
)
|
||||
.unwrap();
|
||||
let consumed =
|
||||
ManifestPath::from_reference(root.parent().unwrap(), "../child/workflow.fabro")
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(produced, consumed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn serializes_as_plain_string() {
|
||||
let serialized = serde_json::to_string(&manifest_path("../workflow.fabro")).unwrap();
|
||||
|
||||
assert_eq!(serialized, "\"../workflow.fabro\"");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deserializes_from_plain_string() {
|
||||
let path: ManifestPath = serde_json::from_str("\"workflow.fabro\"").unwrap();
|
||||
|
||||
assert_eq!(path.as_path(), Path::new("workflow.fabro"));
|
||||
}
|
||||
}
|
||||
|
|
@ -20,6 +20,7 @@ use fabro_util::json::normalize_json_value;
|
|||
use tokio::task::spawn_blocking;
|
||||
|
||||
use super::source::{ResolveWorkflowInput, WorkflowInput, resolve_workflow};
|
||||
use crate::ManifestPath;
|
||||
use crate::error::Error;
|
||||
use crate::event::{Event, append_event, to_run_event_at};
|
||||
use crate::file_resolver::FileResolver;
|
||||
|
|
@ -37,7 +38,7 @@ pub struct CreateRunInput {
|
|||
pub settings: WorkflowSettings,
|
||||
pub cwd: PathBuf,
|
||||
pub workflow_slug: Option<String>,
|
||||
pub workflow_path: Option<PathBuf>,
|
||||
pub workflow_path: Option<ManifestPath>,
|
||||
pub workflow_bundle: Option<WorkflowBundle>,
|
||||
pub submitted_manifest_bytes: Option<Vec<u8>>,
|
||||
pub run_id: Option<RunId>,
|
||||
|
|
@ -643,8 +644,8 @@ mod tests {
|
|||
fn validate_from_bundle_resolves_nested_import_files_relative_to_imported_graph() {
|
||||
let validated = validate(ValidateInput {
|
||||
workflow: WorkflowInput::Bundled(BundledWorkflow {
|
||||
logical_path: PathBuf::from("workflow.fabro"),
|
||||
source: r#"digraph Test {
|
||||
path: ManifestPath::from_wire("workflow.fabro").unwrap(),
|
||||
source: r#"digraph Test {
|
||||
graph [goal="Ship"]
|
||||
start [shape=Mdiamond]
|
||||
validate [import="./child/validate.fabro"]
|
||||
|
|
@ -652,9 +653,10 @@ mod tests {
|
|||
start -> validate -> exit
|
||||
}"#
|
||||
.to_string(),
|
||||
files: HashMap::from([
|
||||
config: None,
|
||||
files: HashMap::from([
|
||||
(
|
||||
PathBuf::from("child/validate.fabro"),
|
||||
ManifestPath::from_wire("child/validate.fabro").unwrap(),
|
||||
r#"digraph Validate {
|
||||
start [shape=Mdiamond]
|
||||
lint [prompt="@../prompts/lint.md"]
|
||||
|
|
@ -664,7 +666,7 @@ mod tests {
|
|||
.to_string(),
|
||||
),
|
||||
(
|
||||
PathBuf::from("prompts/lint.md"),
|
||||
ManifestPath::from_wire("prompts/lint.md").unwrap(),
|
||||
"Lint {{ goal }}".to_string(),
|
||||
),
|
||||
]),
|
||||
|
|
|
|||
|
|
@ -120,9 +120,9 @@ pub(crate) fn resolve_workflow(request: ResolveWorkflowInput) -> anyhow::Result<
|
|||
Ok(ResolvedWorkflow {
|
||||
raw_source: workflow.source.clone(),
|
||||
settings,
|
||||
workflow_slug: workflow_slug_from_path(&workflow.logical_path),
|
||||
workflow_slug: workflow_slug_from_path(workflow.path.as_path()),
|
||||
workflow_toml_path: None,
|
||||
dot_path: Some(workflow.logical_path.clone()),
|
||||
dot_path: Some(workflow.path.as_path().to_path_buf()),
|
||||
current_dir: Some(workflow.current_dir()),
|
||||
file_resolver: Some(workflow.file_resolver()),
|
||||
goal_override,
|
||||
|
|
|
|||
|
|
@ -31,6 +31,7 @@ use fabro_vault::Vault;
|
|||
use tokio::runtime::Handle;
|
||||
use tokio::sync::RwLock as AsyncRwLock;
|
||||
|
||||
use crate::ManifestPath;
|
||||
use crate::artifact_upload::ArtifactSink;
|
||||
use crate::context::Context;
|
||||
use crate::error::Error;
|
||||
|
|
@ -77,7 +78,7 @@ struct RunSession {
|
|||
pr_github_app: Option<fabro_github::GitHubCredentials>,
|
||||
pr_origin_url: Option<String>,
|
||||
pr_model: String,
|
||||
workflow_path: Option<PathBuf>,
|
||||
workflow_path: Option<ManifestPath>,
|
||||
workflow_bundle: Option<Arc<WorkflowBundle>>,
|
||||
run_control: Option<Arc<RunControlState>>,
|
||||
vault: Option<Arc<AsyncRwLock<Vault>>>,
|
||||
|
|
@ -1007,6 +1008,7 @@ mod tests {
|
|||
use object_store::memory::InMemory;
|
||||
|
||||
use super::*;
|
||||
use crate::ManifestPath;
|
||||
use crate::context::Context;
|
||||
use crate::event::{Emitter, EventBody};
|
||||
use crate::handler::HandlerRegistry;
|
||||
|
|
@ -1201,9 +1203,11 @@ mod tests {
|
|||
let registry = Arc::new(test_registry());
|
||||
let store = memory_store();
|
||||
let workflow_bundle = WorkflowBundle::new(HashMap::from([
|
||||
(PathBuf::from("workflow.fabro"), BundledWorkflow {
|
||||
logical_path: PathBuf::from("workflow.fabro"),
|
||||
source: r#"digraph Root {
|
||||
(
|
||||
ManifestPath::from_wire("workflow.fabro").unwrap(),
|
||||
BundledWorkflow {
|
||||
path: ManifestPath::from_wire("workflow.fabro").unwrap(),
|
||||
source: r#"digraph Root {
|
||||
graph [goal="Bundle child"]
|
||||
start [shape=Mdiamond]
|
||||
manager [
|
||||
|
|
@ -1215,19 +1219,25 @@ mod tests {
|
|||
exit [shape=Msquare]
|
||||
start -> manager -> exit
|
||||
}"#
|
||||
.to_string(),
|
||||
files: HashMap::new(),
|
||||
}),
|
||||
(PathBuf::from("children/review.fabro"), BundledWorkflow {
|
||||
logical_path: PathBuf::from("children/review.fabro"),
|
||||
source: r"digraph Review {
|
||||
.to_string(),
|
||||
config: None,
|
||||
files: HashMap::new(),
|
||||
},
|
||||
),
|
||||
(
|
||||
ManifestPath::from_wire("children/review.fabro").unwrap(),
|
||||
BundledWorkflow {
|
||||
path: ManifestPath::from_wire("children/review.fabro").unwrap(),
|
||||
source: r"digraph Review {
|
||||
start [shape=Mdiamond]
|
||||
exit [shape=Msquare]
|
||||
start -> exit
|
||||
}"
|
||||
.to_string(),
|
||||
files: HashMap::new(),
|
||||
}),
|
||||
.to_string(),
|
||||
config: None,
|
||||
files: HashMap::new(),
|
||||
},
|
||||
),
|
||||
]));
|
||||
|
||||
crate::operations::create(
|
||||
|
|
@ -1235,7 +1245,7 @@ mod tests {
|
|||
crate::operations::CreateRunInput {
|
||||
workflow: crate::operations::WorkflowInput::Bundled(
|
||||
workflow_bundle
|
||||
.workflow(Path::new("workflow.fabro"))
|
||||
.workflow(&ManifestPath::from_wire("workflow.fabro").unwrap())
|
||||
.unwrap()
|
||||
.clone(),
|
||||
),
|
||||
|
|
@ -1248,7 +1258,7 @@ mod tests {
|
|||
}),
|
||||
cwd: temp.path().to_path_buf(),
|
||||
workflow_slug: Some("bundle-child".to_string()),
|
||||
workflow_path: Some(PathBuf::from("workflow.fabro")),
|
||||
workflow_path: Some(ManifestPath::from_wire("workflow.fabro").unwrap()),
|
||||
workflow_bundle: Some(workflow_bundle),
|
||||
submitted_manifest_bytes: None,
|
||||
run_id: Some(fixtures::RUN_1),
|
||||
|
|
|
|||
|
|
@ -178,11 +178,11 @@ pub async fn write_finalize_commit(
|
|||
);
|
||||
match writer.write_snapshot(&dump, "finalize run").await {
|
||||
Ok(snapshot) => {
|
||||
if !snapshot.pushed {
|
||||
if let Some(detail) = snapshot.push_error.as_deref() {
|
||||
emit_metadata_warning(
|
||||
services,
|
||||
"checkpoint_metadata_push_failed",
|
||||
format!("failed to push metadata ref refs/heads/{meta_branch}"),
|
||||
format!("failed to push metadata ref refs/heads/{meta_branch}: {detail}"),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ use fabro_validate::{Diagnostic, Severity};
|
|||
use fabro_vault::Vault;
|
||||
use tokio::sync::RwLock as AsyncRwLock;
|
||||
|
||||
use crate::ManifestPath;
|
||||
use crate::artifact_upload::ArtifactSink;
|
||||
use crate::context::Context;
|
||||
use crate::error::Error;
|
||||
|
|
@ -241,7 +242,7 @@ pub struct InitOptions {
|
|||
pub interviewer: Arc<dyn Interviewer>,
|
||||
pub lifecycle: LifecycleOptions,
|
||||
pub run_options: RunOptions,
|
||||
pub workflow_path: Option<PathBuf>,
|
||||
pub workflow_path: Option<ManifestPath>,
|
||||
pub workflow_bundle: Option<Arc<WorkflowBundle>>,
|
||||
pub hooks: fabro_hooks::HookSettings,
|
||||
pub sandbox_env: SandboxEnvSpec,
|
||||
|
|
|
|||
|
|
@ -942,7 +942,7 @@ mod tests {
|
|||
self.inner.os_version()
|
||||
}
|
||||
|
||||
async fn git_push_ref(&self, refspec: &str) -> bool {
|
||||
async fn git_push_ref(&self, refspec: &str) -> fabro_sandbox::Result<()> {
|
||||
self.pushes.lock().unwrap().push(refspec.to_string());
|
||||
self.inner.git_push_ref(refspec).await
|
||||
}
|
||||
|
|
@ -1351,6 +1351,7 @@ mod tests {
|
|||
);
|
||||
|
||||
let snapshot = writer.write_snapshot(&dump, "checkpoint").await.unwrap();
|
||||
assert_eq!(snapshot.push_error, None);
|
||||
let commit_sha = snapshot.commit_sha;
|
||||
|
||||
let current = std::process::Command::new("git")
|
||||
|
|
@ -1447,6 +1448,64 @@ mod tests {
|
|||
assert_eq!(sandbox.pushes(), vec![refspec.clone(), refspec]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn sandbox_metadata_writer_records_log_safe_push_error() {
|
||||
let repo_dir = tempfile::tempdir().unwrap();
|
||||
let repo = repo_dir.path();
|
||||
init_git_repo(repo);
|
||||
std::fs::write(repo.join("tracked.txt"), "seed\n").unwrap();
|
||||
git_commit_all(repo, "initial");
|
||||
let missing_origin = repo_dir.path().join("missing-origin.git");
|
||||
let add_remote = std::process::Command::new("git")
|
||||
.args(["remote", "add", "origin", missing_origin.to_str().unwrap()])
|
||||
.current_dir(repo)
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(add_remote.status.success());
|
||||
|
||||
let sandbox = fabro_agent::LocalSandbox::new(repo.to_path_buf());
|
||||
let run_id = fabro_types::fixtures::RUN_2.to_string();
|
||||
let branch = crate::sandbox_metadata::metadata_branch_name(&run_id);
|
||||
let mut projection = fabro_store::RunProjection::default();
|
||||
projection.spec = Some(fabro_types::RunSpec {
|
||||
run_id: fabro_types::fixtures::RUN_2,
|
||||
settings: fabro_types::WorkflowSettings::default(),
|
||||
graph: fabro_types::Graph::new("metadata"),
|
||||
workflow_slug: Some("metadata".to_string()),
|
||||
source_directory: Some("/Users/client/project".to_string()),
|
||||
git: None,
|
||||
labels: HashMap::new(),
|
||||
provenance: None,
|
||||
manifest_blob: None,
|
||||
definition_blob: None,
|
||||
fork_source_ref: None,
|
||||
in_place: false,
|
||||
});
|
||||
let dump = crate::run_dump::RunDump::from_projection(&projection);
|
||||
let runtime = crate::sandbox_metadata::SandboxGitRuntime::new();
|
||||
let writer = crate::sandbox_metadata::SandboxMetadataWriter::new(
|
||||
&sandbox,
|
||||
&runtime,
|
||||
&run_id,
|
||||
&branch,
|
||||
crate::git::GitAuthor::default(),
|
||||
);
|
||||
|
||||
let snapshot = writer.write_snapshot(&dump, "checkpoint").await.unwrap();
|
||||
|
||||
let push_error = snapshot.push_error.unwrap();
|
||||
assert!(push_error.contains("git push origin"));
|
||||
assert!(push_error.contains("hint:"));
|
||||
assert!(
|
||||
!push_error.contains("fatal:"),
|
||||
"push error should be log-safe: {push_error}"
|
||||
);
|
||||
assert!(
|
||||
!push_error.contains(missing_origin.to_str().unwrap()),
|
||||
"push error should not include raw git stderr paths: {push_error}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn list_changed_files_raw_classifies_add_modify_delete() {
|
||||
let repo_dir = tempfile::tempdir().unwrap();
|
||||
|
|
|
|||
|
|
@ -77,7 +77,7 @@ pub(crate) struct SandboxMetadataWriter<'a> {
|
|||
|
||||
pub(crate) struct MetadataSnapshot {
|
||||
pub commit_sha: String,
|
||||
pub pushed: bool,
|
||||
pub push_error: Option<String>,
|
||||
}
|
||||
|
||||
impl<'a> SandboxMetadataWriter<'a> {
|
||||
|
|
@ -178,10 +178,15 @@ impl<'a> SandboxMetadataWriter<'a> {
|
|||
.await?;
|
||||
let commit = parse_fast_import_mark(&stdout)?;
|
||||
let refspec = format!("{full_ref}:{full_ref}");
|
||||
let pushed = self.sandbox.git_push_ref(&refspec).await;
|
||||
let push_error = self
|
||||
.sandbox
|
||||
.git_push_ref(&refspec)
|
||||
.await
|
||||
.err()
|
||||
.map(|err| err.to_string());
|
||||
Ok(MetadataSnapshot {
|
||||
commit_sha: commit,
|
||||
pushed,
|
||||
push_error,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
use std::collections::HashMap;
|
||||
#[cfg(test)]
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
|
@ -13,6 +14,7 @@ use fabro_model::Provider;
|
|||
use tokio::time;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
use crate::ManifestPath;
|
||||
use crate::event::Emitter;
|
||||
use crate::handler::HandlerRegistry;
|
||||
use crate::runtime_store::RunStoreHandle;
|
||||
|
|
@ -125,8 +127,8 @@ pub struct EngineServices {
|
|||
/// When true, handlers should skip real execution and return simulated
|
||||
/// results.
|
||||
pub dry_run: bool,
|
||||
/// Logical path of the current workflow when running from a bundle.
|
||||
pub workflow_path: Option<PathBuf>,
|
||||
/// Manifest path of the current workflow when running from a bundle.
|
||||
pub workflow_path: Option<ManifestPath>,
|
||||
/// Bundled workflows available for child-workflow resolution.
|
||||
pub workflow_bundle: Option<Arc<WorkflowBundle>>,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -110,10 +110,10 @@ impl ImportTransform {
|
|||
return Ok(());
|
||||
};
|
||||
|
||||
if import_stack.contains(&resolved_file.logical_path) {
|
||||
if import_stack.contains(&resolved_file.path) {
|
||||
let cycle = import_stack
|
||||
.iter()
|
||||
.chain(std::iter::once(&resolved_file.logical_path))
|
||||
.chain(std::iter::once(&resolved_file.path))
|
||||
.map(|path| path.display().to_string())
|
||||
.collect::<Vec<_>>()
|
||||
.join(" -> ");
|
||||
|
|
@ -137,7 +137,7 @@ impl ImportTransform {
|
|||
if let Err(message) = Self::splice_import(
|
||||
graph,
|
||||
placeholder_id,
|
||||
&resolved_file.logical_path,
|
||||
&resolved_file.path,
|
||||
&placeholder,
|
||||
prepared,
|
||||
) {
|
||||
|
|
@ -152,52 +152,47 @@ impl ImportTransform {
|
|||
resolved_file: &ResolvedFile,
|
||||
import_stack: &mut Vec<PathBuf>,
|
||||
) -> Result<PreparedImport, ImportPrepareError> {
|
||||
Self::with_import_stack(
|
||||
import_stack,
|
||||
resolved_file.logical_path.clone(),
|
||||
|import_stack| {
|
||||
let rendered_source = render_template(
|
||||
&resolved_file.content,
|
||||
&TemplateContext::new()
|
||||
.with_goal("{{ goal }}")
|
||||
.with_inputs(self.inputs.clone()),
|
||||
)
|
||||
.map_err(|error| ImportPrepareError::Hard(Error::Validation(error.to_string())))?;
|
||||
Self::with_import_stack(import_stack, resolved_file.path.clone(), |import_stack| {
|
||||
let rendered_source = render_template(
|
||||
&resolved_file.content,
|
||||
&TemplateContext::new()
|
||||
.with_goal("{{ goal }}")
|
||||
.with_inputs(self.inputs.clone()),
|
||||
)
|
||||
.map_err(|error| ImportPrepareError::Hard(Error::Validation(error.to_string())))?;
|
||||
|
||||
let mut graph = parser::parse(&rendered_source).map_err(|error| {
|
||||
ImportPrepareError::Soft(format!(
|
||||
"failed to parse {}: {error}",
|
||||
resolved_file.logical_path.display()
|
||||
))
|
||||
})?;
|
||||
let mut graph = parser::parse(&rendered_source).map_err(|error| {
|
||||
ImportPrepareError::Soft(format!(
|
||||
"failed to parse {}: {error}",
|
||||
resolved_file.path.display()
|
||||
))
|
||||
})?;
|
||||
|
||||
let import_base_dir = resolved_file
|
||||
.logical_path
|
||||
.parent()
|
||||
.map_or_else(|| PathBuf::from("."), Path::to_path_buf);
|
||||
graph =
|
||||
FileInliningTransform::new(import_base_dir.clone(), Arc::clone(&self.resolver))
|
||||
.apply(graph)
|
||||
.map_err(ImportPrepareError::Hard)?;
|
||||
let import_base_dir = resolved_file
|
||||
.path
|
||||
.parent()
|
||||
.map_or_else(|| PathBuf::from("."), Path::to_path_buf);
|
||||
graph = FileInliningTransform::new(import_base_dir.clone(), Arc::clone(&self.resolver))
|
||||
.apply(graph)
|
||||
.map_err(ImportPrepareError::Hard)?;
|
||||
|
||||
if let Some(message) = Self::unresolved_imported_prompt_error(&graph) {
|
||||
return Err(ImportPrepareError::Soft(message));
|
||||
}
|
||||
if let Some(message) = Self::unresolved_imported_prompt_error(&graph) {
|
||||
return Err(ImportPrepareError::Soft(message));
|
||||
}
|
||||
|
||||
let nested_imports = Self::collect_import_nodes(&graph);
|
||||
for (placeholder_id, import_path) in nested_imports {
|
||||
self.expand_import(
|
||||
&mut graph,
|
||||
&placeholder_id,
|
||||
&import_path,
|
||||
&import_base_dir,
|
||||
import_stack,
|
||||
)?;
|
||||
}
|
||||
let nested_imports = Self::collect_import_nodes(&graph);
|
||||
for (placeholder_id, import_path) in nested_imports {
|
||||
self.expand_import(
|
||||
&mut graph,
|
||||
&placeholder_id,
|
||||
&import_path,
|
||||
&import_base_dir,
|
||||
import_stack,
|
||||
)?;
|
||||
}
|
||||
|
||||
Self::validate_imported_graph(graph).map_err(ImportPrepareError::Soft)
|
||||
},
|
||||
)
|
||||
Self::validate_imported_graph(graph).map_err(ImportPrepareError::Soft)
|
||||
})
|
||||
}
|
||||
|
||||
fn splice_import(
|
||||
|
|
|
|||
|
|
@ -1,16 +1,24 @@
|
|||
use std::collections::HashMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::file_resolver::{BundleFileResolver, FileResolver, normalize_logical_path};
|
||||
use crate::ManifestPath;
|
||||
use crate::file_resolver::{BundleFileResolver, FileResolver};
|
||||
|
||||
#[derive(Clone, Debug, Serialize, Deserialize)]
|
||||
pub struct ParsedWorkflowConfig {
|
||||
pub path: ManifestPath,
|
||||
pub source: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
|
||||
pub struct BundledWorkflow {
|
||||
pub logical_path: PathBuf,
|
||||
pub source: String,
|
||||
pub files: HashMap<PathBuf, String>,
|
||||
pub path: ManifestPath,
|
||||
pub source: String,
|
||||
pub config: Option<ParsedWorkflowConfig>,
|
||||
pub files: HashMap<ManifestPath, String>,
|
||||
}
|
||||
|
||||
impl BundledWorkflow {
|
||||
|
|
@ -21,54 +29,49 @@ impl BundledWorkflow {
|
|||
|
||||
#[must_use]
|
||||
pub fn current_dir(&self) -> PathBuf {
|
||||
self.logical_path
|
||||
.parent()
|
||||
.map_or_else(|| PathBuf::from("."), Path::to_path_buf)
|
||||
self.path.parent_or_dot().to_path_buf()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
|
||||
pub struct WorkflowBundle {
|
||||
workflows: HashMap<PathBuf, BundledWorkflow>,
|
||||
workflows: HashMap<ManifestPath, BundledWorkflow>,
|
||||
}
|
||||
|
||||
impl WorkflowBundle {
|
||||
#[must_use]
|
||||
pub fn new(workflows: HashMap<PathBuf, BundledWorkflow>) -> Self {
|
||||
pub fn new(workflows: HashMap<ManifestPath, BundledWorkflow>) -> Self {
|
||||
Self { workflows }
|
||||
}
|
||||
|
||||
pub fn workflow(&self, logical_path: &Path) -> Option<&BundledWorkflow> {
|
||||
self.workflows.get(logical_path)
|
||||
pub fn workflow(&self, path: &ManifestPath) -> Option<&BundledWorkflow> {
|
||||
self.workflows.get(path)
|
||||
}
|
||||
|
||||
pub fn resolve_child(
|
||||
&self,
|
||||
current_workflow_path: &Path,
|
||||
current_workflow_path: &ManifestPath,
|
||||
reference: &str,
|
||||
) -> Option<&BundledWorkflow> {
|
||||
let current_dir = current_workflow_path
|
||||
.parent()
|
||||
.unwrap_or_else(|| Path::new("."));
|
||||
let logical_path = normalize_logical_path(current_dir, reference)?;
|
||||
self.workflows.get(&logical_path)
|
||||
let path = ManifestPath::from_reference(current_workflow_path.parent_or_dot(), reference)?;
|
||||
self.workflows.get(&path)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn workflows(&self) -> &HashMap<PathBuf, BundledWorkflow> {
|
||||
pub fn workflows(&self) -> &HashMap<ManifestPath, BundledWorkflow> {
|
||||
&self.workflows
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Serialize, Deserialize)]
|
||||
pub struct RunDefinition {
|
||||
pub workflow_path: PathBuf,
|
||||
pub workflows: HashMap<PathBuf, BundledWorkflow>,
|
||||
pub workflow_path: ManifestPath,
|
||||
pub workflows: HashMap<ManifestPath, BundledWorkflow>,
|
||||
}
|
||||
|
||||
impl RunDefinition {
|
||||
#[must_use]
|
||||
pub fn new(workflow_path: PathBuf, bundle: WorkflowBundle) -> Self {
|
||||
pub fn new(workflow_path: ManifestPath, bundle: WorkflowBundle) -> Self {
|
||||
Self {
|
||||
workflow_path,
|
||||
workflows: bundle.workflows,
|
||||
|
|
|
|||
|
|
@ -34,5 +34,9 @@ export interface ErrorResponseEntry {
|
|||
* Optional machine-readable error code for structured client handling.
|
||||
*/
|
||||
'code'?: string;
|
||||
/**
|
||||
* Server-generated request identifier; matches the x-request-id response header.
|
||||
*/
|
||||
'request_id'?: string;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -25,6 +25,10 @@ export interface ErrorResponse {
|
|||
* List of error entries.
|
||||
*/
|
||||
'errors': Array<ErrorResponseEntry>;
|
||||
/**
|
||||
* Server-generated request identifier; matches the x-request-id response header.
|
||||
*/
|
||||
'request_id'?: string;
|
||||
/**
|
||||
* Optional list of runtime env keys that were written before an install failure. Currently populated by `POST /install/finish` failure responses only.
|
||||
*/
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue