From a1926c863257130ebf0f9655764beb9382803481 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Thu, 24 Sep 2026 13:53:46 -0400 Subject: [PATCH 1/3] Restore configured artifact storage for workflow captures --- Cargo.lock | 1 + docs/internal/run-directory-keys.md | 12 +- docs/public/agents/outputs.mdx | 6 +- docs/public/api-reference/fabro-api.yaml | 112 ++++++ .../src/commands/run/petri_worker.rs | 5 + .../fabro-cli/tests/it/scenario/artifacts.rs | 121 +++++- .../src/server/handler/artifacts.rs | 176 ++++++++- .../fabro-server/src/server/petri_runs.rs | 5 + lib/apps/fabro-server/src/server/tests.rs | 5 +- .../src/server/tests/artifact_storage.rs | 288 ++++++++++++++ lib/components/fabro-petri/Cargo.toml | 1 + lib/components/fabro-petri/README.md | 10 +- lib/components/fabro-petri/VIEWS.md | 2 +- lib/components/fabro-petri/src/artifacts.rs | 73 ++++ lib/components/fabro-petri/src/engine.rs | 30 +- lib/components/fabro-petri/src/hooks.rs | 356 +++++++++++++++--- lib/components/fabro-petri/src/lib.rs | 1 + .../fabro-petri/src/projection/platform.rs | 2 +- lib/components/fabro-petri/tests/hooks.rs | 77 ++-- .../fabro-petri/tests/support/mod.rs | 1 + .../fabro-store/src/artifact_store.rs | 95 ++++- .../fabro-store/src/platform_records.rs | 71 +++- lib/foundation/fabro-api/build.rs | 2 + lib/foundation/fabro-api/src/lib.rs | 35 +- .../fabro-api/tests/artifact_round_trip.rs | 67 ++++ lib/foundation/fabro-client/src/client.rs | 20 + .../fabro-types/src/artifact_source.rs | 55 +++ lib/foundation/fabro-types/src/lib.rs | 2 + .../fabro-types/src/run_projection.rs | 51 ++- .../src/.openapi-generator/FILES | 2 + .../src/api/run-internals-api.ts | 89 +++++ .../src/models/artifact-source.ts | 29 ++ .../fabro-api-client/src/models/index.ts | 2 + .../src/models/run-artifact.ts | 36 ++ .../src/models/run-projection.ts | 7 + 35 files changed, 1710 insertions(+), 137 deletions(-) create mode 100644 lib/apps/fabro-server/src/server/tests/artifact_storage.rs create mode 100644 lib/components/fabro-petri/src/artifacts.rs create mode 100644 lib/foundation/fabro-api/tests/artifact_round_trip.rs create mode 100644 lib/foundation/fabro-types/src/artifact_source.rs create mode 100644 lib/packages/fabro-api-client/src/models/artifact-source.ts create mode 100644 lib/packages/fabro-api-client/src/models/run-artifact.ts diff --git a/Cargo.lock b/Cargo.lock index 424c56a68..9b7010719 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2539,6 +2539,7 @@ dependencies = [ "fabro-workflow", "httpmock", "lithos-llm", + "object_store", "pebble-coding-agent", "petri-attractor-steps", "petri-execution", diff --git a/docs/internal/run-directory-keys.md b/docs/internal/run-directory-keys.md index b248e84a6..e9daf0ca4 100644 --- a/docs/internal/run-directory-keys.md +++ b/docs/internal/run-directory-keys.md @@ -37,5 +37,15 @@ These names are still real, but they are no longer live scratch files by default ## Notes -- Artifact binaries are no longer stored in the SlateDB keyspace. They live in `ArtifactStore`; the run scratch tree only contains local cached copies when a workflow stage writes them to disk. +- New captured artifact binaries live in `ArtifactStore`, with originals retained in the sandbox. Historical SQLite captures remain in the blob table. - Final diffs for checkpointed runs are projected from the run store; they are no longer written as scratch files. + +## Captured artifact content + +New automatic captures use `//captures/sha256/` +in the configured artifact store. Platform records and the run projection hold +the stage, retry, relative path and content source (`object` for this layout, +`blob` for historical SQLite captures). Content objects alone are not listing +entries. Historical stage-keyed objects remain readable. Run deletion removes +both object layouts; generic blobs and checkpoint patches keep their SQLite +storage contract. diff --git a/docs/public/agents/outputs.mdx b/docs/public/agents/outputs.mdx index 977345f51..5337b21b8 100644 --- a/docs/public/agents/outputs.mdx +++ b/docs/public/agents/outputs.mdx @@ -254,9 +254,11 @@ When `[run.artifacts]` contains include patterns, Fabro scans the sandbox after 1. Fabro compiles and validates the configured workspace-relative globs. 2. The sandbox provider enumerates regular files and their sizes without recursing through symlinks below the workspace root. -3. Fabro applies the globs to normalized relative paths, enforces its collection limits, and downloads the selected files. +3. Fabro applies the globs to normalized relative paths, enforces its collection limits, and copies selected files to the local directory or S3 bucket configured by `server.artifacts`. Originals remain in the sandbox. -Each scan represents the post-stage workspace state; Fabro does not depend on filesystem modification timestamps. The same path and content hash is recorded only once per run, even when it still matches after later stages. Individual files over 10 MB are skipped, and each collection is limited to 100 files and 50 MB total. +Each scan represents the post-stage workspace state; Fabro does not depend on filesystem modification timestamps. The same path and content hash is recorded only once per run, even when it still matches after later stages. Individual files over 10 MiB (10,485,760 bytes) are skipped, and each collection is limited to 100 files and 50 MiB total. A file exactly 10 MiB can be captured. + +New artifact bytes use the configured artifact backend, with stage and filename metadata in SQLite. Existing captures stored as SQLite blobs remain readable without moving their bytes. Generic offloaded context values and checkpoint patches continue to use SQLite blobs. ### What gets captured diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 368751c38..ec066dc94 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -3085,6 +3085,76 @@ paths: schema: $ref: "#/components/schemas/ErrorResponse" + /api/v1/runs/{id}/artifacts/content/{digest}: + put: + operationId: writeRunArtifactContent + tags: [Run Internals] + summary: Write Captured Artifact Content + description: > + Stores captured file bytes in the configured artifact backend for this run. + Requires a worker token belonging to the run. The body is limited to + 10 MiB and must match the SHA-256 digest. Repeating the same upload is + safe. Uploading content alone does not create an artifact listing entry. + parameters: + - $ref: "#/components/parameters/RunId" + - name: digest + in: path + required: true + schema: + $ref: "#/components/schemas/BlobHash" + requestBody: + required: true + content: + application/octet-stream: + schema: + type: string + format: binary + responses: + "204": + description: Complete content stored + "400": + description: Invalid digest or digest does not match the body + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "401": + description: Authentication required + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "403": + description: A worker token belonging to this run is required + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "404": + description: Run not found + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "409": + description: Run is archived + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "413": + description: Content exceeds 10 MiB + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "500": + description: Artifact storage failed + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + /api/v1/runs/{id}/blobs: post: operationId: writeRunBlob @@ -12898,6 +12968,43 @@ components: diff: $ref: "#/components/schemas/RunDiff" + ArtifactSource: + description: Exactly one payload source; object content is owned by the containing run. + type: object + properties: + blob: + $ref: "#/components/schemas/BlobHash" + object: + $ref: "#/components/schemas/BlobHash" + oneOf: + - required: [blob] + - required: [object] + + RunArtifact: + description: A captured workspace file with its durable payload source. + type: object + oneOf: + - required: [blob] + - required: [object] + required: [stage_id, retry, relative_path, size] + properties: + blob: + $ref: "#/components/schemas/BlobHash" + object: + $ref: "#/components/schemas/BlobHash" + stage_id: + $ref: "#/components/schemas/StageId" + retry: + type: integer + format: uint32 + minimum: 1 + relative_path: + type: string + size: + type: integer + format: uint64 + minimum: 0 + RunProjection: description: Raw internal run projection derived from the event log. type: object @@ -12940,6 +13047,11 @@ components: oneOf: - $ref: "#/components/schemas/RunControlAction" - type: "null" + artifacts: + type: array + description: Captured files; older projections may omit this field. + items: + $ref: "#/components/schemas/RunArtifact" checkpoints: type: array description: Sequence-tagged checkpoint history entries. 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 681862e6d..9bfc188ac 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -71,6 +71,7 @@ use fabro_auth::VaultCredentialSource; use fabro_client::{Client, ServerTarget}; use fabro_interview::{ControlInterviewer, WorkerControlMessage, WorkerControlOutcome}; use fabro_llm::credentials::{CredentialProvider, readiness}; +use fabro_petri::artifacts::ClientArtifactWriter; use fabro_petri::blobs::ClientBlobs; use fabro_petri::controls::{RunControls, SteerError}; use fabro_petri::engine::{self, Conclusion, Execution, RunRequest}; @@ -222,6 +223,10 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { worker.client.clone_for_reuse(), run_id, ))), + artifact_writer: Some(Arc::new(ClientArtifactWriter::new( + worker.client.clone_for_reuse(), + run_id, + ))), hooks: Some(hooks), }; let paused_mirror = mirror_paused_state(run_id, &controls, Arc::clone(&records)); diff --git a/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs b/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs index 234c4fe61..81f327746 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs @@ -1,15 +1,132 @@ //! `fabro artifact list` and `fabro artifact cp` over a run whose artifacts //! the engine's hooks collected: every file under `[run.artifacts] include` -//! in a stage's workspace, once per content, into the blob table. +//! in a stage's workspace, once per content, into configured artifact storage. use std::path::PathBuf; use std::time::Duration; use fabro_test::{fabro_snapshot, test_context}; -use super::petri::{RunningServer, host_plugin, run_detached, wait_for_success}; +use super::petri::{RunningServer, host_plugin, run_detached, run_json, wait_for_success}; use crate::cmd::support::{read_text, text_tree}; +#[tokio::test(flavor = "multi_thread")] +async fn artifact_worker_captures_large_files_in_the_configured_local_store() { + if host_plugin().is_none() { + return; + } + let context = test_context!(); + let server = RunningServer::start_with( + "\n[server.artifacts]\nprovider = \"local\"\nprefix = \"selected-prefix\"\n", + &[], + ) + .await; + // RunningServer explicitly selects --storage-dir, which also selects the + // local artifact root. Inspect that resolved backend, outside the sandbox. + let objects = server.storage_dir.join("objects/artifacts"); + let workspace = artifact_workspace(&context); + tokio::fs::write(workspace.join("workflow.fabro"), r#"digraph Capture { + graph [goal="Capture binary files", default_max_retries=0] + start [shape=Mdiamond] + write [shape=parallelogram, script="mkdir -p assets && dd if=/dev/zero of=assets/medium.bin bs=1048576 count=3 && cp assets/medium.bin assets/same.bin && dd if=/dev/zero of=assets/limit.bin bs=1048576 count=10 && cp assets/limit.bin assets/skipped.bin && printf x >> assets/skipped.bin"] + keep [shape=parallelogram, script="test -f assets/skipped.bin"] + exit [shape=Msquare] + start -> write -> keep -> exit + }"#).await.unwrap(); + let run_id = run_detached(&context, &server, &workspace); + wait_for_success(&server, &run_id).await; + let projection = run_json(&server, &format!("runs/{run_id}/state")).await; + let artifacts = projection["artifacts"].as_array().unwrap(); + assert_eq!( + artifacts.len(), + 3, + "unchanged files are captured once; oversize is skipped" + ); + let database = + fabro_db::Database::connect(fabro_config::Storage::new(&server.storage_dir).sqlite_path()) + .await + .unwrap(); + let blobs = fabro_store::BlobStore::new(database.clone_pool()); + let (rebuilt, _, _) = fabro_petri::test_support::rebuild( + database.pool(), + database.pool(), + run_id.parse().unwrap(), + ) + .await + .unwrap(); + assert_eq!( + serde_json::to_value(rebuilt.unwrap().artifacts).unwrap(), + projection["artifacts"] + ); + + for (path, size) in [ + ("medium.bin", 3 * 1024 * 1024), + ("same.bin", 3 * 1024 * 1024), + ("limit.bin", 10 * 1024 * 1024), + ] { + let bytes = vec![0; size]; + let hash = fabro_types::BlobHash::new(&bytes); + let capture = artifacts + .iter() + .find(|entry| entry["relative_path"] == format!("assets/{path}")) + .unwrap(); + assert_eq!(capture["object"], hash.to_string()); + assert!(capture.get("blob").is_none()); + assert_eq!( + tokio::fs::read( + objects.join(format!("selected-prefix/{run_id}/captures/sha256/{hash}")) + ) + .await + .unwrap(), + bytes + ); + assert!(blobs.read(&hash).await.unwrap().is_none()); + let destination = context.temp_dir.join(format!("download-{path}")); + let output = context + .command() + .args([ + "artifact", + "cp", + &format!("{run_id}:assets/{path}"), + destination.to_str().unwrap(), + "--server", + &server.target(), + ]) + .output() + .unwrap(); + assert!(output.status.success(), "artifact download failed"); + assert_eq!( + tokio::fs::read(destination.join(path)).await.unwrap(), + bytes + ); + } + let mut scopes = tokio::fs::read_dir(server.petri_run_dir(&run_id).join("scopes")) + .await + .unwrap(); + let original = scopes + .next_entry() + .await + .unwrap() + .unwrap() + .path() + .join("work/assets"); + assert_eq!( + tokio::fs::metadata(original.join("limit.bin")) + .await + .unwrap() + .len(), + 10 * 1024 * 1024 + ); + assert_eq!( + tokio::fs::metadata(original.join("skipped.bin")) + .await + .unwrap() + .len(), + 10 * 1024 * 1024 + 1 + ); + server.shutdown(); +} + /// Three command stages that leave files under `assets/`. The second and /// third write different contents to the same path, so the path names an /// artifact of each; the third also writes a `summary.txt` that collides diff --git a/lib/apps/fabro-server/src/server/handler/artifacts.rs b/lib/apps/fabro-server/src/server/handler/artifacts.rs index bf505e4c5..91a73675e 100644 --- a/lib/apps/fabro-server/src/server/handler/artifacts.rs +++ b/lib/apps/fabro-server/src/server/handler/artifacts.rs @@ -5,9 +5,12 @@ use std::sync::Arc; use async_zip::base::write::ZipFileWriter; use async_zip::error::ZipError; use async_zip::{Compression, ZipEntryBuilder}; +use axum::extract::DefaultBodyLimit; +use axum::extract::rejection::BytesRejection; use axum::http::HeaderValue; +use axum::routing::put; use fabro_store::{ArtifactStore, BlobStore, Error as StoreError}; -use fabro_types::{BlobHash, RunProjection}; +use fabro_types::{ARTIFACT_MAX_FILE_BYTES, ArtifactSource, BlobHash, RunProjection}; use fabro_util::error::collect_chain; use futures_util::SinkExt as _; use futures_util::io::AsyncWriteExt as _; @@ -26,9 +29,14 @@ use super::super::{ octet_stream_response, parse_run_id_path, parse_stage_id_path, post, reject_if_archived, required_query_param, validate_relative_artifact_path, }; +use crate::principal_middleware::RequireWorkerRunSegment; 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)), + ) .route("/runs/{id}/blobs", post(write_run_blob)) .route("/runs/{id}/blobs/{blobHash}", get(read_run_blob)) .route("/runs/{id}/artifacts", get(list_run_artifacts)) @@ -43,6 +51,47 @@ pub(super) fn routes() -> Router> { ) } +async fn write_run_artifact_content( + RequireWorkerRunSegment(id, digest): RequireWorkerRunSegment, + State(state): State>, + body: Result, +) -> 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.", + ) + .into_response(); + } + }; + if let Some(response) = reject_if_archived(state.as_ref(), &id).await { + return response; + } + if let Err(error) = state.load_run_projection(&id).await { + return error.into_response(); + } + let Ok(expected) = digest.parse::() else { + return ApiError::bad_request("Invalid artifact digest.").into_response(); + }; + if BlobHash::new(&body) != expected { + return ApiError::bad_request("Artifact content does not match its digest.") + .into_response(); + } + match state.artifact_store.put_capture(&id, &body).await { + Ok(_) => StatusCode::NO_CONTENT.into_response(), + Err(error) => { + warn!(run_id = %id, error = %collect_chain(&error).join(": "), "Artifact upload failed"); + ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + "Artifact storage failed.", + ) + .into_response() + } + } +} + #[derive(serde::Deserialize)] struct ArtifactFilenameParams { #[serde(default)] @@ -86,18 +135,16 @@ async fn read_run_blob( } } -/// Where an artifact's bytes are: the blob table, for one the run's -/// hooks collected, or the artifact store, for one written there directly. +/// Recorded capture sources and historical stage-keyed objects. #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum ArtifactBytes { - Blob(BlobHash), + Captured(ArtifactSource), Store, } /// Every artifact of the run, each once: the ones the run's projection -/// records, with their bytes in the blob table, and the ones written to -/// the artifact store. A path in the store for a stage and retry the -/// projection also collected is the projection's. +/// records, and historical stage-keyed objects. A path in the store for a stage +/// and retry the projection also collected is the projection's. async fn run_artifacts( state: &AppState, run_id: &RunId, @@ -117,7 +164,7 @@ async fn run_artifacts( filename: artifact.relative_path.clone(), size: artifact.size, }, - ArtifactBytes::Blob(artifact.blob), + ArtifactBytes::Captured(artifact.source), )); } let uploaded = state @@ -255,7 +302,10 @@ async fn read_artifact( bytes: ArtifactBytes, ) -> Result, StoreError> { match bytes { - ArtifactBytes::Blob(hash) => blobs.read(&hash).await, + ArtifactBytes::Captured(ArtifactSource::SqliteBlob(hash)) => blobs.read(&hash).await, + ArtifactBytes::Captured(ArtifactSource::ObjectStore(hash)) => { + artifact_store.get_capture(run_id, &hash).await + } ArtifactBytes::Store => artifact_store.get(run_id, key).await, } } @@ -484,7 +534,7 @@ async fn get_stage_artifact( && artifact.relative_path == key.relative_path }) .map_or(ArtifactBytes::Store, |artifact| { - ArtifactBytes::Blob(artifact.blob) + ArtifactBytes::Captured(artifact.source) }); match read_artifact( &state.artifact_store, @@ -502,3 +552,109 @@ async fn get_stage_artifact( } } } + +#[cfg(test)] +mod tests { + use async_zip::base::read::mem::ZipFileReader; + use fabro_types::{RunArtifact, StageId, test_support}; + + use super::*; + + #[tokio::test] + async fn artifact_mixed_sources_keep_precedence_and_zip_contents_without_fallback() { + let state = crate::test_support::test_app_state(); + let run = RunId::new(); + let stage = StageId::new("write", 1); + let mut projection = RunProjection::new( + "capture".to_string(), + test_support::test_run_spec(), + chrono::Utc::now(), + ); + let blobs = state.store_ref().blobs(); + let old = blobs.write(b"old SQLite").await.unwrap(); + let new = state + .artifact_store + .put_capture(&run, b"new object") + .await + .unwrap(); + for (path, size, source) in [ + ("old.bin", 10, ArtifactSource::SqliteBlob(old)), + ("new.bin", 10, ArtifactSource::ObjectStore(new)), + ] { + projection.artifacts.push(RunArtifact { + stage_id: stage.clone(), + retry: 1, + relative_path: path.to_string(), + size, + source, + }); + } + // A historical object at the same logical path loses to the recorded capture. + for (path, bytes) in [ + ("new.bin", b"shadow".as_slice()), + ("legacy.bin", b"legacy".as_slice()), + ] { + state + .artifact_store + .put(&run, &ArtifactKey::new(stage.clone(), 1, path), bytes) + .await + .unwrap(); + } + // Unrecorded content is never a file-list entry. + state + .artifact_store + .put_capture(&run, b"orphan") + .await + .unwrap(); + let entries = run_artifacts(&state, &run, &projection).await.unwrap(); + assert_eq!(entries.len(), 3); + let archive = artifact_archive_body( + state.artifact_store.clone(), + blobs.clone(), + run, + latest_run_artifacts(entries, &projection), + ); + let bytes = axum::body::to_bytes(archive, 1024 * 1024).await.unwrap(); + let zip = ZipFileReader::new(bytes.to_vec()).await.unwrap(); + let mut contents = BTreeMap::new(); + for (index, entry) in zip.file().entries().iter().enumerate() { + let mut bytes = Vec::new(); + zip.reader_with_entry(index) + .await + .unwrap() + .read_to_end_checked(&mut bytes) + .await + .unwrap(); + contents.insert(entry.filename().as_str().unwrap().to_string(), bytes); + } + assert_eq!( + contents, + BTreeMap::from([ + ("legacy.bin".to_string(), b"legacy".to_vec()), + ("new.bin".to_string(), b"new object".to_vec()), + ("old.bin".to_string(), b"old SQLite".to_vec()), + ]) + ); + let missing = blobs + .write(b"SQLite is not the selected source") + .await + .unwrap(); + let key = ArtifactKey::new(stage, 1, "new.bin"); + assert!( + read_artifact( + &state.artifact_store, + &blobs, + &run, + &key, + ArtifactBytes::Captured(ArtifactSource::ObjectStore(missing)) + ) + .await + .unwrap() + .is_none() + ); + // Saved projections preserve both physical sources without rewriting history. + let saved = serde_json::to_value(&projection).unwrap(); + let decoded: RunProjection = serde_json::from_value(saved).unwrap(); + assert_eq!(decoded.artifacts, projection.artifacts); + } +} diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index 5a54f187a..be500bb5d 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -40,6 +40,7 @@ use fabro_config::{ EnvironmentImageLayer, EnvironmentLayer, Home, MergeMap, SettingsLayer, Storage, }; use fabro_interview::ControlInterviewer; +use fabro_petri::artifacts::StoreArtifactWriter; use fabro_petri::controls::RunControls; use fabro_petri::engine::{self, Conclusion, Execution, RunRequest}; use fabro_petri::hooks::HooksSpec; @@ -479,6 +480,10 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { observers, secrets: Some(Arc::new(VaultSecrets::from_vault(&vault))), blobs: Some(state.store_ref().blobs()), + artifact_writer: Some(Arc::new(StoreArtifactWriter::new( + state.artifact_store.clone(), + run_id, + ))), hooks: Some(hooks), }; let result = Box::pin(engine::run(request)).await; diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index c6262bb39..2758f8741 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -4958,8 +4958,7 @@ fn named_workflow_dot(name: &str, goal: &str) -> String { ) } -/// Write one artifact for a stage the way the hooks do: straight into the -/// artifact store. +/// Seed a historical stage-keyed object for artifact reader tests. async fn seed_stage_artifact( state: &AppState, run_id: &str, @@ -10974,3 +10973,5 @@ fn an_unsettled_run_takes_the_stores_status_or_a_termination_at_worker_exit() { RunStatus::Running ); } + +mod artifact_storage; diff --git a/lib/apps/fabro-server/src/server/tests/artifact_storage.rs b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs new file mode 100644 index 000000000..cb1546057 --- /dev/null +++ b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs @@ -0,0 +1,288 @@ +use std::io; + +use axum::http::HeaderValue; +use fabro_petri::artifacts::{ArtifactWriter, ClientArtifactWriter}; +use fabro_types::ARTIFACT_MAX_FILE_BYTES; + +use super::*; + +fn upload(run: RunId, digest: &str, token: &str, body: Body) -> Request { + let mut request = bearer_request( + Method::PUT, + &format!("/runs/{run}/artifacts/content/{digest}"), + token, + body, + ); + request.headers_mut().insert( + header::CONTENT_TYPE, + HeaderValue::from_static("application/octet-stream"), + ); + request +} + +#[tokio::test] +async fn artifact_upload_accepts_large_and_exact_limit_content_without_sqlite_or_listing_entries() { + let (state, app) = jwt_auth_app(); + let user = issue_test_user_jwt(); + let run = create_run_with_bearer(&app, &user).await; + let worker = issue_test_worker_token(&run); + for size in [3 * 1024 * 1024, ARTIFACT_MAX_FILE_BYTES] { + let bytes = vec![0xa3; size]; + let hash = BlobHash::new(&bytes); + for _ in 0..2 { + let response = app + .clone() + .oneshot(upload( + run, + &hash.to_string(), + &worker, + Body::from(bytes.clone()), + )) + .await + .unwrap(); + assert_status!(response, StatusCode::NO_CONTENT).await; + } + assert_eq!( + state + .artifact_store + .get_capture(&run, &hash) + .await + .unwrap() + .unwrap(), + bytes + ); + assert!( + state + .store_ref() + .blobs() + .read(&hash) + .await + .unwrap() + .is_none() + ); + } + assert!( + state + .artifact_store + .list_for_run(&run) + .await + .unwrap() + .is_empty() + ); + let response = app + .oneshot(bearer_request( + Method::GET, + &format!("/runs/{run}/artifacts"), + &user, + Body::empty(), + )) + .await + .unwrap(); + let body = response_json!(response, StatusCode::OK).await; + assert_eq!(body["data"], json!([])); +} + +#[tokio::test] +async fn artifact_upload_rejects_unauthorized_invalid_and_oversized_bodies_without_overwriting() { + let (state, app) = jwt_auth_app(); + let user = issue_test_user_jwt(); + 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 + .artifact_store + .put_capture(&run, b"valid content") + .await + .unwrap(); + let digest = hash.to_string(); + for (token, expected) in [ + (&user, StatusCode::FORBIDDEN), + (&other_worker, StatusCode::FORBIDDEN), + ] { + let response = app + .clone() + .oneshot(upload(run, &digest, token, Body::from("invalid content"))) + .await + .unwrap(); + assert_status!(response, expected).await; + } + let mut missing = upload(run, &digest, &worker, Body::empty()); + missing.headers_mut().remove(header::AUTHORIZATION); + assert_status!( + app.clone().oneshot(missing).await.unwrap(), + StatusCode::UNAUTHORIZED + ) + .await; + for (digest, bytes) in [(&digest[..], "mismatch"), ("invalid", "valid content")] { + let response = app + .clone() + .oneshot(upload(run, digest, &worker, Body::from(bytes))) + .await + .unwrap(); + assert_status!(response, StatusCode::BAD_REQUEST).await; + } + // Stream chunks without a content length: enforce bytes actually read. + let chunks = futures_util::stream::iter([ + Ok::<_, io::Error>(Bytes::from(vec![0; ARTIFACT_MAX_FILE_BYTES])), + Ok(Bytes::from_static(b"x")), + ]); + let response = app + .clone() + .oneshot(upload(run, &digest, &worker, Body::from_stream(chunks))) + .await + .unwrap(); + assert_status!(response, StatusCode::PAYLOAD_TOO_LARGE).await; + let broken = futures_util::stream::iter([ + Ok(Bytes::from_static(b"valid")), + Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "truncated test body", + )), + ]); + let response = app + .clone() + .oneshot(upload(run, &digest, &worker, Body::from_stream(broken))) + .await + .unwrap(); + assert_status!(response, StatusCode::BAD_REQUEST).await; + let missing_run = RunId::new(); + let response = app + .clone() + .oneshot(upload( + missing_run, + &digest, + &issue_test_worker_token(&missing_run), + Body::from("valid content"), + )) + .await + .unwrap(); + assert_status!(response, StatusCode::NOT_FOUND).await; + run_records::append(&state, run, PlatformRecord::RunArchived) + .await + .unwrap(); + let response = app + .oneshot(upload(run, &digest, &worker, Body::from("valid content"))) + .await + .unwrap(); + assert_status!(response, StatusCode::CONFLICT).await; + + assert_eq!( + state + .artifact_store + .get_capture(&run, &hash) + .await + .unwrap() + .unwrap() + .as_ref(), + b"valid content" + ); + assert!( + state + .store_ref() + .blobs() + .read(&hash) + .await + .unwrap() + .is_none() + ); +} + +#[tokio::test] +async fn artifact_worker_client_uses_the_configured_s3_backend_and_prefix() { + let s3 = MockServer::start_async().await; + let settings = fabro_types::settings::server::ObjectStoreSettings::S3 { + bucket: "capture-bucket".to_string(), + region: "us-east-1".to_string(), + endpoint: Some(s3.base_url()), + path_style: true, + }; + let objects = crate::serve::build_object_store_from_settings_with_lookup( + &settings, + &|name| match name { + "AWS_ACCESS_KEY_ID" => Some("fake-access-key".to_string()), + "AWS_SECRET_ACCESS_KEY" => Some("fake-secret-key".to_string()), + _ => None, + }, + Some(&crate::serve::ObjectStoreBuildOptions { + client_options: object_store::ClientOptions::new().with_allow_http(true), + retry_config: object_store::RetryConfig { + max_retries: 0, + ..Default::default() + }, + }), + ) + .unwrap(); + let (database, _) = test_store_bundle(); + let state = TestAppStateBuilder::new() + .vault_entries([("OPENAI_API_KEY", "test-openai-api-key")]) + .server_secret_env(HashMap::from([( + "SESSION_SECRET".to_string(), + TEST_SESSION_SECRET.to_string(), + )])) + .store_bundle(database, ArtifactStore::new(objects, "selected-prefix")) + .build(); + let app = build_router(state.clone(), jwt_auth_mode()); + let run = create_run_with_bearer(&app, &issue_test_user_jwt()).await; + let bytes = vec![0x82; 3 * 1024 * 1024]; + let hash = BlobHash::new(&bytes); + let path = format!("/capture-bucket/selected-prefix/{run}/captures/sha256/{hash}"); + let expected = bytes.clone(); + let put = s3 + .mock_async(|when, then| { + when.method(httpmock::Method::PUT) + .path(&path) + .is_true(move |request| request.body_ref() == expected.as_slice()); + then.status(200).header("etag", "\"test-etag\""); + }) + .await; + let get = s3 + .mock_async(|when, then| { + when.method(GET).path(&path); + then.status(200) + .header("etag", "\"test-etag\"") + .header("last-modified", "Thu, 24 Sep 2026 12:00:00 GMT") + .body(bytes.clone()); + }) + .await; + let server = WorkerControlWsTestServer::spawn(app).await; + let worker_token = issue_test_worker_token(&run); + let mut headers = fabro_http::header::HeaderMap::new(); + headers.insert( + fabro_http::header::AUTHORIZATION, + format!("Bearer {worker_token}").parse().unwrap(), + ); + let transport = fabro_http::HttpClientBuilder::new() + .no_proxy() + .default_headers(headers) + .build() + .unwrap(); + + let client = fabro_client::Client::builder() + .transport(server.base_url.replacen("ws://", "http://", 1), transport) + .credential(fabro_client::Credential::Worker(worker_token)) + .connect() + .await + .unwrap(); + let writer = ClientArtifactWriter::new(client, run); + assert_eq!(writer.write(&bytes).await.unwrap(), hash); + assert_eq!( + state + .artifact_store + .get_capture(&run, &hash) + .await + .unwrap() + .unwrap(), + bytes + ); + assert!( + state + .store_ref() + .blobs() + .read(&hash) + .await + .unwrap() + .is_none() + ); + put.assert_async().await; + get.assert_async().await; +} diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index bb429a499..9fe1f0df4 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -55,6 +55,7 @@ tracing.workspace = true fabro-petri = { path = ".", features = ["test-support"] } fabro-tool = { path = "../fabro-tool" } httpmock = "0.8" +object_store.workspace = true pebble-coding-agent = { workspace = true, features = ["test-util"] } fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"] } fabro-llm = { path = "../fabro-llm", features = ["test-support"] } diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 1b55eca77..dcd020918 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -106,8 +106,9 @@ the stream places them with its finish), `checkpoint` after every route with the stage's diff from its parent commit (`diff_summary`, and the patch as a text blob under `patch_blob`), `artifact.collected` for every file under `[run.artifacts] include` a stage left in its workspace (the bytes go to -the blob table; a file unchanged since an earlier capture is not recorded -again), and `run.diff` at the run's end (the run branch's last checkpoint +the configured local/S3 artifact store through a run-bound writer; a file +unchanged since an earlier capture is not recorded again), and `run.diff` +at the run's end (the run branch's last checkpoint against the base, summary and patch blob). The projection folds them into `start`, `git_identity`, `checkpoints[].diff`, `StageProjection.diff`, `artifacts` and `Conclusion.diff`; a patch is carried as its @@ -189,8 +190,9 @@ Integration tests live under `tests/`: interoperation with Fabro's `BlobStore`. - `hooks.rs` runs command-only bundles through the engine assembly with Fabro's hooks over the memory store, in-memory platform records and an - in-memory blob table: every finish is committed and recorded, a failed - stage's route sees its files, a failed checkpoint ends the run, the run + in-memory blob table and separate local artifact store: every finish is + committed and recorded, a failed stage's route sees its files, a failed + checkpoint ends the run, the run branch, identity, artifacts, per-checkpoint diffs and the run diff are recorded, and the Docker and Daytona variants commit inside their sandboxes. diff --git a/lib/components/fabro-petri/VIEWS.md b/lib/components/fabro-petri/VIEWS.md index 788978640..b42327a80 100644 --- a/lib/components/fabro-petri/VIEWS.md +++ b/lib/components/fabro-petri/VIEWS.md @@ -148,7 +148,7 @@ stages live in the child invocation and list under the fork (see Parallel). | notes | `StageCompletion.notes` | `step.finished` `outcome` notes; `parsed.note {result_prepared, transition}` | attempt | | files touched | `stage.completed` `files_touched` | Pebble's fold of envelope `ToolCallCompleted` (see Agent activity) | session | | stage diff | `StageProjection.diff` | platform record `checkpoint {execution, firing, patch_blob}` | stage | -| artifacts | `RunArtifactEntry {stage_id, node_slug, retry, relative_path, size}`, the stage artifact endpoints | platform record `artifact.collected {execution, firing, attempt, path, blob, bytes, digest}` from the `transition` hook, the bytes in the blob table; a file unchanged since an earlier capture is not recorded again | attempt | +| artifacts | `RunArtifactEntry {stage_id, node_slug, retry, relative_path, size}`, the stage artifact endpoints | platform record `artifact.collected {execution, firing, attempt, path, object (or historical blob), bytes, digest}` from the `transition` hook, new bytes in configured artifact storage, historical blob sources in SQLite; a file unchanged since an earlier capture is not recorded again | attempt | | checkout | `setup.*` lines, `attractor.checkout` | the root `start` stage's `custom attractor.checkout {repository, commit, depth, files}` and its log lines; `[run.prepare]` commands are `run_prepare_N` stages | stage | | hook decisions | none today | `parsed.note {kind: hook}` (`HookReport`), `custom attractor.hook` (a point a step asks itself), `parsed.hook_activity`; `run.note.recorded` for run-level points | attempt | | budget pause | none today | `parsed.budget {state, attempt, remaining_ms, pending_questions}` | attempt | diff --git a/lib/components/fabro-petri/src/artifacts.rs b/lib/components/fabro-petri/src/artifacts.rs new file mode 100644 index 000000000..d7eca8376 --- /dev/null +++ b/lib/components/fabro-petri/src/artifacts.rs @@ -0,0 +1,73 @@ +//! Captured workspace files go to the server's configured artifact store. +//! Engine values and patches keep using the separate blob capability. + +use fabro_client::Client; +use fabro_store::ArtifactStore; +use fabro_types::{BlobHash, RunId}; + +#[derive(Debug, thiserror::Error)] +pub enum ArtifactWriteError { + #[error("artifact storage failed")] + 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. +#[async_trait::async_trait] +pub trait ArtifactWriter: Send + Sync { + async fn write(&self, bytes: &[u8]) -> Result; +} + +pub struct StoreArtifactWriter { + store: ArtifactStore, + run_id: RunId, +} + +impl StoreArtifactWriter { + #[must_use] + pub fn new(store: ArtifactStore, run_id: RunId) -> Self { + Self { store, run_id } + } +} + +#[async_trait::async_trait] +impl ArtifactWriter for StoreArtifactWriter { + async fn write(&self, bytes: &[u8]) -> Result { + self.store + .put_capture(&self.run_id, bytes) + .await + .map_err(ArtifactWriteError::from) + } +} + +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 } + } +} + +#[async_trait::async_trait] +impl ArtifactWriter for ClientArtifactWriter { + async fn write(&self, bytes: &[u8]) -> Result { + let hash = BlobHash::new(bytes); + self.client + .write_run_artifact_content(&self.run_id, &hash, bytes) + .await + .map_err(ArtifactWriteError::Upload)?; + Ok(hash) + } +} diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index b322cce1e..768af58cd 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -63,6 +63,7 @@ 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}; @@ -84,36 +85,38 @@ pub enum Execution { pub struct RunRequest { /// The Fabro run id, which becomes Petri's run key: the run's identity /// in the store and the label on every sandbox of the run. - pub run_id: String, + pub run_id: String, /// Where the run's workspaces, step output and blobs live. - pub run_dir: PathBuf, - pub execution: Execution, + pub run_dir: PathBuf, + pub execution: Execution, /// The run's durable record: the worker's HTTP store, or the server's /// SQLite store under the test override. - pub store: Arc, - pub runtime: RuntimeSpec, + pub store: Arc, + pub runtime: RuntimeSpec, /// The sandbox provider Fabro resolved for the run's environment. - pub provider: SandboxProviderKind, + pub provider: SandboxProviderKind, /// Fires to cancel the run. - pub cancel: CancellationToken, + pub cancel: CancellationToken, /// The run's pause, unpause and steer controls, which the caller keeps /// a clone of to drive them while the run is live. - pub controls: RunControls, + pub controls: RunControls, /// Where the run's questions go. - pub interviewer: Arc, + pub interviewer: Arc, /// The caller's observers of every record, registered ahead of the /// interview dispatcher: the interviewer's own expiry observer among /// them. - pub observers: Vec>, + pub observers: Vec>, /// Where `{{ secrets.NAME }}` references resolve from; `None` leaves /// every secret unknown. - pub secrets: Option>, + pub secrets: Option>, /// Where offloaded stage values go; `None` keeps Petri's local store /// under the run directory. - pub blobs: Option>, + pub blobs: Option>, + /// Where captured workspace files go. Required when capture writes files. + pub artifact_writer: Option>, /// Fabro's hooks: the checkpoint commit and its record. `None` runs /// with Petri's local hook service alone. - pub hooks: Option, + pub hooks: Option, } /// The recorded status of a finished run. @@ -210,6 +213,7 @@ pub async fn run(request: RunRequest) -> Result { Arc::clone(&request.store), resumed, request.blobs.clone(), + request.artifact_writer.clone(), )) }); if let Some(hooks) = &fabro_hooks { diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 94bb4be6e..29667dd05 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -25,10 +25,10 @@ //! and the checkpoint's operation identity, with the stage's diff from its //! parent commit (`diff_summary`, and the patch as a blob); then the stage's //! artifacts: every file under `[run.artifacts] include` in the stage's -//! workspace goes to the blob table and gets an `artifact.collected` record, -//! unless the same file with the same content was already collected earlier -//! in the run. A failed write is a recorded problem on the transition, never -//! a blocked route. +//! workspace goes to configured artifact storage and gets an +//! `artifact.collected` record, unless the same file with the same content +//! was already collected earlier in the run. A failed write is a recorded +//! problem on the transition, never a blocked route. //! - `run_finished`: the run's diff, its run branch against its base commit, as //! the `run.diff` platform record with the patch as a blob; then the //! forwarded point, so the local service runs `run_complete` and `run_failed` @@ -81,7 +81,7 @@ use fabro_store::platform_records::{ }; use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition}; use fabro_types::settings::run::RunNamespace; -use fabro_types::{BlobHash, DiffSummary, GitIdentity, RunId}; +use fabro_types::{ArtifactSource, BlobHash, DiffSummary, GitIdentity, RunId}; use fabro_util::error::collect_chain; use fabro_util::sync; use fabro_util::workspace_glob::{WorkspaceGlobError, WorkspaceGlobSet}; @@ -98,6 +98,7 @@ use tokio::sync::{Mutex as AsyncMutex, OnceCell}; use tokio::{fs, time}; use tracing::{debug, info, warn}; +use crate::artifacts::{ArtifactWriteError, ArtifactWriter}; use crate::blobs::Blobs; use crate::checkpoint::{ CHECKPOINT_FAILED_CLASS, CheckpointError, CheckpointKey, EXCLUDE_DIRS, RunGitSettings, @@ -120,7 +121,7 @@ const GATE_POLL: Duration = Duration::from_millis(50); /// budget. const ARTIFACT_MAX_FILES: usize = 100; /// The largest file collected, the legacy executor's budget. -const ARTIFACT_MAX_FILE_BYTES: u64 = 10 * 1024 * 1024; +const ARTIFACT_MAX_FILE_BYTES: u64 = fabro_types::ARTIFACT_MAX_FILE_BYTES as u64; /// The most bytes one stage's collection keeps, the legacy executor's /// budget. const ARTIFACT_MAX_TOTAL_BYTES: u64 = 50 * 1024 * 1024; @@ -184,8 +185,10 @@ pub enum HookError { }, #[error("invalid run.artifacts.include pattern")] Globs(#[source] Arc), - #[error("the run has no blob table to collect artifacts into")] - NoBlobs, + #[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}`")] List { root: String, @@ -367,9 +370,9 @@ pub struct FabroHooks { inner: Arc, run_id: RunId, records: Arc, - /// Where an artifact's bytes and a diff's patch go; `None` records - /// summaries alone. + /// Where diff patches go; `None` records summaries alone. blobs: Option>, + artifact_writer: Option>, workspaces: RunWorkspaces, lookup: WorkspaceLookup, identity: GitIdentity, @@ -396,8 +399,8 @@ 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` is where artifact bytes and - /// diff patches go. + /// its scope's first acquisition. `blobs` holds diff patches; + /// `artifact_writer` holds captured workspace files. #[must_use] pub fn new( spec: HooksSpec, @@ -408,6 +411,7 @@ impl FabroHooks { store: Arc, resumed: bool, blobs: Option>, + artifact_writer: Option>, ) -> Self { let identity = GitIdentity { name: spec.git.author.name.clone(), @@ -425,6 +429,7 @@ impl FabroHooks { run_id, records: spec.records, blobs, + artifact_writer, workspaces, lookup: WorkspaceLookup::new(Arc::clone(&store), run_key), identity, @@ -943,10 +948,6 @@ impl FabroHooks { // environment; there is no workspace to collect from. return Ok(0); }; - let Some(blobs) = &self.blobs else { - return Err(HookError::NoBlobs); - }; - let already = self.collected_artifacts().await?; 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; @@ -963,49 +964,68 @@ impl FabroHooks { continue; } }; - let digest = BlobHash::new(&bytes); - let identity = (path.clone(), digest.to_string()); - if sync::lock(already).contains(&identity) { + if !self.store_artifact(key, path, &bytes).await? { continue; } - let blob = blobs - .write(&bytes) - .await - .map_err(|source| HookError::Blob { - what: format!("artifact `{path}`"), - source, - })?; - let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord { - execution: key.execution, - firing: key.firing, - attempt: key.attempt, - path: path.clone(), - blob, - bytes: u64::try_from(bytes.len()).unwrap_or(u64::MAX), - digest: digest.to_string(), - operation: Some(key.operation_for(ARTIFACT_EFFECT)), - }); - self.records - .append( - &self.run_id, - &record, - Some(StagePosition { - execution: key.execution, - firing: key.firing, - }), - ) - .await - .map_err(|source| HookError::Write { - kind: "artifact", - source, - })?; - sync::lock(already).insert(identity); total_bytes = total_bytes.saturating_add(size); collected += 1; } Ok(collected) } + /// Publish bytes before their record; only a recorded capture enters + /// the ledger. Failed writes remain retryable, including after restart. + async fn store_artifact( + &self, + key: CheckpointKey, + path: String, + bytes: &[u8], + ) -> Result { + let already = self.collected_artifacts().await?; + let digest = BlobHash::new(bytes); + let identity = (path.clone(), 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, + })); + } + let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord { + execution: key.execution, + firing: key.firing, + attempt: key.attempt, + path: path.clone(), + source: ArtifactSource::ObjectStore(stored), + bytes: u64::try_from(bytes.len()).unwrap_or(u64::MAX), + digest: digest.to_string(), + operation: Some(key.operation_for(ARTIFACT_EFFECT)), + }); + self.records + .append( + &self.run_id, + &record, + Some(StagePosition { + execution: key.execution, + firing: key.firing, + }), + ) + .await + .map_err(|source| HookError::Write { + kind: "artifact", + source, + })?; + sync::lock(already).insert(identity); + Ok(true) + } + /// 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> { @@ -1351,7 +1371,241 @@ impl ExecutionHooks for FabroHooks { #[cfg(test)] mod tests { + use std::sync::atomic::{AtomicBool, Ordering}; + + use fabro_store::{ArtifactStore, StoredPlatformRecord}; + use object_store::memory::InMemory; + use super::*; + use crate::artifacts::StoreArtifactWriter; + use crate::test_support::MemoryPlatformRecords; + + struct NoHooks; + #[async_trait::async_trait] + impl ExecutionHooks for NoHooks {} + + struct FlakyRecords { + records: MemoryPlatformRecords, + fail: AtomicBool, + } + + #[async_trait::async_trait] + impl PlatformRecords for FlakyRecords { + async fn append( + &self, + run_id: &RunId, + record: &PlatformRecord, + position: Option, + ) -> Result { + if self.fail.swap(false, Ordering::SeqCst) { + return Err(PlatformRecordError::Store(fabro_store::Error::Io( + std::io::Error::other("test append unavailable"), + ))); + } + self.records.append(run_id, record, position).await + } + async fn read_kind( + &self, + run_id: &RunId, + kind: PlatformRecordKind, + ) -> Result, PlatformRecordError> { + self.records.read_kind(run_id, kind).await + } + } + + struct FlakyWriter { + writer: StoreArtifactWriter, + fail: AtomicBool, + } + + #[async_trait::async_trait] + impl ArtifactWriter for FlakyWriter { + async fn write(&self, bytes: &[u8]) -> Result { + 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 + } + } + + fn artifact_hooks( + root: &Path, + run_id: RunId, + records: Arc, + writer: Option>, + ) -> FabroHooks { + FabroHooks::new( + HooksSpec { + records, + git: RunGitSettings::default(), + artifacts: vec!["assets/**".to_string()], + test_gates: None, + }, + Arc::new(NoHooks), + run_id, + RunKey::new(run_id.to_string()), + root.to_path_buf(), + Arc::new(petri_store::MemoryRunStore::new()), + false, + None, + writer, + ) + } + + #[tokio::test] + async fn artifact_failures_preserve_sources_and_retry_without_false_ledger_entries() { + for fail_upload in [true, false] { + 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(!fail_upload), + }); + let writer = Arc::new(FlakyWriter { + writer: StoreArtifactWriter::new(store.clone(), run), + fail: AtomicBool::new(fail_upload), + }); + let hooks = artifact_hooks(root.path(), run, records.clone(), Some(writer.clone())); + let key = CheckpointKey { + execution: 0, + firing: 1, + 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.clone(), bytes) + .await + .unwrap_err(); + assert!(error.render().contains(if fail_upload { + "test store unavailable" + } else { + "test append unavailable" + })); + assert!(records.records.records(&run).is_empty()); + assert!(sync::lock(hooks.collected_artifacts().await.unwrap()).is_empty()); + assert_eq!( + store.get_capture(&run, &hash).await.unwrap().is_some(), + !fail_upload + ); + assert!(store.list_for_run(&run).await.unwrap().is_empty()); + assert!( + hooks + .store_artifact(key, path.clone(), bytes) + .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); + } + } + + #[tokio::test] + async fn artifact_legacy_resume_preserves_sources_and_new_run_isolation() { + 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 key = CheckpointKey { + execution: 0, + firing: 1, + attempt: 1, + }; + let bytes = b"old payload"; + let hash = BlobHash::new(bytes); + let path = "assets/report.bin".to_string(); + records + .append( + &run, + &PlatformRecord::ArtifactCollected(ArtifactCollectedRecord { + execution: key.execution, + firing: key.firing, + attempt: key.attempt, + path: path.clone(), + source: ArtifactSource::SqliteBlob(hash), + bytes: bytes.len() as u64, + digest: hash.to_string(), + operation: None, + }), + None, + ) + .await + .unwrap(); + let resumed = artifact_hooks( + root.path(), + run, + records.clone(), + Some(Arc::new(StoreArtifactWriter::new(store.clone(), run))), + ); + assert!( + !resumed + .store_artifact(key, path.clone(), bytes) + .await + .unwrap() + ); + assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); + assert!( + resumed + .store_artifact(key, path.clone(), b"changed") + .await + .unwrap() + ); + assert_eq!(records.records(&run).len(), 2); + + // A fork has its own run ID and no inherited capture records. The same + // payload must be captured under that run, without a cross-run reference. + let fork = RunId::new(); + let fork_hooks = artifact_hooks( + root.path(), + fork, + records.clone(), + Some(Arc::new(StoreArtifactWriter::new(store.clone(), fork))), + ); + assert!(fork_hooks.store_artifact(key, path, bytes).await.unwrap()); + assert_eq!( + store.get_capture(&fork, &hash).await.unwrap().unwrap(), + bytes.as_slice() + ); + assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); + 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() { diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index 80b73dce8..1e4fb7239 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -62,6 +62,7 @@ //! under `petri_*` keys. pub mod admission; +pub mod artifacts; pub mod blobs; pub mod check; pub mod checkpoint; diff --git a/lib/components/fabro-petri/src/projection/platform.rs b/lib/components/fabro-petri/src/projection/platform.rs index 7109e25fd..ca8c2c8a9 100644 --- a/lib/components/fabro-petri/src/projection/platform.rs +++ b/lib/components/fabro-petri/src/projection/platform.rs @@ -121,7 +121,7 @@ impl RunView { retry: record.attempt, relative_path: record.path.clone(), size: record.bytes, - blob: record.blob, + source: record.source, }); } PlatformRecord::RunDiff(record) => { diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 5849b4633..a2f0f7f11 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -22,6 +22,7 @@ use std::sync::Arc; use fabro_checkpoint::author::GitAuthor; use fabro_petri::admission::AdmittedGraphs; +use fabro_petri::artifacts::StoreArtifactWriter; use fabro_petri::blobs::Blobs; use fabro_petri::check::{self, Bundle, CheckRequest, Launch}; use fabro_petri::checkpoint::{ @@ -34,9 +35,10 @@ use fabro_petri::platform_records::PlatformRecords; use fabro_petri::recovery::{self, Recovery, RecoveryRequest}; use fabro_petri::runtime::RuntimeSpec; use fabro_petri::test_support::{MemoryBlobs, MemoryPlatformRecords}; -use fabro_store::{PlatformRecord, PlatformRecordKind}; +use fabro_store::{ArtifactStore, PlatformRecord, PlatformRecordKind}; use fabro_types::settings::run::RunCheckpointSettings; use fabro_types::{GitIdentitySource, RunId, SandboxProviderKind}; +use object_store::local::LocalFileSystem; use petri_execution::inspect::{self, RunInspection}; use petri_store::{Access, MemoryRunStore, RunKey, RunStore as _}; use tokio::fs; @@ -142,27 +144,38 @@ fn docker_plugin() -> Option { /// One run's pieces: the store, its platform records, where it ran. struct Harness { - run_id: RunId, - run_dir: PathBuf, - store: Arc, - records: Arc, - blobs: Arc, + run_id: RunId, + run_dir: PathBuf, + store: Arc, + records: Arc, + blobs: Arc, /// The `[run.artifacts] include` patterns the hooks collect under. - artifacts: Vec, - _root: tempfile::TempDir, + artifacts: Vec, + artifact_store: ArtifactStore, + _root: tempfile::TempDir, } impl Harness { fn new() -> Self { let root = tempfile::tempdir().expect("a temp dir"); + let artifact_root = root.path().join("artifacts"); + std::fs::create_dir(&artifact_root).expect("the isolated artifact directory creates"); + let artifact_store = ArtifactStore::new( + Arc::new( + LocalFileSystem::new_with_prefix(artifact_root) + .expect("the local artifact backend builds"), + ), + "captures-test", + ); Self { - run_id: RunId::new(), - run_dir: root.path().join("run"), - store: Arc::new(MemoryRunStore::new()), - records: Arc::new(MemoryPlatformRecords::new()), - blobs: Arc::new(MemoryBlobs::new()), + artifact_store, + run_id: RunId::new(), + run_dir: root.path().join("run"), + store: Arc::new(MemoryRunStore::new()), + records: Arc::new(MemoryPlatformRecords::new()), + blobs: Arc::new(MemoryBlobs::new()), artifacts: Vec::new(), - _root: root, + _root: root, } } @@ -207,6 +220,10 @@ impl Harness { observers, secrets: None, blobs: Some(Arc::clone(&self.blobs) as Arc), + artifact_writer: Some(Arc::new(StoreArtifactWriter::new( + self.artifact_store.clone(), + self.run_id, + ))), hooks: Some(hooks), }; engine::run(request).await.expect("the run executes") @@ -430,10 +447,10 @@ async fn every_finish_is_committed_and_recorded() { } /// The files under `[run.artifacts] include` are collected once per -/// content into the blob table, the run branch and the author identity are -/// recorded when the branch is created, every checkpoint after the first -/// carries its diff from its parent, and the run's diff is recorded at the -/// end. +/// content into configured artifact storage, the run branch and the author +/// identity are recorded when the branch is created, every checkpoint after the +/// first carries its diff from its parent, and the run's diff is recorded at +/// the end. #[tokio::test] async fn artifacts_the_branch_and_the_diffs_are_recorded() { if host_plugin().is_none() { @@ -467,14 +484,27 @@ async fn artifacts_the_branch_and_the_diffs_are_recorded() { "{artifacts:?}" ); assert_eq!(artifacts[0].bytes, 3); - assert_eq!(artifacts[0].digest, artifacts[0].blob.to_string()); + assert_eq!(artifacts[0].digest, artifacts[0].source.hash().to_string()); assert_ne!(artifacts[0].digest, artifacts[1].digest); + assert!( + artifacts + .iter() + .all(|artifact| matches!(artifact.source, fabro_types::ArtifactSource::ObjectStore(_))) + ); + assert!( + harness + .blobs + .read(&artifacts[1].source.hash()) + .await + .unwrap() + .is_none() + ); let bytes = harness - .blobs - .read(&artifacts[1].blob) + .artifact_store + .get_capture(&harness.run_id, &artifacts[1].source.hash()) .await - .expect("the blob reads") - .expect("the blob exists"); + .expect("the object reads") + .expect("the object exists"); assert_eq!(bytes.as_ref(), b"two"); // The first capture belongs to `write`, the second to `change`; `keep` // saw the file unchanged and recorded nothing. @@ -821,6 +851,7 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() { observers, secrets: None, blobs: None, + artifact_writer: None, hooks: Some(harness.hooks(&SandboxProviderKind::LOCAL)), }; let outcome = engine::run(request).await.expect("the run executes"); diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs index 7cc8fd0c3..9fd47f914 100644 --- a/lib/components/fabro-petri/tests/support/mod.rs +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -121,6 +121,7 @@ pub(crate) fn run_request( interviewer: Arc::new(interviewer), secrets: None, blobs: None, + artifact_writer: None, hooks: None, } } diff --git a/lib/components/fabro-store/src/artifact_store.rs b/lib/components/fabro-store/src/artifact_store.rs index 072e0434c..2193d2a34 100644 --- a/lib/components/fabro-store/src/artifact_store.rs +++ b/lib/components/fabro-store/src/artifact_store.rs @@ -2,7 +2,7 @@ use std::sync::Arc; use bytes::Bytes; use chrono::Utc; -use fabro_types::RunId; +use fabro_types::{BlobHash, RunId}; use futures::StreamExt; use futures::stream::BoxStream; use object_store::ObjectStore; @@ -81,6 +81,33 @@ impl ArtifactStore { Ok(()) } + /// Publish a complete captured file under its content digest within the + /// run. Repeating a put of the same content is safe; metadata is recorded + /// separately only after this operation succeeds. + pub async fn put_capture(&self, run_id: &RunId, data: &[u8]) -> Result { + let hash = BlobHash::new(data); + let path = self.capture_prefix(run_id)?.child(hash.to_string()); + self.object_store + .put(&path, Bytes::copy_from_slice(data).into()) + .await?; + Ok(hash) + } + + /// Read run-owned content named by a capture record, without consulting + /// historical stage keys or the SQLite blob table. + pub async fn get_capture(&self, run_id: &RunId, hash: &BlobHash) -> Result> { + let path = self.capture_prefix(run_id)?.child(hash.to_string()); + match self.object_store.get(&path).await { + Ok(result) => Ok(Some(result.bytes().await?)), + Err(object_store::Error::NotFound { .. }) => Ok(None), + Err(error) => Err(error.into()), + } + } + + fn capture_prefix(&self, run_id: &RunId) -> Result { + Ok(self.run_prefix(run_id)?.child("captures").child("sha256")) + } + pub fn writer(&self, run_id: &RunId, key: &ArtifactKey) -> Result { let path = self.artifact_path(run_id, key)?; Ok(BufWriter::with_capacity( @@ -152,9 +179,15 @@ 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 mut stream = self.object_store.list(Some(&prefix)); let mut artifacts = Vec::new(); while let Some(meta) = stream.next().await.transpose()? { + // Capture metadata lives in platform records. Content objects + // must not be decoded as historical stage/retry/path entries. + if meta.location.prefix_match(&captures).is_some() { + continue; + } artifacts.push(decode_artifact_location( &prefix, &meta.location, @@ -427,6 +460,66 @@ mod tests { ArtifactStore::new(object_store, "artifacts") } + #[tokio::test] + async fn captures_are_run_owned_hidden_from_legacy_listing_and_deleted_with_the_run() { + let objects: Arc = Arc::new(InMemory::new()); + let store = ArtifactStore::new(objects.clone(), "artifacts"); + let run = RunId::new(); + let other = RunId::new(); + let bytes = b"binary\0capture"; + let hash = fabro_types::BlobHash::new(bytes); + let legacy = ArtifactKey::new(StageId::new("captures", 1), 1, "sha256/report.bin"); + 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); + } + let location = ObjectPath::from(format!("artifacts/{run}/captures/sha256/{hash}")); + assert_eq!( + objects.get(&location).await.unwrap().bytes().await.unwrap(), + bytes.as_slice() + ); + assert_eq!( + store.get_capture(&run, &hash).await.unwrap().unwrap(), + bytes.as_slice() + ); + assert_eq!(store.list_for_run(&run).await.unwrap().len(), 1); + assert_eq!( + store + .list_for_node(&run, &legacy.stage_id) + .await + .unwrap() + .len(), + 1 + ); + store.delete_for_run(&run).await.unwrap(); + assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); + assert!(store.get(&run, &legacy).await.unwrap().is_none()); + assert_eq!( + store.get_capture(&other, &hash).await.unwrap().unwrap(), + bytes.as_slice() + ); + assert!( + objects + .head(&ObjectPath::from("artifacts/store-metadata.json")) + .await + .is_ok() + ); + } + + #[tokio::test] + async fn capture_namespace_does_not_hide_malformed_legacy_objects() { + 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()) + .await + .unwrap(); + assert!(store.list_for_run(&run).await.is_err()); + } + #[tokio::test] async fn write_metadata_persists_store_marker() { let object_store: Arc = Arc::new(InMemory::new()); diff --git a/lib/components/fabro-store/src/platform_records.rs b/lib/components/fabro-store/src/platform_records.rs index d83ef9cbb..11e1742ba 100644 --- a/lib/components/fabro-store/src/platform_records.rs +++ b/lib/components/fabro-store/src/platform_records.rs @@ -27,8 +27,9 @@ use std::sync::Arc; use fabro_types::{ - BlobHash, DiffSummary, GitIdentity, PairId, PairTarget, Principal, PullRequestCreationId, - PullRequestLink, RunControlAction, RunId, RunNoticeLevel, RunSpec, RunStatus, + ArtifactSource, BlobHash, DiffSummary, GitIdentity, PairId, PairTarget, Principal, + PullRequestCreationId, PullRequestLink, RunControlAction, RunId, RunNoticeLevel, RunSpec, + RunStatus, }; use serde::{Deserialize, Serialize}; use sqlx::sqlite::{SqliteConnection, SqliteRow}; @@ -190,7 +191,7 @@ pub enum PlatformRecord { #[serde(rename = "checkpoint")] Checkpoint(CheckpointRecord), /// A file a stage's attempt left in its workspace, collected under - /// `[run.artifacts] include` into the blob table. + /// `[run.artifacts] include` into configured artifact storage. #[serde(rename = "artifact.collected")] ArtifactCollected(ArtifactCollectedRecord), /// The run's whole diff, its run branch against its base commit, written @@ -449,6 +450,7 @@ 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, @@ -456,8 +458,9 @@ pub struct ArtifactCollectedRecord { pub attempt: u32, /// The file's path relative to the workspace root. pub path: String, - /// The blob that holds the file's bytes. - pub blob: BlobHash, + /// Where the file's bytes are stored. + #[serde(flatten)] + 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. @@ -466,6 +469,41 @@ pub struct ArtifactCollectedRecord { 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, +} + +impl TryFrom for ArtifactCollectedRecord { + type Error = &'static str; + + 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"); + } + 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, + }) + } +} + /// The run's diff: its run branch's head against its base commit. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct RunDiffRecord { @@ -744,6 +782,27 @@ mod tests { ])) } + #[test] + fn artifact_records_preserve_old_json_and_validate_sources_and_checksums() { + let hash = BlobHash::new(b"report"); + for source in ["blob", "object"] { + let mut wire = json!({ + "kind":"artifact.collected", "execution":0, "firing":3, + "attempt":1, "path":"assets/report.txt", "bytes":6, + "digest":hash.to_string() + }); + wire[source] = json!(hash); + let record: PlatformRecord = serde_json::from_value(wire.clone()).unwrap(); + assert_eq!(serde_json::to_value(&record).unwrap(), wire); + wire["digest"] = json!(BlobHash::new(b"different").to_string()); + assert!(serde_json::from_value::(wire.clone()).is_err()); + wire["digest"] = json!(hash.to_string()); + wire["blob"] = json!(hash); + wire["object"] = json!(hash); + assert!(serde_json::from_value::(wire).is_err()); + } + } + fn sample(kind: PlatformRecordKind) -> PlatformRecord { match kind { PlatformRecordKind::RunCreated => PlatformRecord::RunCreated(RunCreatedRecord { @@ -826,7 +885,7 @@ mod tests { firing: 3, attempt: 1, path: "assets/report.txt".to_string(), - blob: BlobHash::new(b"report"), + source: ArtifactSource::SqliteBlob(BlobHash::new(b"report")), bytes: 6, digest: BlobHash::new(b"report").to_string(), operation: Some(OperationKey { diff --git a/lib/foundation/fabro-api/build.rs b/lib/foundation/fabro-api/build.rs index 7823c1762..551968fc1 100644 --- a/lib/foundation/fabro-api/build.rs +++ b/lib/foundation/fabro-api/build.rs @@ -623,6 +623,8 @@ fn main() { ("ModelCosts", "fabro_types::ModelCosts", &[]), ("ModelTestMode", "fabro_types::ModelTestMode", &[]), ("RunProjection", "fabro_types::RunProjection", &[]), + ("ArtifactSource", "fabro_types::ArtifactSource", &[]), + ("RunArtifact", "fabro_types::RunArtifact", &[]), ("PairId", "fabro_types::PairId", &[]), ("PairMessageId", "fabro_types::PairMessageId", &[]), ("PairStatus", "fabro_types::PairStatus", &[]), diff --git a/lib/foundation/fabro-api/src/lib.rs b/lib/foundation/fabro-api/src/lib.rs index 6961088e7..8f4707739 100644 --- a/lib/foundation/fabro-api/src/lib.rs +++ b/lib/foundation/fabro-api/src/lib.rs @@ -36,8 +36,8 @@ pub mod types { BlockedReason, FailureReason, PendingReason, RunControlAction, RunStatus, SuccessReason, }; pub use fabro_types::{ - AgentEventProps, AgentSessionActivatedProps, AgentToolsAvailableProps, AskFabro, - AuthMethod, AutomationRef, BlobHash, CommandTermination, Conclusion, + AgentEventProps, AgentSessionActivatedProps, AgentToolsAvailableProps, ArtifactSource, + AskFabro, AuthMethod, AutomationRef, BlobHash, CommandTermination, Conclusion, ContextWindowBreakdownItem, ContextWindowCategory, ContextWindowCountMethod, ContextWindowSnapshot, ContextWindowStaleness, ContextWindowWarning, CreateVariableRequest, DiffStats, DiffSummary, DirtyStatus, ExecOutputTail, FailureCategory, FailureDetail, @@ -54,21 +54,22 @@ pub mod types { PullRequestCreation, PullRequestCreationId, PullRequestCreationStatus, PullRequestDetails, PullRequestDetailsStatus, PullRequestDetailsUnavailableReason, PullRequestLink, PullRequestMeta, PullRequestResponse, QuestionType, RepositoryRef, ReviewTarget, - ReviewTargetKind, Run, RunApproval, RunApprovalState, RunClientProvenance, RunFailure, - RunGraph, RunGraphEdge, RunGraphNode, RunIntent, RunIntentArgs, RunPairStatusResponse, - RunProjection, RunProvenance, RunRunnableSource, RunSandbox, RunSandboxFailure, - RunSandboxInstance, RunSandboxKind, RunSandboxPlan, RunSandboxRuntime, RunServerProvenance, - RunSessionMetadata, RunSize, RunStreamItem, RunStreamItemKind, RunTarget, SandboxDetails, - SandboxInfo, SandboxListMeta, SandboxListResponse, SandboxProviderKind, - SandboxProviderLookupError, SandboxService, SandboxServiceListResponse, SecretMetadata, - SecretType, ServerSettings, SessionDetail, SessionEvent, SessionEventBody, SessionId, - SessionStatus, SessionSummary, SessionTurn, SkillActivationSource, SkillSummary, - StageCompletion, StageContextWindow, StageContextWindowUnavailableReason, StageHandler, - StageId, StageInferenceProjection, StageModelUsage, StageOutcome, StageProjection, - StageState, StageToolBatchProjection, SystemActorKind, SystemIntegrationStatus, - SystemIntegrationsResponse, TodoListProjection, ToolCategory, ToolSource, ToolSummary, - TurnId, UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowPath, - WorkflowSettings, WorkflowVersion, WorkflowVersionId, + ReviewTargetKind, Run, RunApproval, RunApprovalState, RunArtifact, RunClientProvenance, + RunFailure, RunGraph, RunGraphEdge, RunGraphNode, RunIntent, RunIntentArgs, + RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource, RunSandbox, + RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan, RunSandboxRuntime, + RunServerProvenance, RunSessionMetadata, RunSize, RunStreamItem, RunStreamItemKind, + RunTarget, SandboxDetails, SandboxInfo, SandboxListMeta, SandboxListResponse, + SandboxProviderKind, SandboxProviderLookupError, SandboxService, + SandboxServiceListResponse, SecretMetadata, SecretType, ServerSettings, SessionDetail, + SessionEvent, SessionEventBody, SessionId, SessionStatus, SessionSummary, SessionTurn, + SkillActivationSource, SkillSummary, StageCompletion, StageContextWindow, + StageContextWindowUnavailableReason, StageHandler, StageId, StageInferenceProjection, + StageModelUsage, StageOutcome, StageProjection, StageState, StageToolBatchProjection, + SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, + ToolCategory, ToolSource, ToolSummary, TurnId, UpdateVariableRequest, UserPrincipal, + Variable, VariableListResponse, WorkflowPath, WorkflowSettings, WorkflowVersion, + WorkflowVersionId, }; pub use lithos_llm::catalog::{ModelHandle, ProviderId}; pub use lithos_llm::types::{ diff --git a/lib/foundation/fabro-api/tests/artifact_round_trip.rs b/lib/foundation/fabro-api/tests/artifact_round_trip.rs new file mode 100644 index 000000000..53e519c0d --- /dev/null +++ b/lib/foundation/fabro-api/tests/artifact_round_trip.rs @@ -0,0 +1,67 @@ +use std::any::TypeId; + +use fabro_api::types; +use fabro_types::{ArtifactSource, BlobHash, RunArtifact, StageId}; +use serde_json::{Value, json}; + +fn validator(name: &str) -> jsonschema::Validator { + let yaml: serde_yaml::Value = serde_yaml::from_str(include_str!( + "../../../../docs/public/api-reference/fabro-api.yaml" + )) + .unwrap(); + let mut spec = serde_json::to_value(yaml).unwrap(); + spec["$ref"] = json!(format!("#/components/schemas/{name}")); + jsonschema::validator_for(&spec).unwrap() +} + +#[test] +fn artifact_api_types_reuse_the_canonical_types_and_wire_format() { + assert_eq!( + TypeId::of::(), + TypeId::of::() + ); + assert_eq!( + TypeId::of::(), + TypeId::of::() + ); + let hash = BlobHash::new(b"payload"); + for source in [ + ArtifactSource::SqliteBlob(hash), + ArtifactSource::ObjectStore(hash), + ] { + let source_json = serde_json::to_value(source).unwrap(); + assert!(validator("ArtifactSource").is_valid(&source_json)); + assert_eq!( + serde_json::from_value::(source_json).unwrap(), + source + ); + let artifact = RunArtifact { + stage_id: StageId::new("write", 1), + retry: 1, + relative_path: "assets/report.bin".to_string(), + size: 7, + source, + }; + let json = serde_json::to_value(&artifact).unwrap(); + assert!(validator("RunArtifact").is_valid(&json)); + assert_eq!( + serde_json::from_value::(json).unwrap(), + artifact + ); + } +} + +#[test] +fn artifact_schema_and_rust_reject_the_same_invalid_sources() { + let schema = validator("ArtifactSource"); + let hash = BlobHash::new(b"payload"); + for json in [ + json!({}), + json!({"blob": hash, "object": hash}), + json!({"blob": hash, "object": Value::Null}), + json!({"object":"invalid"}), + ] { + assert!(!schema.is_valid(&json), "{json}"); + assert!(serde_json::from_value::(json).is_err()); + } +} diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index 6f661d017..c52f0dd39 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1842,6 +1842,26 @@ impl Client { Ok(()) } + /// Upload complete captured file content using this run's worker token. + pub async fn write_run_artifact_content( + &self, + run_id: &RunId, + digest: &BlobHash, + data: &[u8], + ) -> Result<()> { + self.send_api(|client| async move { + client + .write_run_artifact_content() + .id(run_id.to_string()) + .digest(*digest) + .body(data.to_vec()) + .send() + .await + }) + .await?; + Ok(()) + } + pub async fn write_run_blob(&self, run_id: &RunId, data: &[u8]) -> Result { let response = self .send_api(|client| async move { diff --git a/lib/foundation/fabro-types/src/artifact_source.rs b/lib/foundation/fabro-types/src/artifact_source.rs new file mode 100644 index 000000000..0397408b9 --- /dev/null +++ b/lib/foundation/fabro-types/src/artifact_source.rs @@ -0,0 +1,55 @@ +use serde::de::Error as _; +use serde::{Deserialize, Deserializer, Serialize}; + +use crate::BlobHash; + +/// The durable payload location of a captured workspace file. +/// +/// Object content belongs to the containing run in its configured artifact +/// store. SQLite sources retain the wire shape of earlier captures. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)] +pub enum ArtifactSource { + #[serde(rename = "blob")] + SqliteBlob(BlobHash), + #[serde(rename = "object")] + ObjectStore(BlobHash), +} + +impl ArtifactSource { + #[must_use] + pub fn hash(self) -> BlobHash { + match self { + Self::SqliteBlob(hash) | Self::ObjectStore(hash) => hash, + } + } +} + +impl<'de> Deserialize<'de> for ArtifactSource { + fn deserialize>(deserializer: D) -> Result { + // A flattened externally tagged enum can accept one source and + // silently ignore a conflicting second source. Consume both known + // keys in a private wire shape before choosing the canonical enum. + #[derive(Deserialize)] + struct Fields { + #[serde(default, deserialize_with = "present_hash")] + blob: Option, + #[serde(default, deserialize_with = "present_hash")] + object: Option, + } + let fields = Fields::deserialize(deserializer)?; + match (fields.blob, fields.object) { + (Some(hash), None) => Ok(Self::SqliteBlob(hash)), + (None, Some(hash)) => Ok(Self::ObjectStore(hash)), + _ => Err(D::Error::custom( + "exactly one artifact source, blob or object, is required", + )), + } + } +} + +fn present_hash<'de, D: Deserializer<'de>>(deserializer: D) -> Result, D::Error> { + BlobHash::deserialize(deserializer).map(Some) +} + +/// Maximum bytes in one automatically captured workspace file. +pub const ARTIFACT_MAX_FILE_BYTES: usize = 10 * 1024 * 1024; diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index c10532063..e08f61bfa 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -1,6 +1,7 @@ extern crate self as fabro_types; pub mod agent_props; +mod artifact_source; pub mod auth; pub mod blob_hash; pub mod blob_ref; @@ -69,6 +70,7 @@ pub use agent_props::{ AgentEventProps, AgentSessionActivatedProps, AgentToolsAvailableProps, CODING_EVENT_NAMES, SessionCapability, StagePromptProps, coding_event_name, is_coding_event_name, }; +pub use artifact_source::{ARTIFACT_MAX_FILE_BYTES, ArtifactSource}; pub use auth::{IdpIdentity, IdpIdentityError}; pub use blob_hash::BlobHash; pub use blob_ref::{ diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 89a884ff5..62a007122 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -13,7 +13,7 @@ use strum::{Display, EnumString, IntoStaticStr}; use crate::agent_props::{AgentSessionActivatedProps, StagePromptProps}; use crate::{ - AgentBackend, BlobHash, Checkpoint, Conclusion, GitIdentity, InterviewQuestionRecord, + AgentBackend, ArtifactSource, Checkpoint, Conclusion, GitIdentity, InterviewQuestionRecord, InvalidTransition, ModelRef, ModelUsage, ParallelBranchId, PullRequestCreation, PullRequestLink, RunApproval, RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus, RunTiming, StageCompletion, StageHandler, StageId, StageState, StageTiming, StartRecord, @@ -75,8 +75,9 @@ pub struct RunArtifact { /// The file's path relative to the workspace root. pub relative_path: String, pub size: u64, - /// The blob that holds the file's bytes. - pub blob: BlobHash, + /// Where the file's bytes are stored. + #[serde(flatten)] + pub source: ArtifactSource, } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -922,6 +923,50 @@ impl RunProjection { } } +#[cfg(test)] +mod artifact_tests { + use super::RunArtifact; + use crate::BlobHash; + + fn artifact() -> serde_json::Value { + serde_json::json!({ + "stage_id": "write@1", "retry": 1, + "relative_path": "assets/report.bin", "size": 7 + }) + } + + #[test] + fn artifact_sources_round_trip_old_and_new_projection_json() { + for source in ["blob", "object"] { + let mut json = artifact(); + json[source] = serde_json::json!(BlobHash::new(b"payload")); + let decoded: RunArtifact = serde_json::from_value(json.clone()).unwrap(); + assert_eq!(serde_json::to_value(decoded).unwrap(), json); + } + } + + #[test] + fn artifact_sources_reject_ambiguous_missing_and_invalid_hashes() { + let hash = serde_json::json!(BlobHash::new(b"payload")); + for fields in [ + serde_json::json!({}), + serde_json::json!({"blob": hash, "object": hash}), + serde_json::json!({"blob": hash, "object": null}), + serde_json::json!({"object": "not-a-hash"}), + serde_json::json!({"blob": "not-a-hash"}), + ] { + let mut json = artifact(); + json.as_object_mut() + .unwrap() + .extend(fields.as_object().unwrap().clone()); + assert!( + serde_json::from_value::(json.clone()).is_err(), + "{json}" + ); + } + } +} + #[cfg(test)] mod title_tests { use chrono::Utc; diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index ff556cb1a..c97704a3c 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -57,6 +57,7 @@ models/api-question.ts models/approval-mode.ts models/artifact-entry.ts models/artifact-list-response.ts +models/artifact-source.ts models/artifacts-settings.ts models/ask-fabro.ts models/auth-config-response.ts @@ -375,6 +376,7 @@ models/run-approval-state.ts models/run-approval.ts models/run-artifact-entry.ts models/run-artifact-list-response.ts +models/run-artifact.ts models/run-branch-settings.ts models/run-checkpoint-settings.ts models/run-checkpoint.ts diff --git a/lib/packages/fabro-api-client/src/api/run-internals-api.ts b/lib/packages/fabro-api-client/src/api/run-internals-api.ts index 6ba34e729..b9e415cb5 100644 --- a/lib/packages/fabro-api-client/src/api/run-internals-api.ts +++ b/lib/packages/fabro-api-client/src/api/run-internals-api.ts @@ -1019,6 +1019,55 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, + /** + * Stores captured file bytes in the configured artifact backend for this run. Requires a worker token belonging to the run. The body is limited to 10 MiB and must match the SHA-256 digest. Repeating the same upload is safe. Uploading content alone does not create an artifact listing entry. + * @summary Write Captured Artifact Content + * @param {string} id Unique run identifier (ULID). + * @param {string} digest + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + writeRunArtifactContent: async (id: string, digest: string, body: File, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('writeRunArtifactContent', 'id', id) + // verify required parameter 'digest' is not null or undefined + assertParamExists('writeRunArtifactContent', 'digest', digest) + // verify required parameter 'body' is not null or undefined + assertParamExists('writeRunArtifactContent', 'body', body) + const localVarPath = `/api/v1/runs/{id}/artifacts/content/{digest}` + .replace(`{${"id"}}`, encodeURIComponent(String(id))) + .replace(`{${"digest"}}`, encodeURIComponent(String(digest))); + // use dummy base URL string because the URL constructor only accepts absolute URLs. + const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL); + let baseOptions; + if (configuration) { + baseOptions = configuration.baseOptions; + } + + const localVarRequestOptions = { method: 'PUT', ...baseOptions, ...options}; + const localVarHeaderParameter = {} as any; + const localVarQueryParameter = {} as any; + + // authentication SessionCookie required + + // authentication BearerAuth required + // http bearer authentication required + await setBearerAuthToObject(localVarHeaderParameter, configuration) + + localVarHeaderParameter['Content-Type'] = 'application/octet-stream'; + localVarHeaderParameter['Accept'] = 'application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + localVarRequestOptions.data = serializeDataIfNeeded(body, localVarRequestOptions, configuration) + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Writes an opaque binary blob and returns its content-addressed blob hash. * @summary Write Run Blob @@ -1370,6 +1419,21 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.writePetriBlob']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, + /** + * Stores captured file bytes in the configured artifact backend for this run. Requires a worker token belonging to the run. The body is limited to 10 MiB and must match the SHA-256 digest. Repeating the same upload is safe. Uploading content alone does not create an artifact listing entry. + * @summary Write Captured Artifact Content + * @param {string} id Unique run identifier (ULID). + * @param {string} digest + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async writeRunArtifactContent(id: string, digest: string, body: File, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.writeRunArtifactContent(id, digest, body, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.writeRunArtifactContent']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Writes an opaque binary blob and returns its content-addressed blob hash. * @summary Write Run Blob @@ -1627,6 +1691,18 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b writePetriBlob(id: string, owner: string, body: File, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.writePetriBlob(id, owner, body, options).then((request) => request(axios, basePath)); }, + /** + * Stores captured file bytes in the configured artifact backend for this run. Requires a worker token belonging to the run. The body is limited to 10 MiB and must match the SHA-256 digest. Repeating the same upload is safe. Uploading content alone does not create an artifact listing entry. + * @summary Write Captured Artifact Content + * @param {string} id Unique run identifier (ULID). + * @param {string} digest + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + writeRunArtifactContent(id: string, digest: string, body: File, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.writeRunArtifactContent(id, digest, body, options).then((request) => request(axios, basePath)); + }, /** * Writes an opaque binary blob and returns its content-addressed blob hash. * @summary Write Run Blob @@ -1900,6 +1976,19 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).writePetriBlob(id, owner, body, options).then((request) => request(this.axios, this.basePath)); } + /** + * Stores captured file bytes in the configured artifact backend for this run. Requires a worker token belonging to the run. The body is limited to 10 MiB and must match the SHA-256 digest. Repeating the same upload is safe. Uploading content alone does not create an artifact listing entry. + * @summary Write Captured Artifact Content + * @param {string} id Unique run identifier (ULID). + * @param {string} digest + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public writeRunArtifactContent(id: string, digest: string, body: File, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).writeRunArtifactContent(id, digest, body, options).then((request) => request(this.axios, this.basePath)); + } + /** * Writes an opaque binary blob and returns its content-addressed blob hash. * @summary Write Run Blob diff --git a/lib/packages/fabro-api-client/src/models/artifact-source.ts b/lib/packages/fabro-api-client/src/models/artifact-source.ts new file mode 100644 index 000000000..425afdc77 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/artifact-source.ts @@ -0,0 +1,29 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * Exactly one payload source; object content is owned by the containing run. + */ +export interface ArtifactSource { + /** + * Content-addressed SHA-256 hash of a stored blob. Hex input is case-insensitive; Fabro emits the canonical lowercase form. + */ + 'blob'?: string; + /** + * Content-addressed SHA-256 hash of a stored blob. Hex input is case-insensitive; Fabro emits the canonical lowercase form. + */ + 'object'?: string; +} diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 6d9ddbb77..f57e1de38 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -28,6 +28,7 @@ export * from './api-question'; export * from './approval-mode'; export * from './artifact-entry'; export * from './artifact-list-response'; +export * from './artifact-source'; export * from './artifacts-settings'; export * from './ask-fabro'; export * from './auth-config-response'; @@ -344,6 +345,7 @@ export * from './run'; export * from './run-agent-settings'; export * from './run-approval'; export * from './run-approval-state'; +export * from './run-artifact'; export * from './run-artifact-entry'; export * from './run-artifact-list-response'; export * from './run-branch-settings'; diff --git a/lib/packages/fabro-api-client/src/models/run-artifact.ts b/lib/packages/fabro-api-client/src/models/run-artifact.ts new file mode 100644 index 000000000..48982046e --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/run-artifact.ts @@ -0,0 +1,36 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * A captured workspace file with its durable payload source. + */ +export interface RunArtifact { + /** + * Content-addressed SHA-256 hash of a stored blob. Hex input is case-insensitive; Fabro emits the canonical lowercase form. + */ + 'blob'?: string; + /** + * Content-addressed SHA-256 hash of a stored blob. Hex input is case-insensitive; Fabro emits the canonical lowercase form. + */ + 'object'?: string; + /** + * Canonical stage execution identifier in `node_id@visit` form. + */ + 'stage_id': string; + 'retry': number; + 'relative_path': string; + 'size': number; +} diff --git a/lib/packages/fabro-api-client/src/models/run-projection.ts b/lib/packages/fabro-api-client/src/models/run-projection.ts index d0209f38c..ef16ca2de 100644 --- a/lib/packages/fabro-api-client/src/models/run-projection.ts +++ b/lib/packages/fabro-api-client/src/models/run-projection.ts @@ -36,6 +36,9 @@ import type { PullRequestCreation } from './pull-request-creation'; import type { PullRequestLink } from './pull-request-link'; // May contain unused imports in some cases // @ts-ignore +import type { RunArtifact } from './run-artifact'; +// May contain unused imports in some cases +// @ts-ignore import type { RunControlAction } from './run-control-action'; // May contain unused imports in some cases // @ts-ignore @@ -76,6 +79,10 @@ export interface RunProjection { 'status_updated_at': string; 'last_event_at': string; 'pending_control'?: RunControlAction | null; + /** + * Captured files; older projections may omit this field. + */ + 'artifacts'?: Array; /** * Sequence-tagged checkpoint history entries. */ From 71b08b61a190f8b54d1bdcd5a7c78a8f0a4c6729 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Fri, 25 Sep 2026 11:56:59 -0400 Subject: [PATCH 2/3] Preserve custom local artifact roots when overriding storage --- .../administration/server-configuration.mdx | 4 ++ .../fabro-cli/tests/it/scenario/artifacts.rs | 23 +++++--- lib/apps/fabro-server/src/install.rs | 11 ++-- lib/apps/fabro-server/src/serve.rs | 52 ++++++++++++++++++- lib/foundation/fabro-types/src/dense.rs | 24 +++------ 5 files changed, 85 insertions(+), 29 deletions(-) diff --git a/docs/public/administration/server-configuration.mdx b/docs/public/administration/server-configuration.mdx index 48315c0e6..d6850d3d2 100644 --- a/docs/public/administration/server-configuration.mdx +++ b/docs/public/administration/server-configuration.mdx @@ -247,6 +247,10 @@ prefix = "artifacts" root = "/var/lib/fabro/objects" ``` +A custom local artifact root is preserved at startup and when `--storage-dir` +changes the server's storage directory. When `local.root` is omitted, the +default `/objects/artifacts` location follows that override. + The wizard only covers AWS S3 bucket/region plus one of: - runtime credentials already supplied by the deployment environment diff --git a/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs b/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs index 81f327746..952a99621 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/artifacts.rs @@ -16,14 +16,14 @@ async fn artifact_worker_captures_large_files_in_the_configured_local_store() { return; } let context = test_context!(); - let server = RunningServer::start_with( - "\n[server.artifacts]\nprovider = \"local\"\nprefix = \"selected-prefix\"\n", - &[], - ) - .await; - // RunningServer explicitly selects --storage-dir, which also selects the - // local artifact root. Inspect that resolved backend, outside the sandbox. - let objects = server.storage_dir.join("objects/artifacts"); + let objects = context.temp_dir.join("selected-artifact-root"); + let settings = format!( + "\n[server.artifacts]\nprovider = \"local\"\nprefix = \"selected-prefix\"\n[server.artifacts.local]\nroot = {:?}\n", + objects.to_str().unwrap() + ); + // RunningServer also passes --storage-dir. The explicit artifact directory + // must still receive the captures, independently of the database root. + let server = RunningServer::start_with(&settings, &[]).await; let workspace = artifact_workspace(&context); tokio::fs::write(workspace.join("workflow.fabro"), r#"digraph Capture { graph [goal="Capture binary files", default_max_retries=0] @@ -58,6 +58,13 @@ async fn artifact_worker_captures_large_files_in_the_configured_local_store() { serde_json::to_value(rebuilt.unwrap().artifacts).unwrap(), projection["artifacts"] ); + assert!( + !server + .storage_dir + .join(format!("objects/artifacts/selected-prefix/{run_id}")) + .exists(), + "captures must not be redirected to the default artifact directory" + ); for (path, size) in [ ("medium.bin", 3 * 1024 * 1024), diff --git a/lib/apps/fabro-server/src/install.rs b/lib/apps/fabro-server/src/install.rs index 3d3ee4547..49a3416a6 100644 --- a/lib/apps/fabro-server/src/install.rs +++ b/lib/apps/fabro-server/src/install.rs @@ -2262,8 +2262,7 @@ async fn write_artifact_store_metadata( settings: &ServerSettings, storage_dir: &Path, ) -> anyhow::Result<()> { - let mut settings = settings.clone(); - settings.server.storage.root = storage_dir.display().to_string(); + let settings = settings.clone().with_storage_override(storage_dir); let (object_store, prefix) = serve::build_artifact_object_store(&settings.server)?; let artifact_store = ArtifactStore::new(object_store, prefix); artifact_store.write_metadata(FABRO_VERSION).await?; @@ -2596,8 +2595,12 @@ methods = ["dev-token"] .await .unwrap(); - let mut overridden = settings.clone(); - overridden.server.storage.root = dir.path().display().to_string(); + assert!( + dir.path() + .join("objects/artifacts/store-metadata.json") + .is_file() + ); + let overridden = settings.clone().with_storage_override(dir.path()); let (object_store, prefix) = crate::serve::build_artifact_object_store(&overridden.server).unwrap(); let marker = if prefix.is_empty() { diff --git a/lib/apps/fabro-server/src/serve.rs b/lib/apps/fabro-server/src/serve.rs index 3b969de84..57f2c15b2 100644 --- a/lib/apps/fabro-server/src/serve.rs +++ b/lib/apps/fabro-server/src/serve.rs @@ -1151,7 +1151,7 @@ fn server_bind_title(bind: &Bind) -> String { )] mod tests { use std::io; - use std::path::PathBuf; + use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::task::Poll; @@ -1347,6 +1347,56 @@ mod tests { assert_eq!(root, "/srv/fabro-storage/objects/artifacts"); } + #[test] + fn runtime_server_settings_preserve_custom_artifact_root() { + let settings = server_settings( + r#" +_version = 1 +[server.storage] +root = "/srv/from-disk" +[server.artifacts] +prefix = "selected-prefix" +[server.artifacts.local] +root = "/mnt/artifact-files" +"#, + ); + for storage_root in ["/srv/from-disk", "/srv/from-runtime"] { + let resolved = settings + .clone() + .with_storage_override(Path::new(storage_root)); + assert_eq!(resolved.server.storage.root, storage_root); + assert_eq!(resolved.server.artifacts, settings.server.artifacts); + // Startup and config reload can apply the same override again. + assert_eq!( + resolved + .clone() + .with_storage_override(Path::new(storage_root)), + resolved + ); + } + } + + #[test] + fn runtime_server_settings_preserve_s3_artifact_configuration() { + let settings = server_settings( + r#" +_version = 1 +[server.artifacts] +provider = "s3" +prefix = "selected-prefix" +[server.artifacts.s3] +bucket = "artifact-bucket" +region = "us-east-1" +endpoint = "https://objects.example.test" +path_style = true +"#, + ); + let resolved = settings + .clone() + .with_storage_override(Path::new("/srv/from-runtime")); + assert_eq!(resolved.server.artifacts, settings.server.artifacts); + } + #[test] fn runtime_server_settings_keep_disk_defaults_out_of_manifest_defaults() { let mut resolved = resolved_runtime_settings( diff --git a/lib/foundation/fabro-types/src/dense.rs b/lib/foundation/fabro-types/src/dense.rs index c3355b3b4..4be1bc1c6 100644 --- a/lib/foundation/fabro-types/src/dense.rs +++ b/lib/foundation/fabro-types/src/dense.rs @@ -16,27 +16,19 @@ pub struct ServerSettings { impl ServerSettings { #[must_use] 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"); + 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(); + } + } self.server.storage.root = path.display().to_string(); - override_local_object_store_root(&mut self.server.artifacts.store, path, "artifacts"); self } } -fn override_local_object_store_root( - store: &mut ObjectStoreSettings, - storage_root: &Path, - domain: &str, -) { - let ObjectStoreSettings::Local { root } = store else { - return; - }; - *root = storage_root - .join("objects") - .join(domain) - .display() - .to_string(); -} - #[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] pub struct UserSettings { pub cli: CliNamespace, From 2e500b0e2fe5bd517d783cc28d8bb8bf8c4bf833 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Fri, 25 Sep 2026 13:32:13 -0400 Subject: [PATCH 3/3] Simplify artifact capture writer and storage helpers - The artifact writer takes the digest the hooks already computed, so captured bytes are hashed once on each side, and the dead integrity error goes away. - The writer is a required part of HooksSpec, not an optional field on RunRequest, so a run with capture globs always has a writer and the no-writer error goes away. - ArtifactStore routes put/get and the capture methods through shared put_at/get_at helpers. - The upload handler parses the digest with parse_blob_hash_path before any store reads, and builds its size-limit message from the constant. - The default local artifact root comes from one helper used by both config resolution and the storage-dir override. Co-Authored-By: Claude Opus 5.5 --- .../src/commands/run/petri_worker.rs | 15 +- .../src/server/handler/artifacts.rs | 35 +++-- .../fabro-server/src/server/petri_runs.rs | 8 +- .../src/server/tests/artifact_storage.rs | 7 +- lib/components/fabro-petri/src/artifacts.rs | 26 ++-- lib/components/fabro-petri/src/engine.rs | 30 ++-- lib/components/fabro-petri/src/hooks.rs | 132 ++++++------------ lib/components/fabro-petri/tests/hooks.rs | 17 ++- .../fabro-petri/tests/support/mod.rs | 1 - .../fabro-store/src/artifact_store.rs | 54 +++---- .../fabro-config/src/resolve/server.rs | 10 +- lib/foundation/fabro-types/src/dense.rs | 11 +- .../fabro-types/src/settings/server.rs | 12 ++ 13 files changed, 161 insertions(+), 197 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 9bfc188ac..d4b1ab37d 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -198,8 +198,15 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { } runner::set_worker_title(&run_id, WorkerTitlePhase::Running); - let hooks = HooksSpec::for_run(Arc::clone(&records), &worker.run_state.spec.settings.run) - .with_test_gates(test_checkpoint_gates()); + let hooks = HooksSpec::for_run( + Arc::clone(&records), + &worker.run_state.spec.settings.run, + Arc::new(ClientArtifactWriter::new( + worker.client.clone_for_reuse(), + run_id, + )), + ) + .with_test_gates(test_checkpoint_gates()); let request = RunRequest { run_id: run_id.to_string(), run_dir: worker.run_dir.join("petri"), @@ -223,10 +230,6 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> { worker.client.clone_for_reuse(), run_id, ))), - artifact_writer: Some(Arc::new(ClientArtifactWriter::new( - worker.client.clone_for_reuse(), - run_id, - ))), hooks: Some(hooks), }; let paused_mirror = mirror_paused_state(run_id, &controls, Arc::clone(&records)); diff --git a/lib/apps/fabro-server/src/server/handler/artifacts.rs b/lib/apps/fabro-server/src/server/handler/artifacts.rs index 91a73675e..46ba7fab2 100644 --- a/lib/apps/fabro-server/src/server/handler/artifacts.rs +++ b/lib/apps/fabro-server/src/server/handler/artifacts.rs @@ -26,8 +26,8 @@ use super::super::{ Bytes, HashMap, IntoResponse, Json, NodeArtifact, Path, Query, RequireRunBlob, RequireRunScoped, RequiredUser, Response, Router, RunArtifactEntry, RunArtifactListResponse, RunId, StageArtifactEntry, State, StatusCode, WriteBlobResponse, get, header, - octet_stream_response, parse_run_id_path, parse_stage_id_path, post, reject_if_archived, - required_query_param, validate_relative_artifact_path, + octet_stream_response, parse_blob_hash_path, parse_run_id_path, parse_stage_id_path, post, + reject_if_archived, required_query_param, validate_relative_artifact_path, }; use crate::principal_middleware::RequireWorkerRunSegment; @@ -56,12 +56,19 @@ async fn write_run_artifact_content( State(state): State>, body: Result, ) -> Response { + let expected = match parse_blob_hash_path(&digest) { + Ok(expected) => expected, + Err(response) => return response, + }; let body = match body { Ok(body) => body, Err(error) => { return ApiError::new( error.status(), - "Artifact request body could not be read within the 10 MiB limit.", + format!( + "Artifact request body could not be read within the {} MiB limit.", + ARTIFACT_MAX_FILE_BYTES / (1024 * 1024) + ), ) .into_response(); } @@ -72,15 +79,16 @@ async fn write_run_artifact_content( if let Err(error) = state.load_run_projection(&id).await { return error.into_response(); } - let Ok(expected) = digest.parse::() else { - return ApiError::bad_request("Invalid artifact digest.").into_response(); - }; if BlobHash::new(&body) != expected { return ApiError::bad_request("Artifact content does not match its digest.") .into_response(); } - match state.artifact_store.put_capture(&id, &body).await { - Ok(_) => StatusCode::NO_CONTENT.into_response(), + match state + .artifact_store + .put_capture(&id, &expected, &body) + .await + { + Ok(()) => StatusCode::NO_CONTENT.into_response(), Err(error) => { warn!(run_id = %id, error = %collect_chain(&error).join(": "), "Artifact upload failed"); ApiError::new( @@ -572,9 +580,10 @@ mod tests { ); let blobs = state.store_ref().blobs(); let old = blobs.write(b"old SQLite").await.unwrap(); - let new = state + let new = BlobHash::new(b"new object"); + state .artifact_store - .put_capture(&run, b"new object") + .put_capture(&run, &new, b"new object") .await .unwrap(); for (path, size, source) in [ @@ -603,7 +612,7 @@ mod tests { // Unrecorded content is never a file-list entry. state .artifact_store - .put_capture(&run, b"orphan") + .put_capture(&run, &BlobHash::new(b"orphan"), b"orphan") .await .unwrap(); let entries = run_artifacts(&state, &run, &projection).await.unwrap(); @@ -652,9 +661,5 @@ mod tests { .unwrap() .is_none() ); - // Saved projections preserve both physical sources without rewriting history. - let saved = serde_json::to_value(&projection).unwrap(); - let decoded: RunProjection = serde_json::from_value(saved).unwrap(); - assert_eq!(decoded.artifacts, projection.artifacts); } } diff --git a/lib/apps/fabro-server/src/server/petri_runs.rs b/lib/apps/fabro-server/src/server/petri_runs.rs index be500bb5d..1055f7215 100644 --- a/lib/apps/fabro-server/src/server/petri_runs.rs +++ b/lib/apps/fabro-server/src/server/petri_runs.rs @@ -456,6 +456,10 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { &state.stores.run_summaries, ))), &run_state.spec.settings.run, + Arc::new(StoreArtifactWriter::new( + state.artifact_store.clone(), + run_id, + )), ); let request = RunRequest { run_id: run_id.to_string(), @@ -480,10 +484,6 @@ pub(crate) async fn execute(state: Arc, run_id: RunId) { observers, secrets: Some(Arc::new(VaultSecrets::from_vault(&vault))), blobs: Some(state.store_ref().blobs()), - artifact_writer: Some(Arc::new(StoreArtifactWriter::new( - state.artifact_store.clone(), - run_id, - ))), hooks: Some(hooks), }; let result = Box::pin(engine::run(request)).await; diff --git a/lib/apps/fabro-server/src/server/tests/artifact_storage.rs b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs index cb1546057..5d51cbdb6 100644 --- a/lib/apps/fabro-server/src/server/tests/artifact_storage.rs +++ b/lib/apps/fabro-server/src/server/tests/artifact_storage.rs @@ -89,9 +89,10 @@ async fn artifact_upload_rejects_unauthorized_invalid_and_oversized_bodies_witho let run = create_run_with_bearer(&app, &user).await; let worker = issue_test_worker_token(&run); let other_worker = issue_test_worker_token(&RunId::new()); - let hash = state + let hash = BlobHash::new(b"valid content"); + state .artifact_store - .put_capture(&run, b"valid content") + .put_capture(&run, &hash, b"valid content") .await .unwrap(); let digest = hash.to_string(); @@ -264,7 +265,7 @@ async fn artifact_worker_client_uses_the_configured_s3_backend_and_prefix() { .await .unwrap(); let writer = ClientArtifactWriter::new(client, run); - assert_eq!(writer.write(&bytes).await.unwrap(), hash); + writer.write(&hash, &bytes).await.unwrap(); assert_eq!( state .artifact_store diff --git a/lib/components/fabro-petri/src/artifacts.rs b/lib/components/fabro-petri/src/artifacts.rs index d7eca8376..758b87ce3 100644 --- a/lib/components/fabro-petri/src/artifacts.rs +++ b/lib/components/fabro-petri/src/artifacts.rs @@ -11,19 +11,15 @@ pub enum ArtifactWriteError { Store(#[from] fabro_store::Error), #[error("artifact upload failed")] Upload(#[source] anyhow::Error), - #[error("artifact writer returned {actual}, expected {expected}")] - Integrity { - expected: BlobHash, - actual: BlobHash, - }, } -/// Run-bound storage for captured files. Implementations publish complete -/// objects before returning their content hash, preserve errors, and permit -/// concurrent and repeated writes of identical content. +/// Run-bound storage for captured files, keyed by the `digest` the caller +/// computed over `bytes`. Implementations publish complete objects before +/// returning, preserve errors, and permit concurrent and repeated writes of +/// identical content. #[async_trait::async_trait] pub trait ArtifactWriter: Send + Sync { - async fn write(&self, bytes: &[u8]) -> Result; + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError>; } pub struct StoreArtifactWriter { @@ -40,9 +36,9 @@ impl StoreArtifactWriter { #[async_trait::async_trait] impl ArtifactWriter for StoreArtifactWriter { - async fn write(&self, bytes: &[u8]) -> Result { + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { self.store - .put_capture(&self.run_id, bytes) + .put_capture(&self.run_id, digest, bytes) .await .map_err(ArtifactWriteError::from) } @@ -62,12 +58,10 @@ impl ClientArtifactWriter { #[async_trait::async_trait] impl ArtifactWriter for ClientArtifactWriter { - async fn write(&self, bytes: &[u8]) -> Result { - let hash = BlobHash::new(bytes); + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { self.client - .write_run_artifact_content(&self.run_id, &hash, bytes) + .write_run_artifact_content(&self.run_id, digest, bytes) .await - .map_err(ArtifactWriteError::Upload)?; - Ok(hash) + .map_err(ArtifactWriteError::Upload) } } diff --git a/lib/components/fabro-petri/src/engine.rs b/lib/components/fabro-petri/src/engine.rs index 768af58cd..b322cce1e 100644 --- a/lib/components/fabro-petri/src/engine.rs +++ b/lib/components/fabro-petri/src/engine.rs @@ -63,7 +63,6 @@ use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; use crate::admission::AdmittedGraphs; -use crate::artifacts::ArtifactWriter; use crate::blobs::{Blobs, RunBlobs}; use crate::controls::RunControls; use crate::hooks::{FabroHooks, HooksSpec}; @@ -85,38 +84,36 @@ pub enum Execution { pub struct RunRequest { /// The Fabro run id, which becomes Petri's run key: the run's identity /// in the store and the label on every sandbox of the run. - pub run_id: String, + pub run_id: String, /// Where the run's workspaces, step output and blobs live. - pub run_dir: PathBuf, - pub execution: Execution, + pub run_dir: PathBuf, + pub execution: Execution, /// The run's durable record: the worker's HTTP store, or the server's /// SQLite store under the test override. - pub store: Arc, - pub runtime: RuntimeSpec, + pub store: Arc, + pub runtime: RuntimeSpec, /// The sandbox provider Fabro resolved for the run's environment. - pub provider: SandboxProviderKind, + pub provider: SandboxProviderKind, /// Fires to cancel the run. - pub cancel: CancellationToken, + pub cancel: CancellationToken, /// The run's pause, unpause and steer controls, which the caller keeps /// a clone of to drive them while the run is live. - pub controls: RunControls, + pub controls: RunControls, /// Where the run's questions go. - pub interviewer: Arc, + pub interviewer: Arc, /// The caller's observers of every record, registered ahead of the /// interview dispatcher: the interviewer's own expiry observer among /// them. - pub observers: Vec>, + pub observers: Vec>, /// Where `{{ secrets.NAME }}` references resolve from; `None` leaves /// every secret unknown. - pub secrets: Option>, + pub secrets: Option>, /// Where offloaded stage values go; `None` keeps Petri's local store /// under the run directory. - pub blobs: Option>, - /// Where captured workspace files go. Required when capture writes files. - pub artifact_writer: Option>, + pub blobs: Option>, /// Fabro's hooks: the checkpoint commit and its record. `None` runs /// with Petri's local hook service alone. - pub hooks: Option, + pub hooks: Option, } /// The recorded status of a finished run. @@ -213,7 +210,6 @@ pub async fn run(request: RunRequest) -> Result { Arc::clone(&request.store), resumed, request.blobs.clone(), - request.artifact_writer.clone(), )) }); if let Some(hooks) = &fabro_hooks { diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 29667dd05..03e662b8b 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -185,8 +185,6 @@ pub enum HookError { }, #[error("invalid run.artifacts.include pattern")] Globs(#[source] Arc), - #[error("the run has no configured artifact writer")] - NoArtifactWriter, #[error("the captured artifact could not be stored")] Artifact(#[source] ArtifactWriteError), #[error("the workspace could not be listed below `{root}`")] @@ -215,26 +213,33 @@ impl HookError { /// What Fabro's hooks need beside the run: where the platform records go, /// the run's Git settings, and which files are the run's artifacts. pub struct HooksSpec { - pub records: Arc, - pub git: RunGitSettings, + pub records: Arc, + pub git: RunGitSettings, /// The `[run.artifacts] include` patterns: which files of a stage's /// workspace are collected after the stage. - pub artifacts: Vec, + pub artifacts: Vec, /// A test's gate directory: a checkpoint point named by a `.hold` file /// there waits for its `.release` file. `None` outside tests. - pub test_gates: Option, + pub test_gates: Option, + /// Where captured workspace files go. + pub artifact_writer: Arc, } impl HooksSpec { /// The spec a run's settings give: its Git settings and its artifact - /// patterns. + /// patterns, captured through `artifact_writer`. #[must_use] - pub fn for_run(records: Arc, settings: &RunNamespace) -> Self { + pub fn for_run( + records: Arc, + settings: &RunNamespace, + artifact_writer: Arc, + ) -> Self { Self { records, git: RunGitSettings::from(settings), artifacts: settings.artifacts.include.clone(), test_gates: None, + artifact_writer, } } @@ -372,7 +377,7 @@ pub struct FabroHooks { records: Arc, /// Where diff patches go; `None` records summaries alone. blobs: Option>, - artifact_writer: Option>, + artifact_writer: Arc, workspaces: RunWorkspaces, lookup: WorkspaceLookup, identity: GitIdentity, @@ -399,8 +404,7 @@ impl FabroHooks { /// run whose records are in `store` under `run_key`, with its /// workspaces under `run_dir`. `resumed` says the run continues from /// its records, so a sandbox workspace is brought to its snapshot at - /// its scope's first acquisition. `blobs` holds diff patches; - /// `artifact_writer` holds captured workspace files. + /// its scope's first acquisition. `blobs` holds diff patches. #[must_use] pub fn new( spec: HooksSpec, @@ -411,7 +415,6 @@ impl FabroHooks { store: Arc, resumed: bool, blobs: Option>, - artifact_writer: Option>, ) -> Self { let identity = GitIdentity { name: spec.git.author.name.clone(), @@ -429,7 +432,7 @@ impl FabroHooks { run_id, records: spec.records, blobs, - artifact_writer, + artifact_writer: spec.artifact_writer, workspaces, lookup: WorkspaceLookup::new(Arc::clone(&store), run_key), identity, @@ -949,14 +952,16 @@ impl FabroHooks { return Ok(0); }; let candidates = list_artifacts(env.as_ref(), globs).await?; - let limit = usize::try_from(ARTIFACT_MAX_FILE_BYTES).unwrap_or(usize::MAX); let mut collected = 0; let mut total_bytes = 0_u64; for (path, size) in select_artifacts(candidates) { if total_bytes.saturating_add(size) > ARTIFACT_MAX_TOTAL_BYTES { break; } - let bytes = match env.read_file_limited(Path::new(&path), limit).await { + let bytes = match env + .read_file_limited(Path::new(&path), fabro_types::ARTIFACT_MAX_FILE_BYTES) + .await + { Ok(Some(bytes)) => bytes, Ok(None) => continue, Err(error) => { @@ -964,7 +969,7 @@ impl FabroHooks { continue; } }; - if !self.store_artifact(key, path, &bytes).await? { + if !self.store_artifact(key, &path, &bytes).await? { continue; } total_bytes = total_bytes.saturating_add(size); @@ -978,32 +983,25 @@ impl FabroHooks { async fn store_artifact( &self, key: CheckpointKey, - path: String, + path: &str, bytes: &[u8], ) -> Result { let already = self.collected_artifacts().await?; let digest = BlobHash::new(bytes); - let identity = (path.clone(), digest.to_string()); + let identity = (path.to_owned(), digest.to_string()); if sync::lock(already).contains(&identity) { return Ok(false); } - let writer = self - .artifact_writer - .as_ref() - .ok_or(HookError::NoArtifactWriter)?; - let stored = writer.write(bytes).await.map_err(HookError::Artifact)?; - if stored != digest { - return Err(HookError::Artifact(ArtifactWriteError::Integrity { - expected: digest, - actual: stored, - })); - } + self.artifact_writer + .write(&digest, bytes) + .await + .map_err(HookError::Artifact)?; let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord { execution: key.execution, firing: key.firing, attempt: key.attempt, - path: path.clone(), - source: ArtifactSource::ObjectStore(stored), + path: path.to_owned(), + source: ArtifactSource::ObjectStore(digest), bytes: u64::try_from(bytes.len()).unwrap_or(u64::MAX), digest: digest.to_string(), operation: Some(key.operation_for(ARTIFACT_EFFECT)), @@ -1420,13 +1418,13 @@ mod tests { #[async_trait::async_trait] impl ArtifactWriter for FlakyWriter { - async fn write(&self, bytes: &[u8]) -> Result { + async fn write(&self, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> { if self.fail.swap(false, Ordering::SeqCst) { return Err(ArtifactWriteError::Store(fabro_store::Error::Io( std::io::Error::other("test store unavailable"), ))); } - self.writer.write(bytes).await + self.writer.write(digest, bytes).await } } @@ -1434,7 +1432,7 @@ mod tests { root: &Path, run_id: RunId, records: Arc, - writer: Option>, + artifact_writer: Arc, ) -> FabroHooks { FabroHooks::new( HooksSpec { @@ -1442,6 +1440,7 @@ mod tests { git: RunGitSettings::default(), artifacts: vec!["assets/**".to_string()], test_gates: None, + artifact_writer, }, Arc::new(NoHooks), run_id, @@ -1450,7 +1449,6 @@ mod tests { Arc::new(petri_store::MemoryRunStore::new()), false, None, - writer, ) } @@ -1468,7 +1466,7 @@ mod tests { writer: StoreArtifactWriter::new(store.clone(), run), fail: AtomicBool::new(fail_upload), }); - let hooks = artifact_hooks(root.path(), run, records.clone(), Some(writer.clone())); + let hooks = artifact_hooks(root.path(), run, records.clone(), writer.clone()); let key = CheckpointKey { execution: 0, firing: 1, @@ -1477,10 +1475,7 @@ mod tests { let path = "assets/report.bin".to_string(); let bytes = b"binary\0payload"; let hash = BlobHash::new(bytes); - let error = hooks - .store_artifact(key, path.clone(), bytes) - .await - .unwrap_err(); + let error = hooks.store_artifact(key, &path, bytes).await.unwrap_err(); assert!(error.render().contains(if fail_upload { "test store unavailable" } else { @@ -1493,26 +1488,16 @@ mod tests { !fail_upload ); assert!(store.list_for_run(&run).await.unwrap().is_empty()); + assert!(hooks.store_artifact(key, &path, bytes).await.unwrap()); + assert!(!hooks.store_artifact(key, &path, bytes).await.unwrap()); + let resumed = artifact_hooks(root.path(), run, records.clone(), writer); + assert!(!resumed.store_artifact(key, &path, bytes).await.unwrap()); assert!( - hooks - .store_artifact(key, path.clone(), bytes) + resumed + .store_artifact(key, &path, b"changed") .await .unwrap() ); - assert!( - !hooks - .store_artifact(key, path.clone(), bytes) - .await - .unwrap() - ); - let resumed = artifact_hooks(root.path(), run, records.clone(), Some(writer)); - assert!( - !resumed - .store_artifact(key, path.clone(), bytes) - .await - .unwrap() - ); - assert!(resumed.store_artifact(key, path, b"changed").await.unwrap()); assert_eq!(records.records.records(&run).len(), 2); } } @@ -1552,18 +1537,13 @@ mod tests { root.path(), run, records.clone(), - Some(Arc::new(StoreArtifactWriter::new(store.clone(), run))), - ); - assert!( - !resumed - .store_artifact(key, path.clone(), bytes) - .await - .unwrap() + Arc::new(StoreArtifactWriter::new(store.clone(), run)), ); + assert!(!resumed.store_artifact(key, &path, bytes).await.unwrap()); assert!(store.get_capture(&run, &hash).await.unwrap().is_none()); assert!( resumed - .store_artifact(key, path.clone(), b"changed") + .store_artifact(key, &path, b"changed") .await .unwrap() ); @@ -1576,9 +1556,9 @@ mod tests { root.path(), fork, records.clone(), - Some(Arc::new(StoreArtifactWriter::new(store.clone(), fork))), + Arc::new(StoreArtifactWriter::new(store.clone(), fork)), ); - assert!(fork_hooks.store_artifact(key, path, bytes).await.unwrap()); + assert!(fork_hooks.store_artifact(key, &path, bytes).await.unwrap()); assert_eq!( store.get_capture(&fork, &hash).await.unwrap().unwrap(), bytes.as_slice() @@ -1587,26 +1567,6 @@ mod tests { assert_eq!(records.records(&fork).len(), 1); } - #[tokio::test] - async fn artifact_capture_requires_its_writer_without_falling_back_to_blobs() { - let root = tempfile::tempdir().unwrap(); - let run = RunId::new(); - let records = Arc::new(MemoryPlatformRecords::new()); - let hooks = artifact_hooks(root.path(), run, records.clone(), None); - let key = CheckpointKey { - execution: 0, - firing: 1, - attempt: 1, - }; - assert!(matches!( - hooks - .store_artifact(key, "assets/file".to_string(), b"payload") - .await, - Err(HookError::NoArtifactWriter) - )); - assert!(records.records(&run).is_empty()); - } - #[test] fn the_selection_keeps_the_smallest_files_within_the_budgets() { let mut candidates: Vec<(String, u64)> = (0..(ARTIFACT_MAX_FILES + 5)) diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index a2f0f7f11..026f4655c 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -181,13 +181,17 @@ impl Harness { fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec { HooksSpec { - records: Arc::clone(&self.records) as Arc, - git: RunGitSettings { + records: Arc::clone(&self.records) as Arc, + git: RunGitSettings { host_workspaces: *provider == SandboxProviderKind::LOCAL, ..RunGitSettings::default() }, - artifacts: self.artifacts.clone(), - test_gates: None, + artifacts: self.artifacts.clone(), + test_gates: None, + artifact_writer: Arc::new(StoreArtifactWriter::new( + self.artifact_store.clone(), + self.run_id, + )), } } @@ -220,10 +224,6 @@ impl Harness { observers, secrets: None, blobs: Some(Arc::clone(&self.blobs) as Arc), - artifact_writer: Some(Arc::new(StoreArtifactWriter::new( - self.artifact_store.clone(), - self.run_id, - ))), hooks: Some(hooks), }; engine::run(request).await.expect("the run executes") @@ -851,7 +851,6 @@ async fn a_run_hook_blocks_a_tool_effect_through_the_forwarded_service() { observers, secrets: None, blobs: None, - artifact_writer: None, hooks: Some(harness.hooks(&SandboxProviderKind::LOCAL)), }; let outcome = engine::run(request).await.expect("the run executes"); diff --git a/lib/components/fabro-petri/tests/support/mod.rs b/lib/components/fabro-petri/tests/support/mod.rs index 9fd47f914..7cc8fd0c3 100644 --- a/lib/components/fabro-petri/tests/support/mod.rs +++ b/lib/components/fabro-petri/tests/support/mod.rs @@ -121,7 +121,6 @@ pub(crate) fn run_request( interviewer: Arc::new(interviewer), secrets: None, blobs: None, - artifact_writer: None, hooks: None, } } diff --git a/lib/components/fabro-store/src/artifact_store.rs b/lib/components/fabro-store/src/artifact_store.rs index 2193d2a34..5a24d5760 100644 --- a/lib/components/fabro-store/src/artifact_store.rs +++ b/lib/components/fabro-store/src/artifact_store.rs @@ -74,40 +74,47 @@ impl ArtifactStore { } pub async fn put(&self, run_id: &RunId, key: &ArtifactKey, data: &[u8]) -> Result<()> { - let path = self.artifact_path(run_id, key)?; - self.object_store - .put(&path, Bytes::copy_from_slice(data).into()) - .await?; - Ok(()) + self.put_at(&self.artifact_path(run_id, key)?, data).await } /// Publish a complete captured file under its content digest within the - /// run. Repeating a put of the same content is safe; metadata is recorded + /// run. `hash` must be `BlobHash::new(data)`; callers verify or compute + /// it. Repeating a put of the same content is safe; metadata is recorded /// separately only after this operation succeeds. - pub async fn put_capture(&self, run_id: &RunId, data: &[u8]) -> Result { - let hash = BlobHash::new(data); - let path = self.capture_prefix(run_id)?.child(hash.to_string()); - self.object_store - .put(&path, Bytes::copy_from_slice(data).into()) - .await?; - Ok(hash) + pub async fn put_capture(&self, run_id: &RunId, hash: &BlobHash, data: &[u8]) -> Result<()> { + debug_assert_eq!(*hash, BlobHash::new(data)); + self.put_at(&self.capture_path(run_id, hash)?, data).await } /// Read run-owned content named by a capture record, without consulting /// historical stage keys or the SQLite blob table. pub async fn get_capture(&self, run_id: &RunId, hash: &BlobHash) -> Result> { - let path = self.capture_prefix(run_id)?.child(hash.to_string()); - match self.object_store.get(&path).await { - Ok(result) => Ok(Some(result.bytes().await?)), - Err(object_store::Error::NotFound { .. }) => Ok(None), - Err(error) => Err(error.into()), - } + self.get_at(&self.capture_path(run_id, hash)?).await } fn capture_prefix(&self, run_id: &RunId) -> Result { Ok(self.run_prefix(run_id)?.child("captures").child("sha256")) } + fn capture_path(&self, run_id: &RunId, hash: &BlobHash) -> Result { + Ok(self.capture_prefix(run_id)?.child(hash.to_string())) + } + + async fn put_at(&self, path: &ObjectPath, data: &[u8]) -> Result<()> { + self.object_store + .put(path, Bytes::copy_from_slice(data).into()) + .await?; + Ok(()) + } + + async fn get_at(&self, path: &ObjectPath) -> Result> { + match self.object_store.get(path).await { + Ok(result) => Ok(Some(result.bytes().await?)), + Err(object_store::Error::NotFound { .. }) => Ok(None), + Err(err) => Err(err.into()), + } + } + pub fn writer(&self, run_id: &RunId, key: &ArtifactKey) -> Result { let path = self.artifact_path(run_id, key)?; Ok(BufWriter::with_capacity( @@ -142,12 +149,7 @@ impl ArtifactStore { } pub async fn get(&self, run_id: &RunId, key: &ArtifactKey) -> Result> { - let path = self.artifact_path(run_id, key)?; - match self.object_store.get(&path).await { - Ok(result) => Ok(Some(result.bytes().await?)), - Err(object_store::Error::NotFound { .. }) => Ok(None), - Err(err) => Err(err.into()), - } + self.get_at(&self.artifact_path(run_id, key)?).await } pub async fn get_stream( @@ -472,7 +474,7 @@ mod tests { store.write_metadata("test").await.unwrap(); store.put(&run, &legacy, b"legacy").await.unwrap(); for id in [&run, &run, &other] { - assert_eq!(store.put_capture(id, bytes).await.unwrap(), hash); + store.put_capture(id, &hash, bytes).await.unwrap(); } let location = ObjectPath::from(format!("artifacts/{run}/captures/sha256/{hash}")); assert_eq!( diff --git a/lib/foundation/fabro-config/src/resolve/server.rs b/lib/foundation/fabro-config/src/resolve/server.rs index 1c472664e..8ce62c516 100644 --- a/lib/foundation/fabro-config/src/resolve/server.rs +++ b/lib/foundation/fabro-config/src/resolve/server.rs @@ -285,7 +285,7 @@ fn resolve_artifacts( provider, layer.and_then(|artifacts| artifacts.local.as_ref()), layer.and_then(|artifacts| artifacts.s3.as_ref()), - &object_store_default_root(storage_root, "artifacts"), + &ServerArtifactsSettings::default_local_root(Path::new(storage_root)), "server.artifacts", errors, ), @@ -330,14 +330,6 @@ fn resolve_object_store( } } -fn object_store_default_root(storage_root: &str, domain: &str) -> String { - Path::new(storage_root) - .join("objects") - .join(domain) - .to_string_lossy() - .into_owned() -} - fn resolve_integrations(layer: Option<&ServerIntegrationsLayer>) -> ServerIntegrationsSettings { ServerIntegrationsSettings { github: layer diff --git a/lib/foundation/fabro-types/src/dense.rs b/lib/foundation/fabro-types/src/dense.rs index 4be1bc1c6..3fe7e5cd2 100644 --- a/lib/foundation/fabro-types/src/dense.rs +++ b/lib/foundation/fabro-types/src/dense.rs @@ -4,8 +4,8 @@ use std::path::Path; use serde::{Deserialize, Serialize}; use crate::settings::{ - CliNamespace, ObjectStoreSettings, ProjectNamespace, RunNamespace, ServerNamespace, - WorkflowNamespace, + CliNamespace, ObjectStoreSettings, ProjectNamespace, RunNamespace, ServerArtifactsSettings, + ServerNamespace, WorkflowNamespace, }; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -18,10 +18,11 @@ impl ServerSettings { pub fn with_storage_override(mut self, path: &Path) -> Self { // Only the derived default follows the storage directory. A custom // artifact location is independent of the database and runtime root. - let default_artifact_root = Path::new(&self.server.storage.root).join("objects/artifacts"); + let default_artifact_root = + ServerArtifactsSettings::default_local_root(Path::new(&self.server.storage.root)); if let ObjectStoreSettings::Local { root } = &mut self.server.artifacts.store { - if Path::new(root) == default_artifact_root { - *root = path.join("objects/artifacts").display().to_string(); + if *root == default_artifact_root { + *root = ServerArtifactsSettings::default_local_root(path); } } self.server.storage.root = path.display().to_string(); diff --git a/lib/foundation/fabro-types/src/settings/server.rs b/lib/foundation/fabro-types/src/settings/server.rs index 4f4279ded..9860db349 100644 --- a/lib/foundation/fabro-types/src/settings/server.rs +++ b/lib/foundation/fabro-types/src/settings/server.rs @@ -216,6 +216,18 @@ pub struct ServerArtifactsSettings { pub store: ObjectStoreSettings, } +impl ServerArtifactsSettings { + /// The local artifact root a storage root gives when none is configured. + #[must_use] + pub fn default_local_root(storage_root: &std::path::Path) -> String { + storage_root + .join("objects") + .join("artifacts") + .to_string_lossy() + .into_owned() + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "type", rename_all = "snake_case")] pub enum ObjectStoreSettings {