This commit is contained in:
Scott Werner 2026-09-25 17:32:19 +00:00 • committed by GitHub
commit 7ba70ca964
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
39 changed files with 1773 additions and 180 deletions

1
Cargo.lock generated
View file

@ -2539,6 +2539,7 @@ dependencies = [
"fabro-workflow",
"httpmock",
"lithos-llm",
"object_store",
"pebble-coding-agent",
"petri-attractor-steps",
"petri-execution",

View file

@ -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 `<artifact prefix>/<run ID>/captures/sha256/<digest>`
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.

View file

@ -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 `<storage_root>/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

View file

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

View file

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

View file

@ -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};
@ -197,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"),

View file

@ -1,15 +1,139 @@
//! `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 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]
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"]
);
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),
("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

View file

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

View file

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

View file

@ -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 _;
@ -23,12 +26,17 @@ 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;
pub(super) fn routes() -> Router<Arc<AppState>> {
Router::new()
.route(
"/runs/{id}/artifacts/content/{digest}",
put(write_run_artifact_content).layer(DefaultBodyLimit::max(ARTIFACT_MAX_FILE_BYTES)),
)
.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,55 @@ pub(super) fn routes() -> Router<Arc<AppState>> {
)
}
async fn write_run_artifact_content(
RequireWorkerRunSegment(id, digest): RequireWorkerRunSegment,
State(state): State<Arc<AppState>>,
body: Result<Bytes, BytesRejection>,
) -> Response {
let expected = match parse_blob_hash_path(&digest) {
Ok(expected) => expected,
Err(response) => return response,
};
let body = match body {
Ok(body) => body,
Err(error) => {
return ApiError::new(
error.status(),
format!(
"Artifact request body could not be read within the {} MiB limit.",
ARTIFACT_MAX_FILE_BYTES / (1024 * 1024)
),
)
.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();
}
if BlobHash::new(&body) != expected {
return ApiError::bad_request("Artifact content does not match its digest.")
.into_response();
}
match state
.artifact_store
.put_capture(&id, &expected, &body)
.await
{
Ok(()) => StatusCode::NO_CONTENT.into_response(),
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 +143,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 +172,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 +310,10 @@ async fn read_artifact(
bytes: ArtifactBytes,
) -> Result<Option<Bytes>, 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 +542,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 +560,106 @@ 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 = BlobHash::new(b"new object");
state
.artifact_store
.put_capture(&run, &new, 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, &BlobHash::new(b"orphan"), 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()
);
}
}

View file

@ -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;
@ -455,6 +456,10 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
&state.stores.run_summaries,
))),
&run_state.spec.settings.run,
Arc::new(StoreArtifactWriter::new(
state.artifact_store.clone(),
run_id,
)),
);
let request = RunRequest {
run_id: run_id.to_string(),

View file

@ -4953,8 +4953,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,
@ -10969,3 +10968,5 @@ fn an_unsettled_run_takes_the_stores_status_or_a_termination_at_worker_exit() {
RunStatus::Running
);
}
mod artifact_storage;

View file

@ -0,0 +1,289 @@
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<Body> {
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 = BlobHash::new(b"valid content");
state
.artifact_store
.put_capture(&run, &hash, 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);
writer.write(&hash, &bytes).await.unwrap();
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;
}

View file

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

View file

@ -105,8 +105,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
@ -188,8 +189,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.

View file

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

View file

@ -0,0 +1,67 @@
//! 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),
}
/// 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, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError>;
}
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, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> {
self.store
.put_capture(&self.run_id, digest, 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, digest: &BlobHash, bytes: &[u8]) -> Result<(), ArtifactWriteError> {
self.client
.write_run_artifact_content(&self.run_id, digest, bytes)
.await
.map_err(ArtifactWriteError::Upload)
}
}

View file

