From 2e500b0e2fe5bd517d783cc28d8bb8bf8c4bf833 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Fri, 25 Sep 2026 13:32:13 -0400 Subject: [PATCH] Simplify artifact capture writer and storage helpers - The artifact writer takes the digest the hooks already computed, so captured bytes are hashed once on each side, and the dead integrity error goes away. - The writer is a required part of HooksSpec, not an optional field on RunRequest, so a run with capture globs always has a writer and the no-writer error goes away. - ArtifactStore routes put/get and the capture methods through shared put_at/get_at helpers. - The upload handler parses the digest with parse_blob_hash_path before any store reads, and builds its size-limit message from the constant. - The default local artifact root comes from one helper used by both config resolution and the storage-dir override. Co-Authored-By: Claude Opus 5.5 --- .../src/commands/run/petri_worker.rs | 15 +- .../src/server/handler/artifacts.rs | 35 +++-- .../fabro-server/src/server/petri_runs.rs | 8 +- .../src/server/tests/artifact_storage.rs | 7 +- lib/components/fabro-petri/src/artifacts.rs | 26 ++-- lib/components/fabro-petri/src/engine.rs | 30 ++-- lib/components/fabro-petri/src/hooks.rs | 132 ++++++------------ lib/components/fabro-petri/tests/hooks.rs | 17 ++- .../fabro-petri/tests/support/mod.rs | 1 - .../fabro-store/src/artifact_store.rs | 54 +++---- .../fabro-config/src/resolve/server.rs | 10 +- lib/foundation/fabro-types/src/dense.rs | 11 +- .../fabro-types/src/settings/server.rs | 12 ++ 13 files changed, 161 insertions(+), 197 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 9bfc188ac..d4b1ab37d 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -198,8 +198,15 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { } runner::set_worker_title(&run_id, WorkerTitlePhase::Running); - let hooks = HooksSpec::for_run(Arc::clone(&records), &worker.run_state.spec.settings.run) - .with_test_gates(test_checkpoint_gates()); + let hooks = HooksSpec::for_run( + Arc::clone(&records), + &worker.run_state.spec.settings.run, + Arc::new(ClientArtifactWriter::new( + worker.client.clone_for_reuse(), + run_id, + )), + ) + .with_test_gates(test_checkpoint_gates()); let request = RunRequest { run_id: run_id.to_string(), run_dir: worker.run_dir.join("petri"), @@ -223,10 +230,6 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { worker.client.clone_for_reuse(), run_id, ))), - artifact_writer: Some(Arc::new(ClientArtifactWriter::new( - worker.client.clone_for_reuse(), - run_id, - ))), hooks: Some(hooks), }; let paused_mirror = mirror_paused_state(run_id, &controls, Arc::clone(&records)); diff --git a/lib/apps/fabro-server/src/server/handler/artifacts.rs b/lib/apps/fabro-server/src/server/handler/artifacts.rs index 91a73675e..46ba7fab2 100644 --- a/lib/apps/fabro-server/src/server/handler/artifacts.rs +++ b/lib/apps/fabro-server/src/server/handler/artifacts.rs @@ -26,8 +26,8 @@ use super::super::{ Bytes, HashMap, IntoResponse, Json, NodeArtifact, Path, Query, RequireRunBlob, RequireRunScoped, RequiredUser, Response, Router, RunArtifactEntry, RunArtifactListResponse, RunId, StageArtifactEntry, State, StatusCode, WriteBlobResponse, get, header, - octet_stream_response, parse_run_id_path, parse_stage_id_path, post, reject_if_archived, - required_query_param, validate_relative_artifact_path, + octet_stream_response, parse_blob_hash_path, parse_run_id_path, parse_stage_id_path, post, + reject_if_archived, required_query_param, validate_relative_artifact_path, }; use crate::principal_middleware::RequireWorkerRunSegment; @@ -56,12 +56,19 @@ async fn write_run_artifact_content( State(state): State>, body: Result, ) -> Response { + let expected = match parse_blob_hash_path(&digest) { + Ok(expected) => expected, + Err(response) => return response, + }; let body = match body { Ok(body) => body, Err(error) => { return ApiError::new( error.status(), - "Artifact request body could not be read within the 10 MiB limit.", + format!( + "Artifact request body could not be read within the {} MiB limit.", + ARTIFACT_MAX_FILE_BYTES / (1024 * 1024) + ), ) .into_response(); } @@ -72,15 +79,16 @@ async fn write_run_artifact_content( if let Err(error) = state.load_run_projection(&id).await { return error.into_response(); } - let Ok(expected) = digest.parse::() else { - return ApiError::bad_request("Invalid artifact digest.").into_response(); - }; if BlobHash::new(&body) != expected { return ApiError::bad_request("Artifact content does not match its digest.") .into_response(); } - match state.artifact_store.put_capture(&id, &body).await { - Ok(_) => StatusCode::NO_CONTENT.into_response(), + match state + .artifact_store + .put_capture(&id, &expected, &body) + .await + { + Ok(()) => StatusCode::NO_CONTENT.into_response(), Err(error) => { warn!(run_id = %id, error = %collect_chain(&error).join(": "), "Artifact upload failed"); ApiError::new( @@ -572,9 +580,10 @@ mod tests { ); let blobs = state.store_ref().blobs(); let old = blobs.write(b"old SQLite").await.unwrap(); - let new = state + let new = BlobHash::new(b"new object"); + state .artifact_store - .put_capture(&run, b"new object") + .put_capture(&run, &new, b"new object") .await .unwrap(); for (path, size, source) in [ @@ -603,7 +612,7 @@ mod tests { // Unrecorded content is never a file-list entry. state .artifact_store - .put_capture(&run, b"orphan") + .put_capture(&run, &BlobHash::new(b"orphan"), b"orphan") .await .unwrap(); let entries = run_artifacts(&state, &run, &projection).await.unwrap(); @@ -652,9 +661,5 @@ mod tests { .unwrap() .is_none() ); - // Saved projections preserve both physical sources without rewriting history. - let saved = serde_json::to_value(&projection).unwrap(); - let decoded: RunProjection = serde_json::from_value(saved).unwrap(); - assert_eq!(decoded.artifacts, projection.artifacts); } } diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index be500bb5d..1055f7215 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -456,6 +456,10 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { &state.stores.run_summaries, ))), &run_state.spec.settings.run, + Arc::new(StoreArtifactWriter::new( + state.artifact_store.clone(), + run_id, + )), ); let request = RunRequest { run_id: run_id.to_string(), @@ -480,10 +484,6 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { observers, secrets: Some(Arc::new(VaultSecrets::from_vault(&vault))), blobs: Some(state.store_ref().blobs()), - artifact_writer: Some(Arc::new(StoreArtifactWriter::new( - state.artifact_store.clone(), - run_id, - ))), hooks: Some(hooks), }; let result = Box::pin(engine::run(request)).await; diff --git a/lib/apps/fabro-server/src/server/tests/artifact_storage.rs b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs index cb1546057..5d51cbdb6 100644 --- a/lib/apps/fabro-server/src/server/tests/artifact_storage.rs +++ b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs @@ -89,9 +89,10 @@ async fn artifact_upload_rejects_unauthorized_invalid_and_oversized_bodies_witho let run = create_run_with_bearer(&app, &user).await; let worker = issue_test_worker_token(&run); let other_worker = issue_test_worker_token(&RunId::new()); - let hash = state + let hash = BlobHash::new(b"valid content"); + state .artifact_store - .put_capture(&run, b"valid content") + .put_capture(&run, &hash, b"valid content") .await .unwrap(); let digest = hash.to_string(); @@ -264,7 +265,7 @@ async fn artifact_worker_client_uses_the_configured_s3_backend_and_prefix() { .await .unwrap(); let writer = ClientArtifactWriter::new(client, run); - assert_eq!(writer.write(&bytes).await.unwrap(), hash); + writer.write(&hash, &bytes).await.unwrap(); assert_eq!( state .artifact_store diff --git a/lib/components/fabro-petri/src/artifacts.rs b/lib/components/fabro-petri/src/artifacts.rs index d7eca8376..758b87ce3 100644 --- a/lib/components/fabro-petri/src/artifacts.rs +++ b/lib/components/fabro-petri/src/artifacts.rs @@ -11,19 +11,15 @@ pub enum ArtifactWriteError { Store(#[from] fabro_store::Error), #[error("artifact upload failed")] Upload(#[source] anyhow::Error), - #[error("artifact writer returned {actual}, expected {expected}")] - Integrity { - expected: BlobHash, - actual: BlobHash, - }, } -/// Run-bound storage for captured files. Implementations publish complete -/// objects before returning their content hash, preserve errors, and permit -/// concurrent and repeated writes of identical content. +/// Run-bound storage for captured files, keyed by the `digest` the caller +/// computed over `bytes`. Implementations publish complete objects before +/// returning, preserve errors, and permit concurrent and repeated writes of +/// identical content. #[async_trait::async_trait] pub trait ArtifactWriter: Send + Sync { - async fn write(&self, bytes: &[u8]) -> Result; + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError>; } pub struct StoreArtifactWriter { @@ -40,9 +36,9 @@ impl StoreArtifactWriter { #[async_trait::async_trait] impl ArtifactWriter for StoreArtifactWriter { - async fn write(&self, bytes: &[u8]) -> Result { + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { self.store - .put_capture(&self.run_id, bytes) + .put_capture(&self.run_id, digest, bytes) .await .map_err(ArtifactWriteError::from) } @@ -62,12 +58,10 @@ impl ClientArtifactWriter { #[async_trait::async_trait] impl ArtifactWriter for ClientArtifactWriter { - async fn write(&self, bytes: &[u8]) -> Result { - let hash = BlobHash::new(bytes); + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { self.client - .write_run_artifact_content(&self.run_id, &hash, bytes) + .write_run_artifact_content(&self.run_id, digest, bytes) .await - .map_err(ArtifactWriteError::Upload)?; - Ok(hash) + .map_err(ArtifactWriteError::Upload) } } diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 768af58cd..b322cce1e 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -63,7 +63,6 @@ use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; use crate::admission::AdmittedGraphs; -use crate::artifacts::ArtifactWriter; use crate::blobs::{Blobs, RunBlobs}; use crate::controls::RunControls; use crate::hooks::{FabroHooks, HooksSpec}; @@ -85,38 +84,36 @@ pub enum Execution { pub struct RunRequest { /// The Fabro run id, which becomes Petri's run key: the run's identity /// in the store and the label on every sandbox of the run. - pub run_id: String, + pub run_id: String, /// Where the run's workspaces, step output and blobs live. - pub run_dir: PathBuf, - pub execution: Execution, + pub run_dir: PathBuf, + pub execution: Execution, /// The run's durable record: the worker's HTTP store, or the server's /// SQLite store under the test override. - pub store: Arc, - pub runtime: RuntimeSpec, + pub store: Arc, + pub runtime: RuntimeSpec, /// The sandbox provider Fabro resolved for the run's environment. - pub provider: SandboxProviderKind, + pub provider: SandboxProviderKind, /// Fires to cancel the run. - pub cancel: CancellationToken, + pub cancel: CancellationToken, /// The run's pause, unpause and steer controls, which the caller keeps /// a clone of to drive them while the run is live. - pub controls: RunControls, + pub controls: RunControls, /// Where the run's questions go. - pub interviewer: Arc, + pub interviewer: Arc, /// The caller's observers of every record, registered ahead of the /// interview dispatcher: the interviewer's own expiry observer among /// them. - pub observers: Vec>, + pub observers: Vec>, /// Where `{{ secrets.NAME }}` references resolve from; `None` leaves /// every secret unknown. - pub secrets: Option>, + pub secrets: Option>, /// Where offloaded stage values go; `None` keeps Petri's local store /// under the run directory. - pub blobs: Option>, - /// Where captured workspace files go. Required when capture writes files. - pub artifact_writer: Option>, + pub blobs: Option>, /// Fabro's hooks: the checkpoint commit and its record. `None` runs /// with Petri's local hook service alone. - pub hooks: Option, + pub hooks: Option, } /// The recorded status of a finished run. @@ -213,7 +210,6 @@ pub async fn run(request: RunRequest) -> Result { Arc::clone(&request.store), resumed, request.blobs.clone(), - request.artifact_writer.clone(), )) }); if let Some(hooks) = &fabro_hooks { diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 29667dd05..03e662b8b 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -185,8 +185,6 @@ pub enum HookError { }, #[error("invalid run.artifacts.include pattern")] Globs(#[source] Arc), - #[error("the run has no configured artifact writer")] - NoArtifactWriter, #[error("the captured artifact could not be stored")] Artifact(#[source] ArtifactWriteError), #[error("the workspace could not be listed below `{root}`")] @@ -215,26 +213,33 @@ impl HookError { /// What Fabro's hooks need beside the run: where the platform records go, /// the run's Git settings, and which files are the run's artifacts. pub struct HooksSpec { - pub records: Arc, - pub git: RunGitSettings, + pub records: Arc, + pub git: RunGitSettings, /// The `[run.artifacts] include` patterns: which files of a stage's /// workspace are collected after the stage. - pub artifacts: Vec, + pub artifacts: Vec, /// A test's gate directory: a checkpoint point named by a `.hold` file /// there waits for its `.release` file. `None` outside tests. - pub test_gates: Option, + pub test_gates: Option, + /// Where captured workspace files go. + pub artifact_writer: Arc, } impl HooksSpec { /// The spec a run's settings give: its Git settings and its artifact - /// patterns. + /// patterns, captured through `artifact_writer`. #[must_use] - pub fn for_run(records: Arc, settings: &RunNamespace) -> Self { + pub fn for_run( + records: Arc, + settings: &RunNamespace, + artifact_writer: Arc, + ) -> Self { Self { records, git: RunGitSettings::from(settings), artifacts: settings.artifacts.include.clone(), test_gates: None, + artifact_writer, } } @@ -372,7 +377,7 @@ pub struct FabroHooks { records: Arc, /// Where diff patches go; `None` records summaries alone. blobs: Option>, - artifact_writer: Option>, + artifact_writer: Arc, workspaces: RunWorkspaces, lookup: WorkspaceLookup, identity: GitIdentity, @@ -399,8 +404,7 @@ impl FabroHooks { /// run whose records are in `store` under `run_key`, with its /// workspaces under `run_dir`. `resumed` says the run continues from /// its records, so a sandbox workspace is brought to its snapshot at - /// its scope's first acquisition. `blobs` holds diff patches; - /// `artifact_writer` holds captured workspace files. + /// its scope's first acquisition. `blobs` holds diff patches. #[must_use] pub fn new( spec: HooksSpec, @@ -411,7 +415,6 @@ impl FabroHooks { store: Arc, resumed: bool, blobs: Option>, - artifact_writer: Option>, ) -> Self { let identity = GitIdentity { name: spec.git.author.name.clone(), @@ -429,7 +432,7 @@ impl FabroHooks { run_id, records: spec.records, blobs, - artifact_writer, + artifact_writer: spec.artifact_writer, workspaces, lookup: WorkspaceLookup::new(Arc::clone(&store), run_key), identity, @@ -949,14 +952,16 @@ impl FabroHooks { return Ok(0); }; let candidates = list_artifacts(env.as_ref(), globs).await?; - let limit = usize::try_from(ARTIFACT_MAX_FILE_BYTES).unwrap_or(usize::MAX); let mut collected = 0; let mut total_bytes = 0_u64; for (path, size) in select_artifacts(candidates) { if total_bytes.saturating_add(size) > ARTIFACT_MAX_TOTAL_BYTES { break; } - let bytes = match env.read_file_limited(Path::new(&path), limit).await { + let bytes = match env + .read_file_limited(Path::new(&path), fabro_types::ARTIFACT_MAX_FILE_BYTES) + .await + { Ok(Some(bytes)) => bytes, Ok(None) => continue, Err(error) => { @@ -964,7 +969,7 @@ impl FabroHooks { continue; } }; - if !self.store_artifact(key, path, &bytes).await? { + if !self.store_artifact(key, &path, &bytes).await? { continue; } total_bytes = total_bytes.saturating_add(size); @@ -978,32 +983,25 @@ impl FabroHooks { async fn store_artifact( &self, key: CheckpointKey, - path: String, + path: &str, bytes: &[u8], ) -> Result { let already = self.collected_artifacts().await?; let digest = BlobHash::new(bytes); - let identity = (path.clone(), digest.to_string()); + let identity = (path.to_owned(), digest.to_string()); if sync::lock(already).contains(&identity) { return Ok(false); } - let writer = self - .artifact_writer - .as_ref() - .ok_or(HookError::NoArtifactWriter)?; - let stored = writer.write(bytes).await.map_err(HookError::Artifact)?; - if stored != digest { - return Err(HookError::Artifact(ArtifactWriteError::Integrity { - expected: digest, - actual: stored, - })); - } + self.artifact_writer + .write(&digest, bytes) + .await + .map_err(HookError::Artifact)?; let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord { execution: key.execution, firing: key.firing, attempt: key.attempt, - path: path.clone(), - source: ArtifactSource::ObjectStore(stored), + path: path.to_owned(), + source: ArtifactSource::ObjectStore(digest), bytes: u64::try_from(bytes.len()).unwrap_or(u64::MAX), digest: digest.to_string(), operation: Some(key.operation_for(ARTIFACT_EFFECT)), @@ -1420,13 +1418,13 @@ mod tests { #[async_trait::async_trait] impl ArtifactWriter for FlakyWriter { - async fn write(&self, bytes: &[u8]) -> Result { + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { if self.fail.swap(false, Ordering::SeqCst) { return Err(ArtifactWriteError::Store(fabro_store::Error::Io( std::io::Error::other("test store unavailable"), ))); } - self.writer.write(bytes).await + self.writer.write(digest, bytes).await } } @@ -1434,7 +1432,7 @@ mod tests { root: &Path, run_id: RunId, records: Arc, - writer: Option>, + artifact_writer: Arc, ) -> FabroHooks { FabroHooks::new( HooksSpec { @@ -1442,6 +1440,7 @@ mod tests { git: RunGitSettings::default(), artifacts: vec!["assets/**".to_string()], test_gates: None, + artifact_writer, }, Arc::new(NoHooks), run_id, @@ -1450,7 +1449,6 @@ mod tests { Arc::new(petri_store::MemoryRunStore::new()), false, None, - writer, ) } @@ -1468,7 +1466,7 @@ mod tests { writer: StoreArtifactWriter::new(store.clone(), run), fail: AtomicBool::new(fail_upload), }); - let hooks = artifact_hooks(root.path(), run, records.clone(), Some(writer.clone())); + let hooks = artifact_hooks(root.path(), run, records.clone(), writer.clone()); let key = CheckpointKey { execution: 0, firing: 1, @@ -1477,10 +1475,7 @@ mod tests { let path = "assets/report.bin".to_string(); let bytes = b"binary\0payload"; let hash = BlobHash::new(bytes); - let error = hooks - .store_artifact(key, path.clone(), bytes) - .await - .unwrap_err(); + let error = hooks.store_artifact(key, &path, bytes).await.unwrap_err(); assert!(error.render().contains(if fail_upload { "test store unavailable" } else { @@ -1493,26 +1488,16 @@ mod tests { !fail_upload ); assert!(store.list_for_run(&run).await.unwrap().is_empty()); + assert!(hooks.store_artifact(key, &path, bytes).await.unwrap()); + assert!(!hooks.store_artifact(key, &path, bytes).await.unwrap()); + let resumed = artifact_hooks(root.path(), run, records.clone(), writer); + assert!(!resumed.store_artifact(key, &path, bytes).await.unwrap()); assert!( - hooks - .store_artifact(key, path.clone(), bytes) + resumed + .store_artifact(key, &path, b"changed") .await .unwrap() ); - assert!( - !hooks - .store_artifact(key, path.clone(), bytes) - .await - .unwrap() - ); - let resumed = artifact_hooks(root.path(), run, records.clone(), Some(writer)); - assert!( - !resumed - .store_artifact(key, path.clone(), bytes) - .await - .unwrap() - ); - assert!(resumed.store_artifact(key, path, b"changed").await.unwrap()); assert_eq!(records.records.records(&run).len(), 2); } } @@ -1552,18 +1537,13 @@ mod tests { root.path(), run, records.clone(), - Some(Arc::new(StoreArtifactWriter::new(store.clone(), run))), - ); - assert!( - !resumed - .store_artifact(key, path.clone(), bytes) - .await - .unwrap() + Arc::new(StoreArtifactWriter::new(store.clone(), run)), ); + assert!(!resumed.store_artifact(key, &path, bytes).await.unwrap()); assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); assert!( resumed - .store_artifact(key, path.clone(), b"changed") + .store_artifact(key, &path, b"changed") .await .unwrap() ); @@ -1576,9 +1556,9 @@ mod tests { root.path(), fork, records.clone(), - Some(Arc::new(StoreArtifactWriter::new(store.clone(), fork))), + Arc::new(StoreArtifactWriter::new(store.clone(), fork)), ); - assert!(fork_hooks.store_artifact(key, path, bytes).await.unwrap()); + assert!(fork_hooks.store_artifact(key, &path, bytes).await.unwrap()); assert_eq!( store.get_capture(&fork, &hash).await.unwrap().unwrap(), bytes.as_slice() @@ -1587,26 +1567,6 @@ mod tests { assert_eq!(records.records(&fork).len(), 1); } - #[tokio::test] - async fn artifact_capture_requires_its_writer_without_falling_back_to_blobs() { - let root = tempfile::tempdir().unwrap(); - let run = RunId::new(); - let records = Arc::new(MemoryPlatformRecords::new()); - let hooks = artifact_hooks(root.path(), run, records.clone(), None); - let key = CheckpointKey { - execution: 0, - firing: 1, - attempt: 1, - }; - assert!(matches!( - hooks - .store_artifact(key, "assets/file".to_string(), b"payload") - .await, - Err(HookError::NoArtifactWriter) - )); - assert!(records.records(&run).is_empty()); - } - #[test] fn the_selection_keeps_the_smallest_files_within_the_budgets() { let mut candidates: Vec<(String, u64)> = (0..(ARTIFACT_MAX_FILES + 5)) diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index a2f0f7f11..026f4655c 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -181,13 +181,17 @@ impl Harness { fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec { HooksSpec { - records: Arc::clone(&self.records) as Arc, - git: RunGitSettings { + records: Arc::clone(&self.records) as Arc, + git: RunGitSettings { host_workspaces: *provider == SandboxProviderKind::LOCAL, ..RunGitSettings::default() }, - artifacts: self.artifacts.clone(), - test_gates: None, + artifacts: self.artifacts.clone(), + test_gates: None, + artifact_writer: Arc::new(StoreArtifactWriter::new( + self.artifact_store.clone(), + self.run_id, + )), } } @@ -220,10 +224,6 @@ impl Harness { observers, secrets: None, blobs: Some(Arc::clone(&self.blobs) as Arc), - artifact_writer: Some(Arc::new(StoreArtifactWriter::new( - self.artifact_store.clone(), - self.run_id, - ))), hooks: Some(hooks), }; engine::run(request).await.expect("the run executes") @@ -851,7 +851,6 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() { observers, secrets: None, blobs: None, - artifact_writer: None, hooks: Some(harness.hooks(&SandboxProviderKind::LOCAL)), }; let outcome = engine::run(request).await.expect("the run executes"); diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs index 9fd47f914..7cc8fd0c3 100644 --- a/lib/components/fabro-petri/tests/support/mod.rs +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -121,7 +121,6 @@ pub(crate) fn run_request( interviewer: Arc::new(interviewer), secrets: None, blobs: None, - artifact_writer: None, hooks: None, } } diff --git a/lib/components/fabro-store/src/artifact_store.rs b/lib/components/fabro-store/src/artifact_store.rs index 2193d2a34..5a24d5760 100644 --- a/lib/components/fabro-store/src/artifact_store.rs +++ b/lib/components/fabro-store/src/artifact_store.rs @@ -74,40 +74,47 @@ impl ArtifactStore { } pub async fn put(&self, run_id: &RunId, key: &ArtifactKey, data: &[u8]) -> Result<()> { - let path = self.artifact_path(run_id, key)?; - self.object_store - .put(&path, Bytes::copy_from_slice(data).into()) - .await?; - Ok(()) + self.put_at(&self.artifact_path(run_id, key)?, data).await } /// Publish a complete captured file under its content digest within the - /// run. Repeating a put of the same content is safe; metadata is recorded + /// run. `hash` must be `BlobHash::new(data)`; callers verify or compute + /// it. Repeating a put of the same content is safe; metadata is recorded /// separately only after this operation succeeds. - pub async fn put_capture(&self, run_id: &RunId, data: &[u8]) -> Result { - let hash = BlobHash::new(data); - let path = self.capture_prefix(run_id)?.child(hash.to_string()); - self.object_store - .put(&path, Bytes::copy_from_slice(data).into()) - .await?; - Ok(hash) + pub async fn put_capture(&self, run_id: &RunId, hash: &BlobHash, data: &[u8]) -> Result<()> { + debug_assert_eq!(*hash, BlobHash::new(data)); + self.put_at(&self.capture_path(run_id, hash)?, data).await } /// Read run-owned content named by a capture record, without consulting /// historical stage keys or the SQLite blob table. pub async fn get_capture(&self, run_id: &RunId, hash: &BlobHash) -> Result> { - let path = self.capture_prefix(run_id)?.child(hash.to_string()); - match self.object_store.get(&path).await { - Ok(result) => Ok(Some(result.bytes().await?)), - Err(object_store::Error::NotFound { .. }) => Ok(None), - Err(error) => Err(error.into()), - } + self.get_at(&self.capture_path(run_id, hash)?).await } fn capture_prefix(&self, run_id: &RunId) -> Result { Ok(self.run_prefix(run_id)?.child("captures").child("sha256")) } + fn capture_path(&self, run_id: &RunId, hash: &BlobHash) -> Result { + Ok(self.capture_prefix(run_id)?.child(hash.to_string())) + } + + async fn put_at(&self, path: &ObjectPath, data: &[u8]) -> Result<()> { + self.object_store + .put(path, Bytes::copy_from_slice(data).into()) + .await?; + Ok(()) + } + + async fn get_at(&self, path: &ObjectPath) -> Result> { + match self.object_store.get(path).await { + Ok(result) => Ok(Some(result.bytes().await?)), + Err(object_store::Error::NotFound { .. }) => Ok(None), + Err(err) => Err(err.into()), + } + } + pub fn writer(&self, run_id: &RunId, key: &ArtifactKey) -> Result { let path = self.artifact_path(run_id, key)?; Ok(BufWriter::with_capacity( @@ -142,12 +149,7 @@ impl ArtifactStore { } pub async fn get(&self, run_id: &RunId, key: &ArtifactKey) -> Result> { - let path = self.artifact_path(run_id, key)?; - match self.object_store.get(&path).await { - Ok(result) => Ok(Some(result.bytes().await?)), - Err(object_store::Error::NotFound { .. }) => Ok(None), - Err(err) => Err(err.into()), - } + self.get_at(&self.artifact_path(run_id, key)?).await } pub async fn get_stream( @@ -472,7 +474,7 @@ mod tests { store.write_metadata("test").await.unwrap(); store.put(&run, &legacy, b"legacy").await.unwrap(); for id in [&run, &run, &other] { - assert_eq!(store.put_capture(id, bytes).await.unwrap(), hash); + store.put_capture(id, &hash, bytes).await.unwrap(); } let location = ObjectPath::from(format!("artifacts/{run}/captures/sha256/{hash}")); assert_eq!( diff --git a/lib/foundation/fabro-config/src/resolve/server.rs b/lib/foundation/fabro-config/src/resolve/server.rs index 1c472664e..8ce62c516 100644 --- a/lib/foundation/fabro-config/src/resolve/server.rs +++ b/lib/foundation/fabro-config/src/resolve/server.rs @@ -285,7 +285,7 @@ fn resolve_artifacts( provider, layer.and_then(|artifacts| artifacts.local.as_ref()), layer.and_then(|artifacts| artifacts.s3.as_ref()), - &object_store_default_root(storage_root, "artifacts"), + &ServerArtifactsSettings::default_local_root(Path::new(storage_root)), "server.artifacts", errors, ), @@ -330,14 +330,6 @@ fn resolve_object_store( } } -fn object_store_default_root(storage_root: &str, domain: &str) -> String { - Path::new(storage_root) - .join("objects") - .join(domain) - .to_string_lossy() - .into_owned() -} - fn resolve_integrations(layer: Option<&ServerIntegrationsLayer>) -> ServerIntegrationsSettings { ServerIntegrationsSettings { github: layer diff --git a/lib/foundation/fabro-types/src/dense.rs b/lib/foundation/fabro-types/src/dense.rs index 4be1bc1c6..3fe7e5cd2 100644 --- a/lib/foundation/fabro-types/src/dense.rs +++ b/lib/foundation/fabro-types/src/dense.rs @@ -4,8 +4,8 @@ use std::path::Path; use serde::{Deserialize, Serialize}; use crate::settings::{ - CliNamespace, ObjectStoreSettings, ProjectNamespace, RunNamespace, ServerNamespace, - WorkflowNamespace, + CliNamespace, ObjectStoreSettings, ProjectNamespace, RunNamespace, ServerArtifactsSettings, + ServerNamespace, WorkflowNamespace, }; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -18,10 +18,11 @@ impl ServerSettings { pub fn with_storage_override(mut self, path: &Path) -> Self { // Only the derived default follows the storage directory. A custom // artifact location is independent of the database and runtime root. - let default_artifact_root = Path::new(&self.server.storage.root).join("objects/artifacts"); + let default_artifact_root = + ServerArtifactsSettings::default_local_root(Path::new(&self.server.storage.root)); if let ObjectStoreSettings::Local { root } = &mut self.server.artifacts.store { - if Path::new(root) == default_artifact_root { - *root = path.join("objects/artifacts").display().to_string(); + if *root == default_artifact_root { + *root = ServerArtifactsSettings::default_local_root(path); } } self.server.storage.root = path.display().to_string(); diff --git a/lib/foundation/fabro-types/src/settings/server.rs b/lib/foundation/fabro-types/src/settings/server.rs index 4f4279ded..9860db349 100644 --- a/lib/foundation/fabro-types/src/settings/server.rs +++ b/lib/foundation/fabro-types/src/settings/server.rs @@ -216,6 +216,18 @@ pub struct ServerArtifactsSettings { pub store: ObjectStoreSettings, } +impl ServerArtifactsSettings { + /// The local artifact root a storage root gives when none is configured. + #[must_use] + pub fn default_local_root(storage_root: &std::path::Path) -> String { + storage_root + .join("objects") + .join("artifacts") + .to_string_lossy() + .into_owned() + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "type", rename_all = "snake_case")] pub enum ObjectStoreSettings {