diff --git a/crates/arc-workflows/src/artifact.rs b/crates/arc-workflows/src/artifact.rs index 21509a87a..08200d493 100644 --- a/crates/arc-workflows/src/artifact.rs +++ b/crates/arc-workflows/src/artifact.rs @@ -78,7 +78,7 @@ impl ArtifactStore { let (stored, file_path) = if is_file_backed { let base = self.base_dir.as_ref().expect("base_dir checked above"); - let artifacts_dir = base.join("artifacts"); + let artifacts_dir = base.join("artifacts").join("values"); std::fs::create_dir_all(&artifacts_dir)?; let path = artifacts_dir.join(format!("{id}.json")); std::fs::write(&path, &serialized)?; @@ -179,7 +179,7 @@ impl ArtifactStore { /// Returns `None` if no `base_dir` is configured. #[must_use] pub fn artifacts_dir(&self) -> Option { - self.base_dir.as_ref().map(|b| b.join("artifacts")) + self.base_dir.as_ref().map(|b| b.join("artifacts").join("values")) } /// Remove all artifacts. Also deletes file-backed data from disk. @@ -377,7 +377,7 @@ mod tests { assert!(info.size_bytes > FILE_BACKING_THRESHOLD); assert_eq!( info.file_path, - Some(dir.path().join("artifacts").join("big.json")) + Some(dir.path().join("artifacts").join("values").join("big.json")) ); let retrieved = store.retrieve("big").unwrap(); @@ -393,7 +393,7 @@ mod tests { let data = serde_json::json!(large_string); store.store("big", "large", data).unwrap(); - let file_path = dir.path().join("artifacts").join("big.json"); + let file_path = dir.path().join("artifacts").join("values").join("big.json"); assert!(file_path.exists()); store.remove("big"); @@ -438,6 +438,7 @@ mod tests { path, dir.path() .join("artifacts") + .join("values") .join("response.plan.json") .to_str() .unwrap() @@ -451,6 +452,7 @@ mod tests { assert!(dir .path() .join("artifacts") + .join("values") .join("response.plan.json") .exists()); } diff --git a/crates/arc-workflows/src/engine.rs b/crates/arc-workflows/src/engine.rs index 2603258eb..34faef828 100644 --- a/crates/arc-workflows/src/engine.rs +++ b/crates/arc-workflows/src/engine.rs @@ -318,14 +318,14 @@ fn write_manifest(logs_root: &Path, graph: &Graph, config: &RunConfig) -> serde_ /// Return the directory for a node's logs. /// /// First visit (`visit <= 1`): `{logs_root}/nodes/{node_id}` -/// Subsequent visits: `{logs_root}/nodes/{node_id}-attempt_{visit}` +/// Subsequent visits: `{logs_root}/nodes/{node_id}-visit_{visit}` pub fn node_dir(logs_root: &Path, node_id: &str, visit: usize) -> PathBuf { if visit <= 1 { logs_root.join("nodes").join(node_id) } else { logs_root .join("nodes") - .join(format!("{node_id}-attempt_{visit}")) + .join(format!("{node_id}-visit_{visit}")) } } @@ -896,9 +896,16 @@ impl WorkflowRunEngine { // Collect assets after handler completes (both success and error) { - let assets_dir = node_dir(logs_root, &node.id, visit) + let node_slug = if visit <= 1 { + node.id.clone() + } else { + format!("{}-visit_{visit}", node.id) + }; + let assets_dir = logs_root + .join("artifacts") .join("assets") - .join(format!("attempt_{attempt}")); + .join(&node_slug) + .join(format!("retry_{attempt}")); match asset_snapshot::collect_assets( self.services.sandbox.as_ref(), &assets_dir, @@ -3622,7 +3629,7 @@ mod tests { let root = Path::new("/tmp/logs"); assert_eq!( node_dir(root, "work", 2), - root.join("nodes").join("work-attempt_2") + root.join("nodes").join("work-visit_2") ); } @@ -3631,7 +3638,7 @@ mod tests { let root = Path::new("/tmp/logs"); assert_eq!( node_dir(root, "work", 5), - root.join("nodes").join("work-attempt_5") + root.join("nodes").join("work-visit_5") ); } diff --git a/crates/arc-workflows/tests/integration.rs b/crates/arc-workflows/tests/integration.rs index 699038d38..a97f23e1e 100644 --- a/crates/arc-workflows/tests/integration.rs +++ b/crates/arc-workflows/tests/integration.rs @@ -7575,6 +7575,7 @@ async fn large_context_values_are_offloaded_to_artifact_store() { let artifact_file = dir .path() .join("artifacts") + .join("values") .join("response.big_output.json"); assert!( artifact_file.exists(), @@ -7901,11 +7902,11 @@ async fn node_dir_uses_visit_count_on_revisit() { first.display() ); - // Second visit: nodes/gated_work-attempt_2/status.json + // Second visit: nodes/gated_work-visit_2/status.json let second = dir .path() .join("nodes") - .join("gated_work-attempt_2") + .join("gated_work-visit_2") .join("status.json"); assert!( second.exists(), @@ -11279,10 +11280,10 @@ async fn asset_collection_local_sandbox_success() { // Check that asset files were collected into the stage directory let assets_dir = logs_dir .path() - .join("nodes") - .join("create_assets") + .join("artifacts") .join("assets") - .join("attempt_1"); + .join("create_assets") + .join("retry_1"); let report_path = assets_dir.join("test-results/report.xml"); assert!( @@ -11382,10 +11383,10 @@ async fn asset_collection_local_sandbox_on_failure() { let assets_dir = logs_dir .path() - .join("nodes") - .join("create_assets") + .join("artifacts") .join("assets") - .join("attempt_1"); + .join("create_assets") + .join("retry_1"); let report_path = assets_dir.join("test-results/report.xml"); assert!( @@ -11470,10 +11471,10 @@ async fn asset_collection_docker_sandbox() { let assets_dir = logs_dir .path() - .join("nodes") - .join("create_assets") + .join("artifacts") .join("assets") - .join("attempt_1"); + .join("create_assets") + .join("retry_1"); let report_path = assets_dir.join("test-results/report.xml"); assert!(