@ -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,8 @@ pub enum HookError {
},
#[error("invalid run.artifacts.include pattern")]
Globs(#[source] Arc<WorkspaceGlobError>),
#[error("the run has no blob table to collect artifacts into")]
NoBlobs,
#[error("the captured artifact could not be stored")]
Artifact(#[source] ArtifactWriteError),
#[error("the workspace could not be listed below `{root}`")]
List {
root: String,
@ -212,26 +213,33 @@ impl HookError {
/// What Fabro's hooks need beside the run: where the platform records go,
/// the run's Git settings, and which files are the run's artifacts.
pub struct HooksSpec {
pub records: Arc<dyn PlatformRecords>,
pub git: RunGitSettings,
pub records: Arc<dyn PlatformRecords>,
pub git: RunGitSettings,
/// The `[run.artifacts] include` patterns: which files of a stage's
/// workspace are collected after the stage.
pub artifacts: Vec<String>,
pub artifacts: Vec<String>,
/// A test's gate directory: a checkpoint point named by a `.hold` file
/// there waits for its `.release` file. `None` outside tests.
pub test_gates: Option<PathBuf>,
pub test_gates: Option<PathBuf>,
/// Where captured workspace files go.
pub artifact_writer: Arc<dyn ArtifactWriter>,
}
impl HooksSpec {
/// The spec a run's settings give: its Git settings and its artifact
/// patterns.
/// patterns, captured through `artifact_writer`.
#[must_use]
pub fn for_run(records: Arc<dyn PlatformRecords>, settings: &RunNamespace) -> Self {
pub fn for_run(
records: Arc<dyn PlatformRecords>,
settings: &RunNamespace,
artifact_writer: Arc<dyn ArtifactWriter>,
) -> Self {
Self {
records,
git: RunGitSettings::from(settings),
artifacts: settings.artifacts.include.clone(),
test_gates: None,
artifact_writer,
}
}
@ -367,9 +375,9 @@ pub struct FabroHooks {
inner: Arc<dyn ExecutionHooks>,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
/// 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<Arc<dyn Blobs>>,
artifact_writer: Arc<dyn ArtifactWriter>,
workspaces: RunWorkspaces,
lookup: WorkspaceLookup,
identity: GitIdentity,
@ -396,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` is where artifact bytes and
/// diff patches go.
/// its scope's first acquisition. `blobs` holds diff patches.
#[must_use]
pub fn new(
spec: HooksSpec,
@ -425,6 +432,7 @@ impl FabroHooks {
run_id,
records: spec.records,
blobs,
artifact_writer: spec.artifact_writer,
workspaces,
lookup: WorkspaceLookup::new(Arc::clone(&store), run_key),
identity,
@ -943,19 +951,17 @@ 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;
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) => {
@ -963,49 +969,61 @@ 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: &str,
bytes: &[u8],
) -> Result<bool, HookError> {
let already = self.collected_artifacts().await?;
let digest = BlobHash::new(bytes);
let identity = (path.to_owned(), digest.to_string());
if sync::lock(already).contains(&identity) {
return Ok(false);
}
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.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)),
});
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<HashSet<ArtifactIdentity>>, HookError> {
@ -1351,7 +1369,203 @@ 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<StagePosition>,
) -> Result<StoredPlatformRecord, PlatformRecordError> {
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<Vec<StoredPlatformRecord>, 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, 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(digest, bytes).await
}
}
fn artifact_hooks(
root: &Path,
run_id: RunId,
records: Arc<dyn PlatformRecords>,
artifact_writer: Arc<dyn ArtifactWriter>,
) -> FabroHooks {
FabroHooks::new(
HooksSpec {
records,
git: RunGitSettings::default(),
artifacts: vec!["assets/**".to_string()],
test_gates: None,
artifact_writer,
},
Arc::new(NoHooks),
run_id,
RunKey::new(run_id.to_string()),
root.to_path_buf(),
Arc::new(petri_store::MemoryRunStore::new()),
false,
None,
)
}
#[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(), 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, 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, 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!(
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(),
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, 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(),
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);
}
#[test]
fn the_selection_keeps_the_smallest_files_within_the_budgets() {

View file

@ -62,6 +62,7 @@
//! workspace `Cargo.toml` for how they are tracked.
pub mod admission;
pub mod artifacts;
pub mod blobs;
pub mod check;
pub mod checkpoint;

View file

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

View file

@ -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,39 +144,54 @@ fn docker_plugin() -> Option<PathBuf> {
/// One run's pieces: the store, its platform records, where it ran.
struct Harness {
run_id: RunId,
run_dir: PathBuf,
store: Arc<MemoryRunStore>,
records: Arc<MemoryPlatformRecords>,
blobs: Arc<MemoryBlobs>,
run_id: RunId,
run_dir: PathBuf,
store: Arc<MemoryRunStore>,
records: Arc<MemoryPlatformRecords>,
blobs: Arc<MemoryBlobs>,
/// The `[run.artifacts] include` patterns the hooks collect under.
artifacts: Vec<String>,
_root: tempfile::TempDir,
artifacts: Vec<String>,
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,
}
}
fn hooks(&self, provider: &SandboxProviderKind) -> HooksSpec {
HooksSpec {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings {
records: Arc::clone(&self.records) as Arc<dyn PlatformRecords>,
git: RunGitSettings {
host_workspaces: *provider == SandboxProviderKind::LOCAL,
..RunGitSettings::default()
},
artifacts: self.artifacts.clone(),
test_gates: None,
artifacts: self.artifacts.clone(),
test_gates: None,
artifact_writer: Arc::new(StoreArtifactWriter::new(
self.artifact_store.clone(),
self.run_id,
)),
}
}
@ -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.

View file

@ -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;
@ -74,13 +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.put_at(&self.artifact_path(run_id, key)?, data).await
}
/// Publish a complete captured file under its content digest within the
/// run. `hash` must be `BlobHash::new(data)`; callers verify or compute
/// it. Repeating a put of the same content is safe; metadata is recorded
/// separately only after this operation succeeds.
pub async fn put_capture(&self, run_id: &RunId, hash: &BlobHash, data: &[u8]) -> Result<()> {
debug_assert_eq!(*hash, BlobHash::new(data));
self.put_at(&self.capture_path(run_id, hash)?, data).await
}
/// Read run-owned content named by a capture record, without consulting
/// historical stage keys or the SQLite blob table.
pub async fn get_capture(&self, run_id: &RunId, hash: &BlobHash) -> Result<Option<Bytes>> {
self.get_at(&self.capture_path(run_id, hash)?).await
}
fn capture_prefix(&self, run_id: &RunId) -> Result<ObjectPath> {
Ok(self.run_prefix(run_id)?.child("captures").child("sha256"))
}
fn capture_path(&self, run_id: &RunId, hash: &BlobHash) -> Result<ObjectPath> {
Ok(self.capture_prefix(run_id)?.child(hash.to_string()))
}
async fn put_at(&self, path: &ObjectPath, data: &[u8]) -> Result<()> {
self.object_store
.put(&path, Bytes::copy_from_slice(data).into())
.put(path, Bytes::copy_from_slice(data).into())
.await?;
Ok(())
}
async fn get_at(&self, path: &ObjectPath) -> Result<Option<Bytes>> {
match self.object_store.get(path).await {
Ok(result) => Ok(Some(result.bytes().await?)),
Err(object_store::Error::NotFound { .. }) => Ok(None),
Err(err) => Err(err.into()),
}
}
pub fn writer(&self, run_id: &RunId, key: &ArtifactKey) -> Result<BufWriter> {
let path = self.artifact_path(run_id, key)?;
Ok(BufWriter::with_capacity(
@ -115,12 +149,7 @@ impl ArtifactStore {
}
pub async fn get(&self, run_id: &RunId, key: &ArtifactKey) -> Result<Option<Bytes>> {
let path = self.artifact_path(run_id, key)?;
match self.object_store.get(&path).await {
Ok(result) => Ok(Some(result.bytes().await?)),
Err(object_store::Error::NotFound { .. }) => Ok(None),
Err(err) => Err(err.into()),
}
self.get_at(&self.artifact_path(run_id, key)?).await
}
pub async fn get_stream(
@ -152,9 +181,15 @@ impl ArtifactStore {
pub async fn list_for_run(&self, run_id: &RunId) -> Result<Vec<NodeArtifact>> {
let prefix = self.run_prefix(run_id)?;
let captures = self.capture_prefix(run_id)?;
let 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 +462,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<dyn ObjectStore> = 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] {
store.put_capture(id, &hash, bytes).await.unwrap();
}
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<dyn ObjectStore> = Arc::new(InMemory::new());
let store = ArtifactStore::new(objects.clone(), "artifacts");
let run = RunId::new();
let location = ObjectPath::from(format!("artifacts/{run}/captures/elsewhere/bad"));
objects
.put(&location, Bytes::from_static(b"bad").into())
.await
.unwrap();
assert!(store.list_for_run(&run).await.is_err());
}
#[tokio::test]
async fn write_metadata_persists_store_marker() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());

View file

@ -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<OperationKey>,
}
/// Private decoding boundary for validating the redundant legacy checksum.
#[derive(Deserialize)]
struct ArtifactCollectedWire {
execution: u64,
firing: u64,
attempt: u32,
path: String,
#[serde(flatten)]
source: ArtifactSource,
bytes: u64,
digest: String,
#[serde(default)]
operation: Option<OperationKey>,
}
impl TryFrom<ArtifactCollectedWire> for ArtifactCollectedRecord {
type Error = &'static str;
fn try_from(wire: ArtifactCollectedWire) -> std::result::Result<Self, Self::Error> {
if wire.digest != wire.source.hash().to_string() {
return Err("artifact checksum does not match its payload source");
}
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::<PlatformRecord>(wire.clone()).is_err());
wire["digest"] = json!(hash.to_string());
wire["blob"] = json!(hash);
wire["object"] = json!(hash);
assert!(serde_json::from_value::<PlatformRecord>(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 {

View file

@ -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", &[]),

View file

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

View file

@ -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::<types::ArtifactSource>(),
TypeId::of::<ArtifactSource>()
);
assert_eq!(
TypeId::of::<types::RunArtifact>(),
TypeId::of::<RunArtifact>()
);
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::<types::ArtifactSource>(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::<types::RunArtifact>(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::<types::ArtifactSource>(json).is_err());
}
}

View file

@ -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<BlobHash> {
let response = self
.send_api(|client| async move {

View file

@ -285,7 +285,7 @@ fn resolve_artifacts(
provider,
layer.and_then(|artifacts| artifacts.local.as_ref()),
layer.and_then(|artifacts| artifacts.s3.as_ref()),
&object_store_default_root(storage_root, "artifacts"),
&ServerArtifactsSettings::default_local_root(Path::new(storage_root)),
"server.artifacts",
errors,
),
@ -330,14 +330,6 @@ fn resolve_object_store(
}
}
fn object_store_default_root(storage_root: &str, domain: &str) -> String {
Path::new(storage_root)
.join("objects")
.join(domain)
.to_string_lossy()
.into_owned()
}
fn resolve_integrations(layer: Option<&ServerIntegrationsLayer>) -> ServerIntegrationsSettings {
ServerIntegrationsSettings {
github: layer

View file

@ -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<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
// 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<BlobHash>,
#[serde(default, deserialize_with = "present_hash")]
object: Option<BlobHash>,
}
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<Option<BlobHash>, 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;

View file

@ -4,8 +4,8 @@ use std::path::Path;
use serde::{Deserialize, Serialize};
use crate::settings::{
CliNamespace, ObjectStoreSettings, ProjectNamespace, RunNamespace, ServerNamespace,
WorkflowNamespace,
CliNamespace, ObjectStoreSettings, ProjectNamespace, RunNamespace, ServerArtifactsSettings,
ServerNamespace, WorkflowNamespace,
};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
@ -16,27 +16,20 @@ 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 =
ServerArtifactsSettings::default_local_root(Path::new(&self.server.storage.root));
if let ObjectStoreSettings::Local { root } = &mut self.server.artifacts.store {
if *root == default_artifact_root {
*root = ServerArtifactsSettings::default_local_root(path);
}
}
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,

View file

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

View file

@ -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::<RunArtifact>(json.clone()).is_err(),
"{json}"
);
}
}
}
#[cfg(test)]
mod title_tests {
use chrono::Utc;

View file

@ -216,6 +216,18 @@ pub struct ServerArtifactsSettings {
pub store: ObjectStoreSettings,
}
impl ServerArtifactsSettings {
/// The local artifact root a storage root gives when none is configured.
#[must_use]
pub fn default_local_root(storage_root: &std::path::Path) -> String {
storage_root
.join("objects")
.join("artifacts")
.to_string_lossy()
.into_owned()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ObjectStoreSettings {

View file

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

View file

@ -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<RequestArgs> => {
// 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<void>> {
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<WriteBlobResponse> {
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<void> {
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

View file

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

View file

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

View file

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

View file

@ -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<RunArtifact>;
/**
* Sequence-tagged checkpoint history entries.
*/