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 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-09-28 14:30:36 -04:00
parent ead2ca53b8
commit 7f55dd87b5
12 changed files with 438 additions and 148 deletions

View file

@ -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 {

View file

@ -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<Arc<AppState>> {
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();

View file

@ -456,10 +456,7 @@ pub(crate) async fn execute(state: Arc<AppState>, 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(),

View file

@ -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

View file

@ -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)
}

View file

@ -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<dyn ExecEnv>);
/// 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<Mutex<HashSet<ArtifactIdentity>>>,
/// One lock per identity being captured, so concurrent transitions that
/// leave the same file upload and record it once.
capturing: Mutex<HashMap<ArtifactIdentity, Arc<AsyncMutex<()>>>>,
}
/// 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<bool, HookError> {
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<HashSet<ArtifactIdentity>>, 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<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError> {
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))

View file

@ -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()

View file

@ -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<ObjectPath> {
Ok(self.run_prefix(run_id)?.child("captures"))
}
fn capture_prefix(&self, run_id: &RunId) -> Result<ObjectPath> {
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<ObjectPath> {
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<Vec<NodeArtifact>> {
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<dyn ObjectStore> = 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<dyn ObjectStore> = 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]

View file

@ -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}")]

View file

@ -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<OperationKey>,
}
/// 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<OperationKey>,
}
/// 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<ArtifactCollectedWire> 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<Self, Self::Error> {
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<S: Serializer>(
source: &ArtifactSource,
serializer: S,
) -> Result<S::Ok, S::Error> {
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<ArtifactSource, D::Error> {
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 {

View file

@ -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
})

View file

@ -59,7 +59,9 @@ pub struct RunProjection {
pub pending_interviews: BTreeMap<String, PendingInterviewRecord>,
/// 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<RunArtifact>,
stages: HashMap<StageId, StageProjection>,