diff --git a/docs/agents/outputs.mdx b/docs/agents/outputs.mdx index ef1927b5c..0b240e14f 100644 --- a/docs/agents/outputs.mdx +++ b/docs/agents/outputs.mdx @@ -223,34 +223,12 @@ Collected assets are written to the run's directory, organized by node and retry ~/.fabro/scratch/{run_id}/ cache/ artifacts/ - assets/ + files/ {node_slug}/ retry_1/ test-results/ screenshot.png video.webm - manifest.json -``` - -Each collection writes a `manifest.json` summarizing what was captured: - -```json -{ - "files_copied": 3, - "total_bytes": 245760, - "files_skipped": 0, - "download_errors": 0, - "hash_errors": 0, - "captured_assets": [ - { - "path": "test-results/screenshot.png", - "mime": "image/png", - "content_md5": "a1b2c3...", - "content_sha256": "d4e5f6...", - "bytes": 81920 - } - ] -} ``` ## Observability diff --git a/docs/reference/run-directory.mdx b/docs/reference/run-directory.mdx index 79d1f3f2c..ad99f7e5f 100644 --- a/docs/reference/run-directory.mdx +++ b/docs/reference/run-directory.mdx @@ -28,7 +28,7 @@ These paths are local runtime state and caches, not the canonical run record. - **`worktree/`** — When running in worktree mode, Fabro creates a Git worktree here as the working directory for agents and commands. - **`runtime/`** — Local runtime files. Today this is mainly materialized blob payloads under `runtime/blobs/`. -- **`cache/artifacts/files/`** — Captured artifact files organized by node and retry, plus a `manifest.json` for each retry directory. +- **`cache/artifacts/files/`** — Captured artifact files organized by node and retry. - **`nodes/{manager_node}_{visit}/child/`** — Nested scratch directories for manager-loop child workflows. Large durable values, event streams, checkpoints, diffs, conclusions, and retros are no longer projected into live scratch by default. Use `fabro logs`, `fabro inspect`, the API, or `fabro store dump` for those surfaces. @@ -68,7 +68,6 @@ fabro ps --filter workflow=my-workflow │ │ └── retry_1/ │ │ ├── test-results/ │ │ │ └── screenshot.png -│ │ └── manifest.json │ ├── nodes/ │ │ └── manager/ │ │ └── child/ diff --git a/lib/crates/fabro-cli/src/commands/run/output.rs b/lib/crates/fabro-cli/src/commands/run/output.rs index 24c7b041e..7c6666421 100644 --- a/lib/crates/fabro-cli/src/commands/run/output.rs +++ b/lib/crates/fabro-cli/src/commands/run/output.rs @@ -1,4 +1,4 @@ -use std::path::Path; +use std::path::{Path, PathBuf}; use std::time::Duration; use anyhow::{Context as _, Result}; @@ -10,7 +10,6 @@ use fabro_types::{ use fabro_util::check_report::{CheckDetail, CheckReport, CheckResult, CheckSection, CheckStatus}; use fabro_util::terminal::Styles; use fabro_util::text::strip_goal_decoration; -use fabro_workflow::artifact_snapshot::collect_artifact_paths; use fabro_workflow::outcome::StageStatus; use fabro_workflow::records::Conclusion; use indicatif::HumanDuration; @@ -155,7 +154,7 @@ pub(crate) async fn print_run_summary_with_client( resolve_final_output_with_client(client, run_id, checkpoint.as_ref()).await?; print_final_output(final_output.as_deref(), styles); if let Some(run_dir) = local_run_dir { - print_assets(run_dir, styles); + print_assets_with_client(client, run_id, run_dir, styles).await?; } Ok(()) } @@ -319,11 +318,33 @@ fn blob_id_from_response(response: &str) -> Option { parse_blob_ref(response).or_else(|| parse_legacy_blob_file_ref(response)) } -pub(crate) fn print_assets(run_dir: &Path, styles: &Styles) { - let run_scratch = RunScratch::new(run_dir); - let paths = collect_artifact_paths(&run_scratch.artifact_files_dir()); +async fn resolve_local_artifact_display_paths_with_client( + client: &server_client::ServerStoreClient, + run_id: &RunId, + run_dir: &Path, +) -> Result> { + let mut paths = Vec::new(); + for entry in client.list_run_artifacts(run_id).await? { + let retry = u32::try_from(entry.retry) + .context("server returned invalid negative artifact retry")?; + let path = RunScratch::new(run_dir) + .artifact_stage_dir(&entry.node_slug, retry) + .join(entry.relative_path); + paths.push(path); + } + paths.sort(); + Ok(paths) +} + +async fn print_assets_with_client( + client: &server_client::ServerStoreClient, + run_id: &RunId, + run_dir: &Path, + styles: &Styles, +) -> Result<()> { + let paths = resolve_local_artifact_display_paths_with_client(client, run_id, run_dir).await?; if paths.is_empty() { - return; + return Ok(()); } let home = dirs::home_dir(); eprintln!("\n{}", styles.bold.apply_to("=== Artifacts ===")); @@ -331,14 +352,70 @@ pub(crate) fn print_assets(run_dir: &Path, styles: &Styles) { let display = match &home { Some(home_dir) => { let home_str = home_dir.to_string_lossy(); - if let Some(rest) = path.strip_prefix(home_str.as_ref()) { + if let Some(rest) = path.to_string_lossy().strip_prefix(home_str.as_ref()) { format!("~{rest}") } else { - path.clone() + path.display().to_string() } } - None => path.clone(), + None => path.display().to_string(), }; eprintln!("{display}"); } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use httpmock::MockServer; + + #[tokio::test] + async fn local_artifact_display_paths_come_from_server_artifact_list() { + let server = MockServer::start(); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let run_dir = tempfile::tempdir().unwrap(); + + server.mock(|when, then| { + when.method("GET") + .path(format!("/api/v1/runs/{run_id}/artifacts")); + then.status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "data": [ + { + "stage_id": "test@2", + "node_slug": "test", + "retry": 2, + "relative_path": "reports/junit.xml", + "size": 123 + } + ] + }) + .to_string(), + ); + }); + + let client = crate::server_client::connect_server_target_direct(&format!( + "{}/api/v1", + server.base_url() + )) + .await + .unwrap(); + + let paths = + resolve_local_artifact_display_paths_with_client(&client, &run_id, run_dir.path()) + .await + .unwrap(); + + assert_eq!( + paths, + vec![ + RunScratch::new(run_dir.path()) + .artifact_stage_dir("test", 2) + .join("reports/junit.xml") + ] + ); + } } diff --git a/lib/crates/fabro-cli/tests/it/cmd/support.rs b/lib/crates/fabro-cli/tests/it/cmd/support.rs index 05062b4f8..4cfb7769f 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/support.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/support.rs @@ -348,7 +348,7 @@ include = ["assets/**"] let run = run_local_workflow(context, &workspace_dir, "run.toml"); assert!( run.run_dir - .join("cache/artifacts/files/retry_assets/retry_2/manifest.json") + .join("cache/artifacts/files/retry_assets/retry_2/assets/retry/report.txt") .exists(), "setup_artifact_run should materialize retry_2 assets" ); diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 4af46f9d3..6cfbdfd6f 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -40,7 +40,6 @@ use fabro_types::{ }; use fabro_util::redact::redact_jsonl_line; use fabro_util::version::FABRO_VERSION; -use fabro_workflow::artifacts as workflow_artifacts; use fabro_workflow::error::FabroError; use fabro_workflow::handler::HandlerRegistry; use futures_util::stream; @@ -2274,26 +2273,6 @@ fn payload_too_large_response(detail: impl Into) -> Response { ApiError::new(StatusCode::PAYLOAD_TOO_LARGE, detail.into()).into_response() } -#[allow(clippy::result_large_err)] -fn run_artifacts_dir(run: &fabro_types::RunRecord, run_id: &RunId) -> PathBuf { - Storage::new(run.settings.storage_dir()) - .run_scratch(run_id) - .artifact_files_dir() -} - -#[allow(clippy::result_large_err)] -fn scan_run_artifacts( - run: &fabro_types::RunRecord, - run_id: &RunId, - node_filter: Option<&str>, - retry_filter: Option, -) -> Result, Response> { - workflow_artifacts::scan_artifacts(&run_artifacts_dir(run, run_id), node_filter, retry_filter) - .map_err(|err| { - ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() - }) -} - fn octet_stream_response(bytes: Bytes) -> Response { ( StatusCode::OK, @@ -4354,47 +4333,27 @@ async fn list_run_artifacts( Ok(id) => id, Err(response) => return response, }; - let run = match load_run_record(state.as_ref(), &id).await { - Ok(run) => run, - Err(response) => return response, - }; - - if run.uses_object_backed_artifacts() { - return match state.artifact_store.list_for_run(&id).await { - Ok(entries) => Json(RunArtifactListResponse { - data: entries - .into_iter() - .map(|entry| RunArtifactEntry { - stage_id: entry.node.to_string(), - node_slug: entry.node.node_id().to_string(), - retry: entry.node.visit().cast_signed(), - relative_path: entry.filename, - size: entry.size.cast_signed(), - }) - .collect(), - }) - .into_response(), - Err(err) => { - ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() - } - }; + if let Err(response) = load_run_record(state.as_ref(), &id).await { + return response; } - match scan_run_artifacts(&run, &id, None, None) { + match state.artifact_store.list_for_run(&id).await { Ok(entries) => Json(RunArtifactListResponse { data: entries .into_iter() .map(|entry| RunArtifactEntry { - stage_id: StageId::new(entry.node_slug.clone(), entry.retry).to_string(), - node_slug: entry.node_slug, - retry: entry.retry.cast_signed(), - relative_path: entry.relative_path, + stage_id: entry.node.to_string(), + node_slug: entry.node.node_id().to_string(), + retry: entry.node.visit().cast_signed(), + relative_path: entry.filename, size: entry.size.cast_signed(), }) .collect(), }) .into_response(), - Err(response) => response, + Err(err) => { + ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() + } } } @@ -4411,35 +4370,18 @@ async fn list_stage_artifacts( Ok(stage_id) => stage_id, Err(response) => return response, }; - let run = match load_run_record(state.as_ref(), &id).await { - Ok(run) => run, - Err(response) => return response, - }; + if let Err(response) = load_run_record(state.as_ref(), &id).await { + return response; + } match state.artifact_store.list_for_node(&id, &stage_id).await { - Ok(filenames) if run.uses_object_backed_artifacts() || !filenames.is_empty() => { - Json(ArtifactListResponse { - data: filenames - .into_iter() - .map(|filename| ArtifactEntry { filename }) - .collect(), - }) - .into_response() - } - Ok(_) => { - match scan_run_artifacts(&run, &id, Some(stage_id.node_id()), Some(stage_id.visit())) { - Ok(entries) => Json(ArtifactListResponse { - data: entries - .into_iter() - .map(|entry| ArtifactEntry { - filename: entry.relative_path, - }) - .collect(), - }) - .into_response(), - Err(response) => response, - } - } + Ok(filenames) => Json(ArtifactListResponse { + data: filenames + .into_iter() + .map(|filename| ArtifactEntry { filename }) + .collect(), + }) + .into_response(), Err(err) => { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() } @@ -4866,10 +4808,9 @@ async fn get_stage_artifact( Ok(path) => path, Err(response) => return response, }; - let run = match load_run_record(state.as_ref(), &id).await { - Ok(run) => run, - Err(response) => return response, - }; + if let Err(response) = load_run_record(state.as_ref(), &id).await { + return response; + } match state .artifact_store @@ -4877,23 +4818,7 @@ async fn get_stage_artifact( .await { Ok(Some(bytes)) => octet_stream_response(bytes), - Ok(None) if run.uses_object_backed_artifacts() => { - ApiError::not_found("Artifact not found.").into_response() - } - Ok(None) => { - let artifact_path = run_artifacts_dir(&run, &id) - .join(stage_id.node_id()) - .join(format!("retry_{}", stage_id.visit())) - .join(&relative_path); - match std::fs::read(&artifact_path) { - Ok(bytes) => octet_stream_response(Bytes::from(bytes)), - Err(err) if err.kind() == std::io::ErrorKind::NotFound => { - ApiError::not_found("Artifact not found.").into_response() - } - Err(err) => ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(), - } - } + Ok(None) => ApiError::not_found("Artifact not found.").into_response(), Err(err) => { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() } @@ -5981,7 +5906,7 @@ mod tests { .body(Body::empty()) .unwrap(); - let response = app.oneshot(req).await.unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); assert_eq!(response.status(), StatusCode::NOT_FOUND); } @@ -6843,7 +6768,7 @@ mod tests { } #[tokio::test] - async fn legacy_runs_fallback_to_scratch_artifacts() { + async fn legacy_runs_do_not_fallback_to_scratch_artifacts() { let temp = tempfile::tempdir().unwrap(); let mut settings = dry_run_settings(); settings.storage_dir = Some(temp.path().join("storage")); @@ -6857,37 +6782,8 @@ mod tests { .join("code") .join("retry_2") .join("src/lib.rs"); - let retry_dir = artifact_path - .parent() - .unwrap() - .parent() - .unwrap() - .to_path_buf(); std::fs::create_dir_all(artifact_path.parent().unwrap()).unwrap(); std::fs::write(&artifact_path, "legacy scratch only").unwrap(); - std::fs::write( - retry_dir.join("manifest.json"), - serde_json::to_string( - &fabro_workflow::artifact_snapshot::ArtifactCollectionSummary { - files_copied: 1, - total_bytes: u64::try_from(b"legacy scratch only".len()).unwrap(), - files_skipped: 0, - download_errors: 0, - hash_errors: 0, - captured_assets: vec![ - fabro_workflow::artifact_snapshot::CapturedArtifactInfo { - path: "src/lib.rs".to_string(), - mime: "text/plain".to_string(), - content_md5: "0".repeat(32), - content_sha256: "0".repeat(64), - bytes: u64::try_from(b"legacy scratch only".len()).unwrap(), - }, - ], - }, - ) - .unwrap(), - ) - .unwrap(); let run_state = state .store .open_run_reader(&run_id) @@ -6903,16 +6799,6 @@ mod tests { .unwrap() .uses_object_backed_artifacts() ); - let scanned = workflow_artifacts::scan_artifacts( - &Storage::new(settings.storage_dir()) - .run_scratch(&run_id) - .artifact_files_dir(), - Some("code"), - Some(2), - ) - .unwrap(); - assert_eq!(scanned.len(), 1); - let req = Request::builder() .method("GET") .uri(api(&format!("/runs/{run_id}/stages/code@2/artifacts"))) @@ -6921,7 +6807,7 @@ mod tests { let response = app.clone().oneshot(req).await.unwrap(); assert_eq!(response.status(), StatusCode::OK); let body = body_json(response.into_body()).await; - assert_eq!(body["data"][0]["filename"], "src/lib.rs"); + assert_eq!(body["data"].as_array().unwrap().len(), 0); let req = Request::builder() .method("GET") @@ -6931,9 +6817,7 @@ mod tests { .body(Body::empty()) .unwrap(); let response = app.oneshot(req).await.unwrap(); - assert_eq!(response.status(), StatusCode::OK); - let bytes = to_bytes(response.into_body(), usize::MAX).await.unwrap(); - assert_eq!(&bytes[..], b"legacy scratch only"); + assert_eq!(response.status(), StatusCode::NOT_FOUND); } #[tokio::test] diff --git a/lib/crates/fabro-workflow/src/artifact_snapshot.rs b/lib/crates/fabro-workflow/src/artifact_snapshot.rs index ebef71da0..3bba9a506 100644 --- a/lib/crates/fabro-workflow/src/artifact_snapshot.rs +++ b/lib/crates/fabro-workflow/src/artifact_snapshot.rs @@ -275,34 +275,6 @@ fn compute_artifact_info( }) } -fn write_artifact_manifest( - artifact_capture_dir: &Path, - summary: &ArtifactCollectionSummary, -) -> Result<(), String> { - let json = serde_json::to_string_pretty(summary) - .map_err(|e| format!("failed to serialize manifest: {e}"))?; - let manifest_path = artifact_capture_dir.join("manifest.json"); - if let Some(parent) = manifest_path.parent() { - std::fs::create_dir_all(parent).map_err(|e| { - format!( - "failed to create manifest directory {}: {e}", - parent.display() - ) - })?; - } - std::fs::write(&manifest_path, json) - .map_err(|e| format!("failed to write {}: {e}", manifest_path.display()))?; - Ok(()) -} - -fn cleanup_artifact_capture_dir(artifact_capture_dir: &Path) -> Result<(), String> { - if !artifact_capture_dir.exists() { - return Ok(()); - } - std::fs::remove_dir_all(artifact_capture_dir) - .map_err(|e| format!("failed to clean up {}: {e}", artifact_capture_dir.display())) -} - /// Collect artifact files matching the configured globs that were created during this stage. pub async fn collect_artifacts( sandbox: &dyn Sandbox, @@ -365,61 +337,14 @@ pub async fn collect_artifacts( } } - // Write manifest.json - let summary = ArtifactCollectionSummary { + Ok(ArtifactCollectionSummary { files_copied, total_bytes, files_skipped, download_errors, hash_errors, captured_assets, - }; - - if files_copied > 0 { - if let Err(e) = write_artifact_manifest(artifact_capture_dir, &summary) { - let cleanup_suffix = match cleanup_artifact_capture_dir(artifact_capture_dir) { - Ok(()) => String::new(), - Err(cleanup_err) => format!("; cleanup failed: {cleanup_err}"), - }; - return Err(format!("{e}{cleanup_suffix}")); - } - } - - Ok(summary) -} - -/// Collect all artifact paths from manifest files under `{artifacts_dir}/*/retry_*/manifest.json`. -/// -/// Returns the full on-disk paths to the downloaded artifact files. -pub fn collect_artifact_paths(artifacts_dir: &Path) -> Vec { - let Ok(nodes) = std::fs::read_dir(artifacts_dir) else { - return Vec::new(); - }; - - let mut all_paths = Vec::new(); - for node_entry in nodes.flatten() { - if !node_entry.path().is_dir() { - continue; - } - let Ok(retries) = std::fs::read_dir(node_entry.path()) else { - continue; - }; - for retry_entry in retries.flatten() { - let manifest = retry_entry.path().join("manifest.json"); - let Ok(contents) = std::fs::read_to_string(&manifest) else { - continue; - }; - let Ok(summary) = serde_json::from_str::(&contents) else { - continue; - }; - let retry_dir = retry_entry.path(); - for asset in &summary.captured_assets { - let full_path = retry_dir.join(&asset.path); - all_paths.push(full_path.to_string_lossy().into_owned()); - } - } - } - all_paths + }) } #[cfg(test)] @@ -743,9 +668,9 @@ mod tests { let content = std::fs::read_to_string(&dest).unwrap(); assert_eq!(content, ""); - // Check manifest + // No manifest is written; the durable artifact record lives elsewhere. let manifest = stage_dir.path().join("manifest.json"); - assert!(manifest.exists()); + assert!(!manifest.exists()); } #[tokio::test] @@ -787,123 +712,6 @@ mod tests { assert_eq!(summary.hash_errors, 0); } - #[cfg(unix)] - #[test] - fn write_artifact_manifest_failure_cleans_up_artifact_capture_dir() { - use std::fs::Permissions; - use std::os::unix::fs::PermissionsExt; - - let parent = tempfile::tempdir().unwrap(); - let stage_dir = parent.path().join("stage"); - fs::create_dir_all(stage_dir.join("test-results")).unwrap(); - fs::write(stage_dir.join("test-results/report.xml"), "").unwrap(); - fs::set_permissions(&stage_dir, Permissions::from_mode(0o555)).unwrap(); - - let summary = ArtifactCollectionSummary { - files_copied: 1, - total_bytes: 7, - files_skipped: 0, - download_errors: 0, - hash_errors: 0, - captured_assets: vec![CapturedArtifactInfo { - path: "test-results/report.xml".to_string(), - mime: "text/xml".to_string(), - content_md5: "f1430934c390c118ed2f148e1d44d36c".to_string(), - content_sha256: "28e51ddac37391b99c2b9053f1122d0bf84b02365e6fd8c6e8667378bd00f436" - .to_string(), - bytes: 7, - }], - }; - - let err = write_artifact_manifest(&stage_dir, &summary).unwrap_err(); - assert!(err.contains("failed to write")); - - fs::set_permissions(&stage_dir, Permissions::from_mode(0o755)).unwrap(); - cleanup_artifact_capture_dir(&stage_dir).unwrap(); - assert!(!stage_dir.exists()); - } - - #[test] - fn collect_artifact_paths_from_manifests() { - let tmp = tempfile::tempdir().unwrap(); - let base = tmp.path(); - let artifacts_dir = base.join("cache/artifacts/files"); - - // Create two node directories with manifests - let node_a = artifacts_dir.join("node_a/retry_1"); - std::fs::create_dir_all(&node_a).unwrap(); - std::fs::write( - node_a.join("manifest.json"), - serde_json::to_string(&ArtifactCollectionSummary { - files_copied: 2, - total_bytes: 2048, - files_skipped: 0, - download_errors: 0, - hash_errors: 0, - captured_assets: vec![ - CapturedArtifactInfo { - path: "test-results/report.xml".to_string(), - mime: "text/xml".to_string(), - content_md5: "md5-report".to_string(), - content_sha256: "sha256-report".to_string(), - bytes: 1024, - }, - CapturedArtifactInfo { - path: "test-results/screenshot.png".to_string(), - mime: "image/png".to_string(), - content_md5: "md5-screenshot".to_string(), - content_sha256: "sha256-screenshot".to_string(), - bytes: 1024, - }, - ], - }) - .unwrap(), - ) - .unwrap(); - - let node_b = artifacts_dir.join("node_b/retry_1"); - std::fs::create_dir_all(&node_b).unwrap(); - std::fs::write( - node_b.join("manifest.json"), - serde_json::to_string(&ArtifactCollectionSummary { - files_copied: 1, - total_bytes: 512, - files_skipped: 0, - download_errors: 0, - hash_errors: 0, - captured_assets: vec![CapturedArtifactInfo { - path: "coverage/lcov.info".to_string(), - mime: "application/octet-stream".to_string(), - content_md5: "md5-lcov".to_string(), - content_sha256: "sha256-lcov".to_string(), - bytes: 512, - }], - }) - .unwrap(), - ) - .unwrap(); - - let paths = collect_artifact_paths(&artifacts_dir); - assert_eq!(paths.len(), 3); - let base_str = base.to_string_lossy(); - assert!(paths.contains(&format!( - "{base_str}/cache/artifacts/files/node_a/retry_1/test-results/report.xml" - ))); - assert!(paths.contains(&format!( - "{base_str}/cache/artifacts/files/node_a/retry_1/test-results/screenshot.png" - ))); - assert!(paths.contains(&format!( - "{base_str}/cache/artifacts/files/node_b/retry_1/coverage/lcov.info" - ))); - } - - #[test] - fn collect_asset_paths_empty_when_no_assets() { - let tmp = tempfile::tempdir().unwrap(); - let paths = collect_artifact_paths(&tmp.path().join("cache/artifacts/files")); - assert!(paths.is_empty()); - } - #[test] fn select_files_enforces_count_limit() { // Create 150 small, recent files — should be capped at MAX_FILE_COUNT (100) diff --git a/lib/crates/fabro-workflow/src/artifacts.rs b/lib/crates/fabro-workflow/src/artifacts.rs deleted file mode 100644 index 4dd83a491..000000000 --- a/lib/crates/fabro-workflow/src/artifacts.rs +++ /dev/null @@ -1,153 +0,0 @@ -use std::path::{Path, PathBuf}; - -use anyhow::Result; - -use crate::artifact_snapshot::ArtifactCollectionSummary; - -/// An individual artifact file discovered from a run's artifact manifests. -#[derive(Debug, Clone, serde::Serialize)] -pub struct ArtifactEntry { - pub node_slug: String, - pub retry: u32, - pub relative_path: String, - #[serde(serialize_with = "serialize_path")] - pub absolute_path: PathBuf, - pub size: u64, -} - -fn serialize_path(path: &Path, serializer: S) -> Result { - serializer.serialize_str(&path.display().to_string()) -} - -/// Walk `{artifacts_dir}/*/retry_*/manifest.json`, stat each file, and return entries. -pub fn scan_artifacts( - artifacts_dir: &Path, - node_filter: Option<&str>, - retry_filter: Option, -) -> Result> { - let Ok(nodes) = std::fs::read_dir(artifacts_dir) else { - return Ok(Vec::new()); - }; - - let mut entries = Vec::new(); - for node_entry in nodes.flatten() { - if !node_entry.path().is_dir() { - continue; - } - let node_slug = node_entry.file_name().to_string_lossy().into_owned(); - - if let Some(filter) = node_filter { - if node_slug != filter { - continue; - } - } - - let Ok(retries) = std::fs::read_dir(node_entry.path()) else { - continue; - }; - for retry_entry in retries.flatten() { - let retry_dir = retry_entry.path(); - let dir_name = retry_entry.file_name().to_string_lossy().into_owned(); - let retry: u32 = dir_name - .strip_prefix("retry_") - .and_then(|value| value.parse().ok()) - .unwrap_or(0); - - if let Some(filter) = retry_filter { - if retry != filter { - continue; - } - } - - let manifest = retry_dir.join("manifest.json"); - let Ok(contents) = std::fs::read_to_string(&manifest) else { - continue; - }; - let Ok(summary) = serde_json::from_str::(&contents) else { - continue; - }; - - for asset in &summary.captured_assets { - let absolute_path = retry_dir.join(&asset.path); - entries.push(ArtifactEntry { - node_slug: node_slug.clone(), - retry, - relative_path: asset.path.clone(), - absolute_path, - size: asset.bytes, - }); - } - } - } - - entries.sort_by(|left, right| { - left.node_slug - .cmp(&right.node_slug) - .then_with(|| left.retry.cmp(&right.retry)) - .then_with(|| left.relative_path.cmp(&right.relative_path)) - }); - - Ok(entries) -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::artifact_snapshot::{ArtifactCollectionSummary, CapturedArtifactInfo}; - - #[test] - fn scan_artifacts_filters_by_node_and_retry() { - let tmp = tempfile::tempdir().unwrap(); - let artifacts_dir = tmp.path().join("cache/artifacts/files"); - - let retry_1 = artifacts_dir.join("work/retry_1"); - std::fs::create_dir_all(&retry_1).unwrap(); - std::fs::write( - retry_1.join("manifest.json"), - serde_json::to_string(&ArtifactCollectionSummary { - files_copied: 1, - total_bytes: 5, - files_skipped: 0, - download_errors: 0, - hash_errors: 0, - captured_assets: vec![CapturedArtifactInfo { - path: "report.txt".to_string(), - mime: "text/plain".to_string(), - content_md5: "a".repeat(32), - content_sha256: "b".repeat(64), - bytes: 5, - }], - }) - .unwrap(), - ) - .unwrap(); - - let retry_2 = artifacts_dir.join("work/retry_2"); - std::fs::create_dir_all(&retry_2).unwrap(); - std::fs::write( - retry_2.join("manifest.json"), - serde_json::to_string(&ArtifactCollectionSummary { - files_copied: 1, - total_bytes: 6, - files_skipped: 0, - download_errors: 0, - hash_errors: 0, - captured_assets: vec![CapturedArtifactInfo { - path: "report.txt".to_string(), - mime: "text/plain".to_string(), - content_md5: "c".repeat(32), - content_sha256: "d".repeat(64), - bytes: 6, - }], - }) - .unwrap(), - ) - .unwrap(); - - let entries = scan_artifacts(&artifacts_dir, Some("work"), Some(2)).unwrap(); - assert_eq!(entries.len(), 1); - assert_eq!(entries[0].retry, 2); - assert_eq!(entries[0].relative_path, "report.txt"); - assert_eq!(entries[0].size, 6); - } -} diff --git a/lib/crates/fabro-workflow/src/lib.rs b/lib/crates/fabro-workflow/src/lib.rs index ca73b5bde..76e2311a1 100644 --- a/lib/crates/fabro-workflow/src/lib.rs +++ b/lib/crates/fabro-workflow/src/lib.rs @@ -115,7 +115,6 @@ pub fn extract_stage_durations_from_events(events: &[EventEnvelope]) -> HashMap< pub mod artifact; pub mod artifact_snapshot; pub mod artifact_upload; -pub mod artifacts; pub(crate) mod condition; pub mod context; pub mod devcontainer_bridge; diff --git a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs index 924445c2b..979c9a1d7 100644 --- a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs +++ b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs @@ -1368,7 +1368,7 @@ async fn daytona_asset_collection() { assert!(content.contains("testsuites")); let manifest_path = artifacts_dir.join("manifest.json"); - assert!(manifest_path.exists(), "manifest.json should exist"); + assert!(!manifest_path.exists(), "manifest.json should not exist"); env.cleanup().await.unwrap(); } diff --git a/lib/crates/fabro-workflow/tests/it/integration.rs b/lib/crates/fabro-workflow/tests/it/integration.rs index ef3b8af49..4adfa3932 100644 --- a/lib/crates/fabro-workflow/tests/it/integration.rs +++ b/lib/crates/fabro-workflow/tests/it/integration.rs @@ -12683,16 +12683,13 @@ async fn asset_collection_local_sandbox_success() { let report_content = std::fs::read_to_string(&report_path).unwrap(); assert!(report_content.contains("testsuites")); - // Check manifest.json was written + // Check manifest.json is no longer written let manifest_path = artifacts_dir.join("manifest.json"); assert!( - manifest_path.exists(), - "manifest.json should exist at {}", + !manifest_path.exists(), + "manifest.json should not exist at {}", manifest_path.display() ); - let manifest: serde_json::Value = - serde_json::from_str(&std::fs::read_to_string(&manifest_path).unwrap()).unwrap(); - assert!(manifest["files_copied"].as_u64().unwrap() >= 1); // Check that ArtifactCaptured events were emitted let captured_events = events.lock().unwrap(); @@ -12893,7 +12890,7 @@ async fn asset_collection_docker_sandbox() { assert!(content.contains("testsuites")); let manifest_path = artifacts_dir.join("manifest.json"); - assert!(manifest_path.exists(), "manifest.json should exist"); + assert!(!manifest_path.exists(), "manifest.json should not exist"); sandbox.cleanup().await.unwrap(); }