From 7f55dd87b577b20da75a7fa3dfeb9c2b0455a27d Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Mon, 28 Sep 2026 14:30:36 -0400 Subject: [PATCH] Harden artifact capture writes - The store checks that captured bytes match their digest in every build, hashing on a blocking thread, and the upload handler relies on that check instead of hashing a second time. - Capture bytes travel as `Bytes` from the hooks through the client and the store, so uploads and retries share one buffer. - The artifact writer takes the run ID from the hooks, so objects are stored under the run their records name. - Concurrent captures of the same file and content wait on one another, so the file is uploaded and recorded once. - When a record append fails, the hooks re-read the run's captures and treat a record that did land as done, so a lost response does not record the capture twice. - An upload that finishes after its run was deleted removes itself, instead of leaving an object nothing references. - Listing a run's stage artifacts skips everything under `captures/`, so an unexpected object there cannot fail the listing or the ZIP. - The capture record derives its `digest` key from its source instead of storing it twice, still writing and checking it on the wire. Co-Authored-By: Claude Opus 5.5 --- .../src/commands/run/petri_worker.rs | 5 +- .../src/server/handler/artifacts.rs | 54 ++-- .../fabro-server/src/server/petri_runs.rs | 5 +- .../src/server/tests/artifact_storage.rs | 9 +- lib/components/fabro-petri/src/artifacts.rs | 47 ++-- lib/components/fabro-petri/src/hooks.rs | 254 +++++++++++++++--- lib/components/fabro-petri/tests/hooks.rs | 8 +- .../fabro-store/src/artifact_store.rs | 113 ++++++-- lib/components/fabro-store/src/error.rs | 2 + .../fabro-store/src/platform_records.rs | 79 +++--- lib/foundation/fabro-client/src/client.rs | 6 +- .../fabro-types/src/run_projection.rs | 4 +- 12 files changed, 438 insertions(+), 148 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 d4b1ab37d..4f10a66e2 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -201,10 +201,7 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { 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, - )), + Arc::new(ClientArtifactWriter::new(worker.client.clone_for_reuse())), ) .with_test_gates(test_checkpoint_gates()); let request = RunRequest { diff --git a/lib/apps/fabro-server/src/server/handler/artifacts.rs b/lib/apps/fabro-server/src/server/handler/artifacts.rs index 46ba7fab2..42daa4032 100644 --- a/lib/apps/fabro-server/src/server/handler/artifacts.rs +++ b/lib/apps/fabro-server/src/server/handler/artifacts.rs @@ -8,9 +8,9 @@ use async_zip::{Compression, ZipEntryBuilder}; use axum::extract::DefaultBodyLimit; use axum::extract::rejection::BytesRejection; use axum::http::HeaderValue; -use axum::routing::put; +use axum::routing; use fabro_store::{ArtifactStore, BlobStore, Error as StoreError}; -use fabro_types::{ARTIFACT_MAX_FILE_BYTES, ArtifactSource, BlobHash, RunProjection}; +use fabro_types::{ARTIFACT_MAX_FILE_BYTES, ArtifactSource, RunProjection}; use fabro_util::error::collect_chain; use futures_util::SinkExt as _; use futures_util::io::AsyncWriteExt as _; @@ -27,7 +27,7 @@ use super::super::{ RequireRunScoped, RequiredUser, Response, Router, RunArtifactEntry, RunArtifactListResponse, RunId, StageArtifactEntry, State, StatusCode, WriteBlobResponse, get, header, 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, + reject_if_archived, required_query_param, run_records, validate_relative_artifact_path, }; use crate::principal_middleware::RequireWorkerRunSegment; @@ -35,7 +35,8 @@ pub(super) fn routes() -> Router> { Router::new() .route( "/runs/{id}/artifacts/content/{digest}", - put(write_run_artifact_content).layer(DefaultBodyLimit::max(ARTIFACT_MAX_FILE_BYTES)), + routing::put(write_run_artifact_content) + .layer(DefaultBodyLimit::max(ARTIFACT_MAX_FILE_BYTES)), ) .route("/runs/{id}/blobs", post(write_run_blob)) .route("/runs/{id}/blobs/{blobHash}", get(read_run_blob)) @@ -79,23 +80,34 @@ async fn write_run_artifact_content( if let Err(error) = state.load_run_projection(&id).await { return error.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, &expected, &body) - .await - { - Ok(()) => StatusCode::NO_CONTENT.into_response(), + match state.artifact_store.put_capture(&id, &expected, body).await { + Ok(()) => {} + Err(StoreError::CaptureDigestMismatch { .. }) => { + return ApiError::bad_request("Artifact content does not match its digest.") + .into_response(); + } Err(error) => { warn!(run_id = %id, error = %collect_chain(&error).join(": "), "Artifact upload failed"); - ApiError::new( + return ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, "Artifact storage failed.", ) - .into_response() + .into_response(); + } + } + // Run deletion removes the run before its objects. A run still present + // now is deleted after this upload landed, and its deletion removes the + // upload; a run already gone never will, so the upload goes here. + match run_records::projection(state.as_ref(), id).await { + Ok(Some(_)) => StatusCode::NO_CONTENT.into_response(), + Ok(None) => { + if let Err(cleanup) = state.artifact_store.delete_capture(&id, &expected).await { + warn!(run_id = %id, error = %collect_chain(&cleanup).join(": "), "Artifact upload for a deleted run could not be removed"); + } + ApiError::not_found("Run not found.").into_response() + } + Err(error) => { + ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response() } } } @@ -564,7 +576,7 @@ async fn get_stage_artifact( #[cfg(test)] mod tests { use async_zip::base::read::mem::ZipFileReader; - use fabro_types::{RunArtifact, StageId, test_support}; + use fabro_types::{BlobHash, RunArtifact, StageId, test_support}; use super::*; @@ -583,7 +595,7 @@ mod tests { let new = BlobHash::new(b"new object"); state .artifact_store - .put_capture(&run, &new, b"new object") + .put_capture(&run, &new, Bytes::from_static(b"new object")) .await .unwrap(); for (path, size, source) in [ @@ -612,7 +624,11 @@ mod tests { // Unrecorded content is never a file-list entry. state .artifact_store - .put_capture(&run, &BlobHash::new(b"orphan"), b"orphan") + .put_capture( + &run, + &BlobHash::new(b"orphan"), + Bytes::from_static(b"orphan"), + ) .await .unwrap(); let entries = run_artifacts(&state, &run, &projection).await.unwrap(); diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 1055f7215..f802d6c36 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -456,10 +456,7 @@ 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, - )), + Arc::new(StoreArtifactWriter::new(state.artifact_store.clone())), ); let request = RunRequest { run_id: run_id.to_string(), 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 5d51cbdb6..84ad08499 100644 --- a/lib/apps/fabro-server/src/server/tests/artifact_storage.rs +++ b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs @@ -92,7 +92,7 @@ async fn artifact_upload_rejects_unauthorized_invalid_and_oversized_bodies_witho let hash = BlobHash::new(b"valid content"); state .artifact_store - .put_capture(&run, &hash, b"valid content") + .put_capture(&run, &hash, Bytes::from_static(b"valid content")) .await .unwrap(); let digest = hash.to_string(); @@ -264,8 +264,11 @@ async fn artifact_worker_client_uses_the_configured_s3_backend_and_prefix() { .connect() .await .unwrap(); - let writer = ClientArtifactWriter::new(client, run); - writer.write(&hash, &bytes).await.unwrap(); + let writer = ClientArtifactWriter::new(client); + writer + .write(&run, &hash, Bytes::from(bytes.clone())) + .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 758b87ce3..c25cc9d0e 100644 --- a/lib/components/fabro-petri/src/artifacts.rs +++ b/lib/components/fabro-petri/src/artifacts.rs @@ -1,6 +1,7 @@ //! Captured workspace files go to the server's configured artifact store. //! Engine values and patches keep using the separate blob capability. +use bytes::Bytes; use fabro_client::Client; use fabro_store::ArtifactStore; use fabro_types::{BlobHash, RunId}; @@ -13,54 +14,68 @@ pub enum ArtifactWriteError { Upload(#[source] anyhow::Error), } -/// 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. +/// Storage for captured files, keyed by the run the hooks record them under +/// and the `digest` of `bytes`. Implementations reject bytes that do not +/// match `digest`, 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, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError>; + async fn write( + &self, + run_id: &RunId, + digest: &BlobHash, + bytes: Bytes, + ) -> Result<(), ArtifactWriteError>; } pub struct StoreArtifactWriter { - store: ArtifactStore, - run_id: RunId, + store: ArtifactStore, } impl StoreArtifactWriter { #[must_use] - pub fn new(store: ArtifactStore, run_id: RunId) -> Self { - Self { store, run_id } + pub fn new(store: ArtifactStore) -> Self { + Self { store } } } #[async_trait::async_trait] impl ArtifactWriter for StoreArtifactWriter { - async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { + async fn write( + &self, + run_id: &RunId, + digest: &BlobHash, + bytes: Bytes, + ) -> Result<(), ArtifactWriteError> { self.store - .put_capture(&self.run_id, digest, bytes) + .put_capture(run_id, digest, bytes) .await .map_err(ArtifactWriteError::from) } } +/// Uploads to the server, which checks the digest before storing. pub struct ClientArtifactWriter { client: Client, - run_id: RunId, } impl ClientArtifactWriter { #[must_use] - pub fn new(client: Client, run_id: RunId) -> Self { - Self { client, run_id } + pub fn new(client: Client) -> Self { + Self { client } } } #[async_trait::async_trait] impl ArtifactWriter for ClientArtifactWriter { - async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { + async fn write( + &self, + run_id: &RunId, + digest: &BlobHash, + bytes: Bytes, + ) -> Result<(), ArtifactWriteError> { self.client - .write_run_artifact_content(&self.run_id, digest, bytes) + .write_run_artifact_content(run_id, digest, bytes) .await .map_err(ArtifactWriteError::Upload) } diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 03e662b8b..892a24838 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -76,8 +76,10 @@ use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex, OnceLock}; use std::time::Duration; +use bytes::Bytes; use fabro_store::platform_records::{ - ArtifactCollectedRecord, CheckpointRecord, GitIdentityRecord, RunBranchRecord, RunDiffRecord, + ArtifactCollectedRecord, CheckpointRecord, GitIdentityRecord, OperationKey, RunBranchRecord, + RunDiffRecord, }; use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition}; use fabro_types::settings::run::RunNamespace; @@ -256,7 +258,7 @@ impl HooksSpec { type AcquiredEnv = (String, Arc); /// The identity of a collected file: its path and content digest. -type ArtifactIdentity = (String, String); +type ArtifactIdentity = (String, BlobHash); /// A checkpoint's workspace and commit. type WorkspaceCommit = (String, String); @@ -314,6 +316,9 @@ struct ArtifactLedger { /// Every artifact collected so far, by path and digest: read from the /// store once, then kept current with every append. collected: OnceCell>>, + /// One lock per identity being captured, so concurrent transitions that + /// leave the same file upload and record it once. + capturing: Mutex>>>, } /// Where each acquired scope's workspace is, and the locks that serialize @@ -442,6 +447,7 @@ impl FabroHooks { checkpoints: CheckpointLedger::default(), artifacts: ArtifactLedger { globs: WorkspaceGlobSet::try_new(&spec.artifacts).map_err(Arc::new), + capturing: Mutex::default(), collected: OnceCell::new(), }, scopes: ScopeEnvs::default(), @@ -969,7 +975,7 @@ impl FabroHooks { continue; } }; - if !self.store_artifact(key, &path, &bytes).await? { + if !self.store_artifact(key, &path, bytes.into()).await? { continue; } total_bytes = total_bytes.saturating_add(size); @@ -984,29 +990,41 @@ impl FabroHooks { &self, key: CheckpointKey, path: &str, - bytes: &[u8], + bytes: Bytes, ) -> Result { let already = self.collected_artifacts().await?; - let digest = BlobHash::new(bytes); - let identity = (path.to_owned(), digest.to_string()); + let digest = BlobHash::new(&bytes); + let identity = (path.to_owned(), digest); if sync::lock(already).contains(&identity) { return Ok(false); } + let capture = Arc::clone( + sync::lock(&self.artifacts.capturing) + .entry(identity.clone()) + .or_default(), + ); + let _capturing = capture.lock().await; + // Whoever held the lock before may have recorded this identity. + if sync::lock(already).contains(&identity) { + return Ok(false); + } + let size = u64::try_from(bytes.len()).unwrap_or(u64::MAX); self.artifact_writer - .write(&digest, bytes) + .write(&self.run_id, &digest, bytes) .await .map_err(HookError::Artifact)?; + let operation = key.operation_for(ARTIFACT_EFFECT); let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord { execution: key.execution, firing: key.firing, attempt: key.attempt, 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)), + bytes: size, + operation: Some(operation.clone()), }); - self.records + let appended = self + .records .append( &self.run_id, &record, @@ -1015,15 +1033,46 @@ impl FabroHooks { firing: key.firing, }), ) - .await - .map_err(|source| HookError::Write { - kind: "artifact", - source, - })?; - sync::lock(already).insert(identity); + .await; + if let Err(source) = appended { + // The append can reach the store and lose only its response. A + // record that landed is this capture; appending it again would + // list the file twice. + if !self.artifact_recorded(&identity, &operation).await { + return Err(HookError::Write { + kind: "artifact", + source, + }); + } + } + sync::lock(already).insert(identity.clone()); + sync::lock(&self.artifacts.capturing).remove(&identity); Ok(true) } + /// Whether the run's records hold this capture's record, read fresh. + async fn artifact_recorded( + &self, + identity: &ArtifactIdentity, + operation: &OperationKey, + ) -> bool { + let Ok(stored) = self + .records + .read_kind(&self.run_id, PlatformRecordKind::ArtifactCollected) + .await + else { + return false; + }; + stored.into_iter().any(|record| match record.record { + PlatformRecord::ArtifactCollected(artifact) => { + artifact.path == identity.0 + && artifact.source.hash() == identity.1 + && artifact.operation.as_ref() == Some(operation) + } + _ => false, + }) + } + /// The artifacts already collected for the run, read once: a file that /// is unchanged since it was collected is not collected again. async fn collected_artifacts(&self) -> Result<&Mutex>, HookError> { @@ -1042,7 +1091,7 @@ impl FabroHooks { .into_iter() .filter_map(|record| match record.record { PlatformRecord::ArtifactCollected(artifact) => { - Some((artifact.path, artifact.digest)) + Some((artifact.path, artifact.source.hash())) } _ => None, }) @@ -1369,10 +1418,11 @@ impl ExecutionHooks for FabroHooks { #[cfg(test)] mod tests { - use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use fabro_store::{ArtifactStore, StoredPlatformRecord}; use object_store::memory::InMemory; + use tokio::sync::Barrier; use super::*; use crate::artifacts::StoreArtifactWriter; @@ -1385,6 +1435,8 @@ mod tests { struct FlakyRecords { records: MemoryPlatformRecords, fail: AtomicBool, + /// Commit the failing append anyway, as a lost response does. + commit: bool, } #[async_trait::async_trait] @@ -1396,6 +1448,9 @@ mod tests { position: Option, ) -> Result { if self.fail.swap(false, Ordering::SeqCst) { + if self.commit { + self.records.append(run_id, record, position).await?; + } return Err(PlatformRecordError::Store(fabro_store::Error::Io( std::io::Error::other("test append unavailable"), ))); @@ -1418,13 +1473,42 @@ mod tests { #[async_trait::async_trait] impl ArtifactWriter for FlakyWriter { - async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { + async fn write( + &self, + run_id: &RunId, + digest: &BlobHash, + bytes: Bytes, + ) -> 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(digest, bytes).await + self.writer.write(run_id, digest, bytes).await + } + } + + /// Holds every upload until the test has started all of them. + struct GatedWriter { + writer: StoreArtifactWriter, + gate: Barrier, + uploads: AtomicUsize, + } + + #[async_trait::async_trait] + impl ArtifactWriter for GatedWriter { + async fn write( + &self, + run_id: &RunId, + digest: &BlobHash, + bytes: Bytes, + ) -> Result<(), ArtifactWriteError> { + self.uploads.fetch_add(1, Ordering::SeqCst); + let write = self.writer.write(run_id, digest, bytes); + // A second capture of the identity must wait on the first, not + // reach this point: time the gate out rather than hang. + let _ = time::timeout(Duration::from_millis(200), self.gate.wait()).await; + write.await } } @@ -1461,9 +1545,10 @@ mod tests { let records = Arc::new(FlakyRecords { records: MemoryPlatformRecords::new(), fail: AtomicBool::new(!fail_upload), + commit: false, }); let writer = Arc::new(FlakyWriter { - writer: StoreArtifactWriter::new(store.clone(), run), + writer: StoreArtifactWriter::new(store.clone()), fail: AtomicBool::new(fail_upload), }); let hooks = artifact_hooks(root.path(), run, records.clone(), writer.clone()); @@ -1473,9 +1558,12 @@ mod tests { attempt: 1, }; let path = "assets/report.bin".to_string(); - let bytes = b"binary\0payload"; - let hash = BlobHash::new(bytes); - let error = hooks.store_artifact(key, &path, bytes).await.unwrap_err(); + let bytes = Bytes::from_static(b"binary\0payload"); + let hash = BlobHash::new(&bytes); + let error = hooks + .store_artifact(key, &path, bytes.clone()) + .await + .unwrap_err(); assert!(error.render().contains(if fail_upload { "test store unavailable" } else { @@ -1488,13 +1576,28 @@ 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()); + assert!( + hooks + .store_artifact(key, &path, bytes.clone()) + .await + .unwrap() + ); + assert!( + !hooks + .store_artifact(key, &path, bytes.clone()) + .await + .unwrap() + ); let resumed = artifact_hooks(root.path(), run, records.clone(), writer); - assert!(!resumed.store_artifact(key, &path, bytes).await.unwrap()); + assert!( + !resumed + .store_artifact(key, &path, bytes.clone()) + .await + .unwrap() + ); assert!( resumed - .store_artifact(key, &path, b"changed") + .store_artifact(key, &path, Bytes::from_static(b"changed")) .await .unwrap() ); @@ -1513,8 +1616,8 @@ mod tests { firing: 1, attempt: 1, }; - let bytes = b"old payload"; - let hash = BlobHash::new(bytes); + let bytes = Bytes::from_static(b"old payload"); + let hash = BlobHash::new(&bytes); let path = "assets/report.bin".to_string(); records .append( @@ -1526,7 +1629,6 @@ mod tests { path: path.clone(), source: ArtifactSource::SqliteBlob(hash), bytes: bytes.len() as u64, - digest: hash.to_string(), operation: None, }), None, @@ -1537,13 +1639,18 @@ mod tests { root.path(), run, records.clone(), - Arc::new(StoreArtifactWriter::new(store.clone(), run)), + Arc::new(StoreArtifactWriter::new(store.clone())), + ); + assert!( + !resumed + .store_artifact(key, &path, bytes.clone()) + .await + .unwrap() ); - 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, b"changed") + .store_artifact(key, &path, Bytes::from_static(b"changed")) .await .unwrap() ); @@ -1556,17 +1663,88 @@ mod tests { root.path(), fork, records.clone(), - Arc::new(StoreArtifactWriter::new(store.clone(), fork)), + Arc::new(StoreArtifactWriter::new(store.clone())), + ); + assert!( + fork_hooks + .store_artifact(key, &path, bytes.clone()) + .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() + bytes ); assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); assert_eq!(records.records(&fork).len(), 1); } + #[tokio::test] + async fn concurrent_captures_of_one_file_upload_and_record_it_once() { + let root = tempfile::tempdir().unwrap(); + let run = RunId::new(); + let records = Arc::new(MemoryPlatformRecords::new()); + let store = ArtifactStore::new(Arc::new(InMemory::new()), "artifacts"); + let writer = Arc::new(GatedWriter { + writer: StoreArtifactWriter::new(store.clone()), + gate: Barrier::new(2), + uploads: AtomicUsize::new(0), + }); + let hooks = artifact_hooks(root.path(), run, records.clone(), writer.clone()); + let key = |firing| CheckpointKey { + execution: 0, + firing, + attempt: 1, + }; + let bytes = Bytes::from_static(b"same payload"); + let (first, second) = tokio::join!( + hooks.store_artifact(key(1), "assets/report.bin", bytes.clone()), + hooks.store_artifact(key(2), "assets/report.bin", bytes.clone()), + ); + let mut outcomes = [first.unwrap(), second.unwrap()]; + outcomes.sort_unstable(); + assert_eq!(outcomes, [false, true]); + assert_eq!(writer.uploads.load(Ordering::SeqCst), 1); + assert_eq!(records.records(&run).len(), 1); + } + + #[tokio::test] + async fn a_committed_append_whose_response_was_lost_is_recorded_once() { + let root = tempfile::tempdir().unwrap(); + let run = RunId::new(); + let store = ArtifactStore::new(Arc::new(InMemory::new()), "artifacts"); + let records = Arc::new(FlakyRecords { + records: MemoryPlatformRecords::new(), + fail: AtomicBool::new(true), + commit: true, + }); + let hooks = artifact_hooks( + root.path(), + run, + records.clone(), + Arc::new(StoreArtifactWriter::new(store)), + ); + let key = CheckpointKey { + execution: 0, + firing: 1, + attempt: 1, + }; + let bytes = Bytes::from_static(b"payload"); + assert!( + hooks + .store_artifact(key, "assets/report.bin", bytes.clone()) + .await + .unwrap() + ); + assert!( + !hooks + .store_artifact(key, "assets/report.bin", bytes) + .await + .unwrap() + ); + assert_eq!(records.records.records(&run).len(), 1); + } + #[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 026f4655c..ac25bd1fd 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -188,10 +188,7 @@ impl Harness { }, artifacts: self.artifacts.clone(), test_gates: None, - artifact_writer: Arc::new(StoreArtifactWriter::new( - self.artifact_store.clone(), - self.run_id, - )), + artifact_writer: Arc::new(StoreArtifactWriter::new(self.artifact_store.clone())), } } @@ -484,8 +481,7 @@ async fn artifacts_the_branch_and_the_diffs_are_recorded() { "{artifacts:?}" ); assert_eq!(artifacts[0].bytes, 3); - assert_eq!(artifacts[0].digest, artifacts[0].source.hash().to_string()); - assert_ne!(artifacts[0].digest, artifacts[1].digest); + assert_ne!(artifacts[0].source.hash(), artifacts[1].source.hash()); assert!( artifacts .iter() diff --git a/lib/components/fabro-store/src/artifact_store.rs b/lib/components/fabro-store/src/artifact_store.rs index 5a24d5760..d5f0fd8d6 100644 --- a/lib/components/fabro-store/src/artifact_store.rs +++ b/lib/components/fabro-store/src/artifact_store.rs @@ -10,6 +10,7 @@ use object_store::buffered::BufWriter; use object_store::path::Path as ObjectPath; use percent_encoding::{AsciiSet, NON_ALPHANUMERIC, percent_decode_str, utf8_percent_encode}; use tokio::io::AsyncWriteExt; +use tokio::task; use crate::{Error, Result, StageId}; @@ -74,16 +75,44 @@ impl ArtifactStore { } pub async fn put(&self, run_id: &RunId, key: &ArtifactKey, data: &[u8]) -> Result<()> { - self.put_at(&self.artifact_path(run_id, key)?, data).await + self.put_at( + &self.artifact_path(run_id, key)?, + Bytes::copy_from_slice(data), + ) + .await } /// Publish a complete captured file under its content digest within the - /// 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, hash: &BlobHash, data: &[u8]) -> Result<()> { - debug_assert_eq!(*hash, BlobHash::new(data)); - self.put_at(&self.capture_path(run_id, hash)?, data).await + /// run, after checking that `data` hashes to `hash`. Repeating a put of + /// the same content is safe; metadata is recorded separately only after + /// this operation succeeds. + /// + /// # Errors + /// + /// [`Error::CaptureDigestMismatch`] when `data` does not hash to `hash`; + /// nothing is written then. + pub async fn put_capture(&self, run_id: &RunId, hash: &BlobHash, data: Bytes) -> Result<()> { + let path = self.capture_path(run_id, hash)?; + // Up to the capture size limit of bytes: hash off the async workers. + let (actual, data) = task::spawn_blocking(move || (BlobHash::new(&data), data)) + .await + .map_err(|err| Error::Other(format!("capture digest task failed: {err}")))?; + if actual != *hash { + return Err(Error::CaptureDigestMismatch { expected: *hash }); + } + self.put_at(&path, data).await + } + + /// Remove one captured file's content. Removing absent content succeeds. + pub async fn delete_capture(&self, run_id: &RunId, hash: &BlobHash) -> Result<()> { + match self + .object_store + .delete(&self.capture_path(run_id, hash)?) + .await + { + Ok(()) | Err(object_store::Error::NotFound { .. }) => Ok(()), + Err(err) => Err(err.into()), + } } /// Read run-owned content named by a capture record, without consulting @@ -92,18 +121,22 @@ impl ArtifactStore { self.get_at(&self.capture_path(run_id, hash)?).await } + /// Everything under a run's `captures/` belongs to recorded captures, + /// never to historical stage keys, whatever digest namespace it uses. + fn captures_root(&self, run_id: &RunId) -> Result { + Ok(self.run_prefix(run_id)?.child("captures")) + } + fn capture_prefix(&self, run_id: &RunId) -> Result { - Ok(self.run_prefix(run_id)?.child("captures").child("sha256")) + Ok(self.captures_root(run_id)?.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?; + async fn put_at(&self, path: &ObjectPath, data: Bytes) -> Result<()> { + self.object_store.put(path, data.into()).await?; Ok(()) } @@ -181,7 +214,7 @@ impl ArtifactStore { pub async fn list_for_run(&self, run_id: &RunId) -> Result> { let prefix = self.run_prefix(run_id)?; - let captures = self.capture_prefix(run_id)?; + let captures = self.captures_root(run_id)?; let mut stream = self.object_store.list(Some(&prefix)); let mut artifacts = Vec::new(); while let Some(meta) = stream.next().await.transpose()? { @@ -474,7 +507,10 @@ mod tests { store.write_metadata("test").await.unwrap(); store.put(&run, &legacy, b"legacy").await.unwrap(); for id in [&run, &run, &other] { - store.put_capture(id, &hash, bytes).await.unwrap(); + store + .put_capture(id, &hash, Bytes::from_static(bytes)) + .await + .unwrap(); } let location = ObjectPath::from(format!("artifacts/{run}/captures/sha256/{hash}")); assert_eq!( @@ -510,16 +546,55 @@ mod tests { } #[tokio::test] - async fn capture_namespace_does_not_hide_malformed_legacy_objects() { + async fn capture_namespace_is_never_listed_as_stage_artifacts() { let objects: Arc = Arc::new(InMemory::new()); let store = ArtifactStore::new(objects.clone(), "artifacts"); let run = RunId::new(); - let location = ObjectPath::from(format!("artifacts/{run}/captures/elsewhere/bad")); - objects - .put(&location, Bytes::from_static(b"bad").into()) + let legacy = ArtifactKey::new(StageId::new("build", 1), 1, "report.txt"); + store.put(&run, &legacy, b"legacy").await.unwrap(); + for location in [ + format!("artifacts/{run}/captures/elsewhere/other"), + format!("artifacts/{run}/captures/stray"), + ] { + objects + .put( + &ObjectPath::from(location), + Bytes::from_static(b"other").into(), + ) + .await + .unwrap(); + } + let listed = store.list_for_run(&run).await.unwrap(); + assert_eq!(listed.len(), 1); + assert_eq!(listed[0].filename, "report.txt"); + } + + #[tokio::test] + async fn put_capture_rejects_content_that_does_not_match_its_digest() { + let objects: Arc = Arc::new(InMemory::new()); + let store = ArtifactStore::new(objects.clone(), "artifacts"); + let run = RunId::new(); + let hash = fabro_types::BlobHash::new(b"expected"); + let err = store + .put_capture(&run, &hash, Bytes::from_static(b"other")) + .await + .unwrap_err(); + assert!(matches!(err, Error::CaptureDigestMismatch { expected } if expected == hash)); + assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); + } + + #[tokio::test] + async fn delete_capture_removes_content_and_tolerates_absence() { + let store = test_store(); + let run = RunId::new(); + let hash = fabro_types::BlobHash::new(b"content"); + store + .put_capture(&run, &hash, Bytes::from_static(b"content")) .await .unwrap(); - assert!(store.list_for_run(&run).await.is_err()); + store.delete_capture(&run, &hash).await.unwrap(); + assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); + store.delete_capture(&run, &hash).await.unwrap(); } #[tokio::test] diff --git a/lib/components/fabro-store/src/error.rs b/lib/components/fabro-store/src/error.rs index 49db7cf24..966448ed5 100644 --- a/lib/components/fabro-store/src/error.rs +++ b/lib/components/fabro-store/src/error.rs @@ -26,6 +26,8 @@ pub enum Error { BlobHashConflict { blob_hash: BlobHash }, #[error("stored blob data does not match requested hash {blob_hash}")] BlobIntegrity { blob_hash: BlobHash }, + #[error("captured content does not match its digest {expected}")] + CaptureDigestMismatch { expected: BlobHash }, #[error("I/O error: {0}")] Io(#[from] std::io::Error), #[error("Invalid event payload: {0}")] diff --git a/lib/components/fabro-store/src/platform_records.rs b/lib/components/fabro-store/src/platform_records.rs index 11e1742ba..2424740c9 100644 --- a/lib/components/fabro-store/src/platform_records.rs +++ b/lib/components/fabro-store/src/platform_records.rs @@ -450,7 +450,6 @@ pub struct CheckpointRecord { /// One file collected from a stage's workspace after its attempt finished. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(try_from = "ArtifactCollectedWire")] pub struct ArtifactCollectedRecord { pub execution: u64, pub firing: u64, @@ -458,49 +457,58 @@ pub struct ArtifactCollectedRecord { pub attempt: u32, /// The file's path relative to the workspace root. pub path: String, - /// Where the file's bytes are stored. - #[serde(flatten)] + /// Where the file's bytes are stored. Its hash is the SHA-256 of the + /// bytes: with `path`, the identity a later capture of the same + /// unchanged file is matched by. The wire also carries it as `digest`. + #[serde(flatten, with = "source_with_digest")] pub source: ArtifactSource, pub bytes: u64, - /// The SHA-256 of the bytes as lowercase hex: with `path`, the identity - /// a later capture of the same unchanged file is matched by. - pub digest: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub operation: Option, } -/// Private decoding boundary for validating the redundant legacy checksum. -#[derive(Deserialize)] -struct ArtifactCollectedWire { - execution: u64, - firing: u64, - attempt: u32, - path: String, - #[serde(flatten)] - source: ArtifactSource, - bytes: u64, - digest: String, - #[serde(default)] - operation: Option, -} +/// The record's `digest` key repeats its source's hash as lowercase hex, as +/// every earlier record wrote it. It is derived on write and checked on read. +mod source_with_digest { + use fabro_types::ArtifactSource; + use serde::de::Error as _; + use serde::{Deserialize, Deserializer, Serialize, Serializer}; -impl TryFrom for ArtifactCollectedRecord { - type Error = &'static str; + #[derive(Serialize)] + struct Written<'a> { + #[serde(flatten)] + source: &'a ArtifactSource, + digest: String, + } - fn try_from(wire: ArtifactCollectedWire) -> std::result::Result { - if wire.digest != wire.source.hash().to_string() { - return Err("artifact checksum does not match its payload source"); + #[derive(Deserialize)] + struct Read { + #[serde(flatten)] + source: ArtifactSource, + digest: String, + } + + pub(super) fn serialize( + source: &ArtifactSource, + serializer: S, + ) -> Result { + Written { + source, + digest: source.hash().to_string(), } - Ok(Self { - execution: wire.execution, - firing: wire.firing, - attempt: wire.attempt, - path: wire.path, - source: wire.source, - bytes: wire.bytes, - digest: wire.digest, - operation: wire.operation, - }) + .serialize(serializer) + } + + pub(super) fn deserialize<'de, D: Deserializer<'de>>( + deserializer: D, + ) -> Result { + let read = Read::deserialize(deserializer)?; + if read.digest != read.source.hash().to_string() { + return Err(D::Error::custom( + "artifact checksum does not match its payload source", + )); + } + Ok(read.source) } } @@ -887,7 +895,6 @@ mod tests { path: "assets/report.txt".to_string(), source: ArtifactSource::SqliteBlob(BlobHash::new(b"report")), bytes: 6, - digest: BlobHash::new(b"report").to_string(), operation: Some(OperationKey { execution: 0, decision: DecisionRef::AttemptStart { diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index c52f0dd39..16c11a194 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1847,14 +1847,16 @@ impl Client { &self, run_id: &RunId, digest: &BlobHash, - data: &[u8], + data: Bytes, ) -> Result<()> { + let data = &data; self.send_api(|client| async move { client .write_run_artifact_content() .id(run_id.to_string()) .digest(*digest) - .body(data.to_vec()) + // A `Bytes` clone shares the buffer, including on a retry. + .body(data.clone()) .send() .await }) diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 62a007122..c4c4ee339 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -59,7 +59,9 @@ pub struct RunProjection { pub pending_interviews: BTreeMap, /// The files collected from the run's workspaces under /// `[run.artifacts] include`, one entry per capture, in the order they - /// were recorded. The bytes are in the blob table under `blob`. + /// were recorded. Each entry's `source` says where its bytes are: the + /// configured artifact store for `object`, the SQLite blob table for + /// earlier captures under `blob`. #[serde(default, skip_serializing_if = "Vec::is_empty")] pub artifacts: Vec, stages: HashMap,