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 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-09-25 13:32:13 -04:00
parent 71b08b61a1
commit 2e500b0e2f
13 changed files with 161 additions and 197 deletions

View file

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

View file

@ -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<Arc<AppState>>,
body: Result<Bytes, BytesRejection>,
) -> 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::<BlobHash>() 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);
}
}

View file

@ -456,6 +456,10 @@ 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,
)),
);
let request = RunRequest {
run_id: run_id.to_string(),
@ -480,10 +484,6 @@ pub(crate) async fn execute(state: Arc<AppState>, 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;

View file

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

View file

@ -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<BlobHash, ArtifactWriteError>;
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<BlobHash, ArtifactWriteError> {
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<BlobHash, ArtifactWriteError> {
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)
}
}

View file

@ -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<dyn RunStore>,
pub runtime: RuntimeSpec,
pub store: Arc<dyn RunStore>,
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<dyn Interviewer>,
pub interviewer: Arc<dyn Interviewer>,
/// The caller's observers of every record, registered ahead of the
/// interview dispatcher: the interviewer's own expiry observer among
/// them.
pub observers: Vec<Arc<dyn ExecutionObserver>>,
pub observers: Vec<Arc<dyn ExecutionObserver>>,
/// Where `{{ secrets.NAME }}` references resolve from; `None` leaves
/// every secret unknown.
pub secrets: Option<Arc<dyn SecretProvider>>,
pub secrets: Option<Arc<dyn SecretProvider>>,
/// Where offloaded stage values go; `None` keeps Petri's local store
/// under the run directory.
pub blobs: Option<Arc<dyn Blobs>>,
/// Where captured workspace files go. Required when capture writes files.
pub artifact_writer: Option<Arc<dyn ArtifactWriter>>,
pub blobs: Option<Arc<dyn Blobs>>,
/// Fabro's hooks: the checkpoint commit and its record. `None` runs
/// with Petri's local hook service alone.
pub hooks: Option<HooksSpec>,
pub hooks: Option<HooksSpec>,
}
/// The recorded status of a finished run.
@ -213,7 +210,6 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
Arc::clone(&request.store),
resumed,
request.blobs.clone(),
request.artifact_writer.clone(),
))
});
if let Some(hooks) = &fabro_hooks {

View file

@ -185,8 +185,6 @@ pub enum HookError {
},
#[error("invalid run.artifacts.include pattern")]
Globs(#[source] Arc<WorkspaceGlobError>),
#[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<dyn PlatformRecords>,
pub git: RunGitSettings,
pub records: Arc<dyn PlatformRecords>,
pub git: RunGitSettings,
/// The `[run.artifacts] include` patterns: which files of a stage's
/// workspace are collected after the stage.
pub artifacts: Vec<String>,
pub artifacts: Vec<String>,
/// 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<PathBuf>,
pub test_gates: Option<PathBuf>,
/// Where captured workspace files go.
pub artifact_writer: Arc<dyn ArtifactWriter>,
}
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<dyn PlatformRecords>, settings: &RunNamespace) -> Self {
pub fn for_run(
records: Arc<dyn PlatformRecords>,
settings: &RunNamespace,
artifact_writer: Arc<dyn ArtifactWriter>,
) -> 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<dyn PlatformRecords>,
/// Where diff patches go; `None` records summaries alone.
blobs: Option<Arc<dyn Blobs>>,
artifact_writer: Option<Arc<dyn ArtifactWriter>>,
artifact_writer: Arc<dyn ArtifactWriter>,
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<dyn RunStore>,
resumed: bool,
blobs: Option<Arc<dyn Blobs>>,
artifact_writer: Option<Arc<dyn ArtifactWriter>>,
) -> 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<bool, HookError> {
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<BlobHash, ArtifactWriteError> {
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<dyn PlatformRecords>,
writer: Option<Arc<dyn ArtifactWriter>>,
artifact_writer: Arc<dyn ArtifactWriter>,
) -> 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))

View file

@ -181,13 +181,17 @@ impl Harness {
fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec {
HooksSpec {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
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<dyn Blobs>),
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");

View file

@ -121,7 +121,6 @@ pub(crate) fn run_request(
interviewer: Arc::new(interviewer),
secrets: None,
blobs: None,
artifact_writer: None,
hooks: None,
}
}

View file

@ -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<BlobHash> {
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<Option<Bytes>> {
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<ObjectPath> {
Ok(self.run_prefix(run_id)?.child("captures").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?;
Ok(())
}
async fn get_at(&self, path: &ObjectPath) -> Result<Option<Bytes>> {
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<BufWriter> {
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<Option<Bytes>> {
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!(

View file

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

View file

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

View file

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