mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Merge pull request #898 from fabro-sh/codex/restore-configured-artifact-storage
Some checks failed
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run
Rust / Test (macOS) (push) Waiting to run
Rust / Process titles (musl) (push) Waiting to run
TypeScript / Typecheck (push) Has been cancelled
TypeScript / Test (push) Has been cancelled
TypeScript / Build (push) Has been cancelled
Some checks failed
Rust / Format (push) Waiting to run
Rust / Clippy (push) Waiting to run
Rust / Rustdoc (push) Waiting to run
Rust / Generated Docs (push) Waiting to run
Rust / Test (Linux) (push) Waiting to run
Rust / Sandbox providers (Docker) (push) Waiting to run
Rust / Test (macOS) (push) Waiting to run
Rust / Process titles (musl) (push) Waiting to run
TypeScript / Typecheck (push) Has been cancelled
TypeScript / Test (push) Has been cancelled
TypeScript / Build (push) Has been cancelled
Restore configured artifact storage for workflow captures
This commit is contained in:
commit
24fb18869c
39 changed files with 2065 additions and 191 deletions
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -2540,6 +2540,7 @@ dependencies = [
|
|||
"fabro-workflow",
|
||||
"httpmock",
|
||||
"lithos-llm",
|
||||
"object_store",
|
||||
"pebble-coding-agent",
|
||||
"petri-attractor-steps",
|
||||
"petri-execution",
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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};
|
||||
|
|
@ -198,8 +199,12 @@ 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())),
|
||||
)
|
||||
.with_test_gates(test_checkpoint_gates());
|
||||
let request = RunRequest {
|
||||
run_id: run_id.to_string(),
|
||||
run_dir: worker.run_dir.join("petri"),
|
||||
|
|
|
|||
|
|
@ -1,15 +1,129 @@
|
|||
//! `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, run_detached, wait_for_success};
|
||||
use super::petri::{RunningServer, 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() {
|
||||
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
|
||||
|
|
|
|||
|
|
@ -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() {
|
||||
|
|
|
|||
|
|
@ -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,64 @@ mod tests {
|
|||
assert_eq!(root, "/srv/fabro-storage/objects/artifacts");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn runtime_server_settings_keep_local_artifacts_under_storage_dir() {
|
||||
// Servers have always kept local artifacts under the storage
|
||||
// directory, including installs whose settings name a `local.root`.
|
||||
let settings = server_settings(
|
||||
r#"
|
||||
_version = 1
|
||||
[server.storage]
|
||||
root = "/srv/from-disk"
|
||||
[server.artifacts]
|
||||
prefix = "artifacts"
|
||||
[server.artifacts.local]
|
||||
root = "/srv/from-disk/objects"
|
||||
"#,
|
||||
);
|
||||
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.prefix, "artifacts");
|
||||
assert_eq!(
|
||||
resolved.server.artifacts.store,
|
||||
fabro_types::settings::ObjectStoreSettings::Local {
|
||||
root: format!("{storage_root}/objects/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(
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
use fabro_store::{ArtifactStore, BlobStore, Error as StoreError};
|
||||
use fabro_types::{BlobHash, RunProjection};
|
||||
use fabro_types::{ARTIFACT_MAX_FILE_BYTES, ArtifactSource, RunProjection};
|
||||
use fabro_util::error::collect_chain;
|
||||
use futures_util::SinkExt as _;
|
||||
use futures_util::io::AsyncWriteExt as _;
|
||||
|
|
@ -23,12 +26,18 @@ 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, run_records, validate_relative_artifact_path,
|
||||
};
|
||||
use crate::principal_middleware::RequireWorkerRunSegment;
|
||||
|
||||
pub(super) fn routes() -> Router<Arc<AppState>> {
|
||||
Router::new()
|
||||
.route(
|
||||
"/runs/{id}/artifacts/content/{digest}",
|
||||
routing::put(write_run_artifact_content)
|
||||
.layer(DefaultBodyLimit::max(ARTIFACT_MAX_FILE_BYTES)),
|
||||
)
|
||||
.route("/runs/{id}/blobs", post(write_run_blob))
|
||||
.route("/runs/{id}/blobs/{blobHash}", get(read_run_blob))
|
||||
.route("/runs/{id}/artifacts", get(list_run_artifacts))
|
||||
|
|
@ -43,6 +52,66 @@ 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();
|
||||
}
|
||||
match state.artifact_store.put_capture(&id, &expected, body).await {
|
||||
Ok(()) => {}
|
||||
Err(StoreError::CaptureDigestMismatch { .. }) => {
|
||||
return ApiError::bad_request("Artifact content does not match its digest.")
|
||||
.into_response();
|
||||
}
|
||||
Err(error) => {
|
||||
warn!(run_id = %id, error = %collect_chain(&error).join(": "), "Artifact upload failed");
|
||||
return ApiError::new(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
"Artifact storage failed.",
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
}
|
||||
// Run deletion removes the run before its objects. A run still present
|
||||
// now is deleted after this upload landed, and its deletion removes the
|
||||
// upload; a run already gone never will, so the upload goes here.
|
||||
match run_records::projection(state.as_ref(), id).await {
|
||||
Ok(Some(_)) => StatusCode::NO_CONTENT.into_response(),
|
||||
Ok(None) => {
|
||||
if let Err(cleanup) = state.artifact_store.delete_capture(&id, &expected).await {
|
||||
warn!(run_id = %id, error = %collect_chain(&cleanup).join(": "), "Artifact upload for a deleted run could not be removed");
|
||||
}
|
||||
ApiError::not_found("Run not found.").into_response()
|
||||
}
|
||||
Err(error) => {
|
||||
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct ArtifactFilenameParams {
|
||||
#[serde(default)]
|
||||
|
|
@ -86,18 +155,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 +184,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 +322,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 +554,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 +572,110 @@ async fn get_stage_artifact(
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use async_zip::base::read::mem::ZipFileReader;
|
||||
use fabro_types::{BlobHash, 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, Bytes::from_static(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"),
|
||||
Bytes::from_static(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()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
@ -464,6 +465,7 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
|
|||
&state.stores.run_summaries,
|
||||
))),
|
||||
&run_state.spec.settings.run,
|
||||
Arc::new(StoreArtifactWriter::new(state.artifact_store.clone())),
|
||||
);
|
||||
let runtime = runtime_spec(
|
||||
&state,
|
||||
|
|
|
|||
|
|
@ -4883,8 +4883,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,
|
||||
|
|
@ -10899,3 +10898,5 @@ fn an_unsettled_run_takes_the_stores_status_or_a_termination_at_worker_exit() {
|
|||
RunStatus::Running
|
||||
);
|
||||
}
|
||||
|
||||
mod artifact_storage;
|
||||
|
|
|
|||
292
lib/apps/fabro-server/src/server/tests/artifact_storage.rs
Normal file
292
lib/apps/fabro-server/src/server/tests/artifact_storage.rs
Normal file
|
|
@ -0,0 +1,292 @@
|
|||
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, Bytes::from_static(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);
|
||||
writer
|
||||
.write(&run, &hash, Bytes::from(bytes.clone()))
|
||||
.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;
|
||||
}
|
||||
|
|
@ -60,6 +60,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"] }
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -187,8 +188,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.
|
||||
|
|
|
|||
|
|
@ -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 |
|
||||
|
|
|
|||
82
lib/components/fabro-petri/src/artifacts.rs
Normal file
82
lib/components/fabro-petri/src/artifacts.rs
Normal file
|
|
@ -0,0 +1,82 @@
|
|||
//! Captured workspace files go to the server's configured artifact store.
|
||||
//! Engine values and patches keep using the separate blob capability.
|
||||
|
||||
use bytes::Bytes;
|
||||
use fabro_client::Client;
|
||||
use fabro_store::ArtifactStore;
|
||||
use fabro_types::{BlobHash, RunId};
|
||||
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum ArtifactWriteError {
|
||||
#[error("artifact storage failed")]
|
||||
Store(#[from] fabro_store::Error),
|
||||
#[error("artifact upload failed")]
|
||||
Upload(#[source] anyhow::Error),
|
||||
}
|
||||
|
||||
/// Storage for captured files, keyed by the run the hooks record them under
|
||||
/// and the `digest` of `bytes`. Implementations reject bytes that do not
|
||||
/// match `digest`, publish complete objects before returning, preserve
|
||||
/// errors, and permit concurrent and repeated writes of identical content.
|
||||
#[async_trait::async_trait]
|
||||
pub trait ArtifactWriter: Send + Sync {
|
||||
async fn write(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
digest: &BlobHash,
|
||||
bytes: Bytes,
|
||||
) -> Result<(), ArtifactWriteError>;
|
||||
}
|
||||
|
||||
pub struct StoreArtifactWriter {
|
||||
store: ArtifactStore,
|
||||
}
|
||||
|
||||
impl StoreArtifactWriter {
|
||||
#[must_use]
|
||||
pub fn new(store: ArtifactStore) -> Self {
|
||||
Self { store }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ArtifactWriter for StoreArtifactWriter {
|
||||
async fn write(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
digest: &BlobHash,
|
||||
bytes: Bytes,
|
||||
) -> Result<(), ArtifactWriteError> {
|
||||
self.store
|
||||
.put_capture(run_id, digest, bytes)
|
||||
.await
|
||||
.map_err(ArtifactWriteError::from)
|
||||
}
|
||||
}
|
||||
|
||||
/// Uploads to the server, which checks the digest before storing.
|
||||
pub struct ClientArtifactWriter {
|
||||
client: Client,
|
||||
}
|
||||
|
||||
impl ClientArtifactWriter {
|
||||
#[must_use]
|
||||
pub fn new(client: Client) -> Self {
|
||||
Self { client }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ArtifactWriter for ClientArtifactWriter {
|
||||
async fn write(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
digest: &BlobHash,
|
||||
bytes: Bytes,
|
||||
) -> Result<(), ArtifactWriteError> {
|
||||
self.client
|
||||
.write_run_artifact_content(run_id, digest, bytes)
|
||||
.await
|
||||
.map_err(ArtifactWriteError::Upload)
|
||||
}
|
||||
}
|
||||
|
|
@ -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`
|
||||
|
|
@ -76,12 +76,14 @@ use std::path::{Path, PathBuf};
|
|||
use std::sync::{Arc, Mutex, OnceLock};
|
||||
use std::time::Duration;
|
||||
|
||||
use bytes::Bytes;
|
||||
use fabro_store::platform_records::{
|
||||
ArtifactCollectedRecord, CheckpointRecord, GitIdentityRecord, RunBranchRecord, RunDiffRecord,
|
||||
ArtifactCollectedRecord, CheckpointRecord, GitIdentityRecord, OperationKey, RunBranchRecord,
|
||||
RunDiffRecord,
|
||||
};
|
||||
use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition};
|
||||
use fabro_types::settings::run::RunNamespace;
|
||||
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 +100,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 +123,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 +187,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 +215,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,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -248,7 +258,7 @@ impl HooksSpec {
|
|||
type AcquiredEnv = (String, Arc<dyn ExecEnv>);
|
||||
|
||||
/// The identity of a collected file: its path and content digest.
|
||||
type ArtifactIdentity = (String, String);
|
||||
type ArtifactIdentity = (String, BlobHash);
|
||||
|
||||
/// A checkpoint's workspace and commit.
|
||||
type WorkspaceCommit = (String, String);
|
||||
|
|
@ -306,6 +316,9 @@ struct ArtifactLedger {
|
|||
/// Every artifact collected so far, by path and digest: read from the
|
||||
/// store once, then kept current with every append.
|
||||
collected: OnceCell<Mutex<HashSet<ArtifactIdentity>>>,
|
||||
/// One lock per identity being captured, so concurrent transitions that
|
||||
/// leave the same file upload and record it once.
|
||||
capturing: Mutex<HashMap<ArtifactIdentity, Arc<AsyncMutex<()>>>>,
|
||||
}
|
||||
|
||||
/// Where each acquired scope's workspace is, and the locks that serialize
|
||||
|
|
@ -367,9 +380,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 +409,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 +437,7 @@ impl FabroHooks {
|
|||
run_id,
|
||||
records: spec.records,
|
||||
blobs,
|
||||
artifact_writer: spec.artifact_writer,
|
||||
workspaces,
|
||||
lookup: WorkspaceLookup::new(Arc::clone(&store), run_key),
|
||||
identity,
|
||||
|
|
@ -434,6 +447,7 @@ impl FabroHooks {
|
|||
checkpoints: CheckpointLedger::default(),
|
||||
artifacts: ArtifactLedger {
|
||||
globs: WorkspaceGlobSet::try_new(&spec.artifacts).map_err(Arc::new),
|
||||
capturing: Mutex::default(),
|
||||
collected: OnceCell::new(),
|
||||
},
|
||||
scopes: ScopeEnvs::default(),
|
||||
|
|
@ -943,19 +957,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 +975,104 @@ 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.into()).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: Bytes,
|
||||
) -> Result<bool, HookError> {
|
||||
let already = self.collected_artifacts().await?;
|
||||
let digest = BlobHash::new(&bytes);
|
||||
let identity = (path.to_owned(), digest);
|
||||
if sync::lock(already).contains(&identity) {
|
||||
return Ok(false);
|
||||
}
|
||||
let capture = Arc::clone(
|
||||
sync::lock(&self.artifacts.capturing)
|
||||
.entry(identity.clone())
|
||||
.or_default(),
|
||||
);
|
||||
let _capturing = capture.lock().await;
|
||||
// Whoever held the lock before may have recorded this identity.
|
||||
if sync::lock(already).contains(&identity) {
|
||||
return Ok(false);
|
||||
}
|
||||
let size = u64::try_from(bytes.len()).unwrap_or(u64::MAX);
|
||||
self.artifact_writer
|
||||
.write(&self.run_id, &digest, bytes)
|
||||
.await
|
||||
.map_err(HookError::Artifact)?;
|
||||
let operation = key.operation_for(ARTIFACT_EFFECT);
|
||||
let record = PlatformRecord::ArtifactCollected(ArtifactCollectedRecord {
|
||||
execution: key.execution,
|
||||
firing: key.firing,
|
||||
attempt: key.attempt,
|
||||
path: path.to_owned(),
|
||||
source: ArtifactSource::ObjectStore(digest),
|
||||
bytes: size,
|
||||
operation: Some(operation.clone()),
|
||||
});
|
||||
let appended = self
|
||||
.records
|
||||
.append(
|
||||
&self.run_id,
|
||||
&record,
|
||||
Some(StagePosition {
|
||||
execution: key.execution,
|
||||
firing: key.firing,
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
if let Err(source) = appended {
|
||||
// The append can reach the store and lose only its response. A
|
||||
// record that landed is this capture; appending it again would
|
||||
// list the file twice.
|
||||
if !self.artifact_recorded(&identity, &operation).await {
|
||||
return Err(HookError::Write {
|
||||
kind: "artifact",
|
||||
source,
|
||||
});
|
||||
}
|
||||
}
|
||||
sync::lock(already).insert(identity.clone());
|
||||
sync::lock(&self.artifacts.capturing).remove(&identity);
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
/// Whether the run's records hold this capture's record, read fresh.
|
||||
async fn artifact_recorded(
|
||||
&self,
|
||||
identity: &ArtifactIdentity,
|
||||
operation: &OperationKey,
|
||||
) -> bool {
|
||||
let Ok(stored) = self
|
||||
.records
|
||||
.read_kind(&self.run_id, PlatformRecordKind::ArtifactCollected)
|
||||
.await
|
||||
else {
|
||||
return false;
|
||||
};
|
||||
stored.into_iter().any(|record| match record.record {
|
||||
PlatformRecord::ArtifactCollected(artifact) => {
|
||||
artifact.path == identity.0
|
||||
&& artifact.source.hash() == identity.1
|
||||
&& artifact.operation.as_ref() == Some(operation)
|
||||
}
|
||||
_ => false,
|
||||
})
|
||||
}
|
||||
|
||||
/// The artifacts already collected for the run, read once: a file that
|
||||
/// is unchanged since it was collected is not collected again.
|
||||
async fn collected_artifacts(&self) -> Result<&Mutex<HashSet<ArtifactIdentity>>, HookError> {
|
||||
|
|
@ -1024,7 +1091,7 @@ impl FabroHooks {
|
|||
.into_iter()
|
||||
.filter_map(|record| match record.record {
|
||||
PlatformRecord::ArtifactCollected(artifact) => {
|
||||
Some((artifact.path, artifact.digest))
|
||||
Some((artifact.path, artifact.source.hash()))
|
||||
}
|
||||
_ => None,
|
||||
})
|
||||
|
|
@ -1351,7 +1418,332 @@ impl ExecutionHooks for FabroHooks {
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
|
||||
use fabro_store::{ArtifactStore, StoredPlatformRecord};
|
||||
use object_store::memory::InMemory;
|
||||
use tokio::sync::Barrier;
|
||||
|
||||
use super::*;
|
||||
use crate::artifacts::StoreArtifactWriter;
|
||||
use crate::test_support::MemoryPlatformRecords;
|
||||
|
||||
struct NoHooks;
|
||||
#[async_trait::async_trait]
|
||||
impl ExecutionHooks for NoHooks {}
|
||||
|
||||
struct FlakyRecords {
|
||||
records: MemoryPlatformRecords,
|
||||
fail: AtomicBool,
|
||||
/// Commit the failing append anyway, as a lost response does.
|
||||
commit: bool,
|
||||
}
|
||||
|
||||
#[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) {
|
||||
if self.commit {
|
||||
self.records.append(run_id, record, position).await?;
|
||||
}
|
||||
return Err(PlatformRecordError::Store(fabro_store::Error::Io(
|
||||
std::io::Error::other("test append unavailable"),
|
||||
)));
|
||||
}
|
||||
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,
|
||||
run_id: &RunId,
|
||||
digest: &BlobHash,
|
||||
bytes: Bytes,
|
||||
) -> Result<(), ArtifactWriteError> {
|
||||
if self.fail.swap(false, Ordering::SeqCst) {
|
||||
return Err(ArtifactWriteError::Store(fabro_store::Error::Io(
|
||||
std::io::Error::other("test store unavailable"),
|
||||
)));
|
||||
}
|
||||
self.writer.write(run_id, digest, bytes).await
|
||||
}
|
||||
}
|
||||
|
||||
/// Holds every upload until the test has started all of them.
|
||||
struct GatedWriter {
|
||||
writer: StoreArtifactWriter,
|
||||
gate: Barrier,
|
||||
uploads: AtomicUsize,
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ArtifactWriter for GatedWriter {
|
||||
async fn write(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
digest: &BlobHash,
|
||||
bytes: Bytes,
|
||||
) -> Result<(), ArtifactWriteError> {
|
||||
self.uploads.fetch_add(1, Ordering::SeqCst);
|
||||
let write = self.writer.write(run_id, digest, bytes);
|
||||
// A second capture of the identity must wait on the first, not
|
||||
// reach this point: time the gate out rather than hang.
|
||||
let _ = time::timeout(Duration::from_millis(200), self.gate.wait()).await;
|
||||
write.await
|
||||
}
|
||||
}
|
||||
|
||||
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),
|
||||
commit: false,
|
||||
});
|
||||
let writer = Arc::new(FlakyWriter {
|
||||
writer: StoreArtifactWriter::new(store.clone()),
|
||||
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 = Bytes::from_static(b"binary\0payload");
|
||||
let hash = BlobHash::new(&bytes);
|
||||
let error = hooks
|
||||
.store_artifact(key, &path, bytes.clone())
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(error.render().contains(if fail_upload {
|
||||
"test store unavailable"
|
||||
} else {
|
||||
"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.clone())
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert!(
|
||||
!hooks
|
||||
.store_artifact(key, &path, bytes.clone())
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
let resumed = artifact_hooks(root.path(), run, records.clone(), writer);
|
||||
assert!(
|
||||
!resumed
|
||||
.store_artifact(key, &path, bytes.clone())
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert!(
|
||||
resumed
|
||||
.store_artifact(key, &path, Bytes::from_static(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 = Bytes::from_static(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,
|
||||
operation: None,
|
||||
}),
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let resumed = artifact_hooks(
|
||||
root.path(),
|
||||
run,
|
||||
records.clone(),
|
||||
Arc::new(StoreArtifactWriter::new(store.clone())),
|
||||
);
|
||||
assert!(
|
||||
!resumed
|
||||
.store_artifact(key, &path, bytes.clone())
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert!(store.get_capture(&run, &hash).await.unwrap().is_none());
|
||||
assert!(
|
||||
resumed
|
||||
.store_artifact(key, &path, Bytes::from_static(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())),
|
||||
);
|
||||
assert!(
|
||||
fork_hooks
|
||||
.store_artifact(key, &path, bytes.clone())
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(
|
||||
store.get_capture(&fork, &hash).await.unwrap().unwrap(),
|
||||
bytes
|
||||
);
|
||||
assert!(store.get_capture(&run, &hash).await.unwrap().is_none());
|
||||
assert_eq!(records.records(&fork).len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn concurrent_captures_of_one_file_upload_and_record_it_once() {
|
||||
let root = tempfile::tempdir().unwrap();
|
||||
let run = RunId::new();
|
||||
let records = Arc::new(MemoryPlatformRecords::new());
|
||||
let store = ArtifactStore::new(Arc::new(InMemory::new()), "artifacts");
|
||||
let writer = Arc::new(GatedWriter {
|
||||
writer: StoreArtifactWriter::new(store.clone()),
|
||||
gate: Barrier::new(2),
|
||||
uploads: AtomicUsize::new(0),
|
||||
});
|
||||
let hooks = artifact_hooks(root.path(), run, records.clone(), writer.clone());
|
||||
let key = |firing| CheckpointKey {
|
||||
execution: 0,
|
||||
firing,
|
||||
attempt: 1,
|
||||
};
|
||||
let bytes = Bytes::from_static(b"same payload");
|
||||
let (first, second) = tokio::join!(
|
||||
hooks.store_artifact(key(1), "assets/report.bin", bytes.clone()),
|
||||
hooks.store_artifact(key(2), "assets/report.bin", bytes.clone()),
|
||||
);
|
||||
let mut outcomes = [first.unwrap(), second.unwrap()];
|
||||
outcomes.sort_unstable();
|
||||
assert_eq!(outcomes, [false, true]);
|
||||
assert_eq!(writer.uploads.load(Ordering::SeqCst), 1);
|
||||
assert_eq!(records.records(&run).len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_committed_append_whose_response_was_lost_is_recorded_once() {
|
||||
let root = tempfile::tempdir().unwrap();
|
||||
let run = RunId::new();
|
||||
let store = ArtifactStore::new(Arc::new(InMemory::new()), "artifacts");
|
||||
let records = Arc::new(FlakyRecords {
|
||||
records: MemoryPlatformRecords::new(),
|
||||
fail: AtomicBool::new(true),
|
||||
commit: true,
|
||||
});
|
||||
let hooks = artifact_hooks(
|
||||
root.path(),
|
||||
run,
|
||||
records.clone(),
|
||||
Arc::new(StoreArtifactWriter::new(store)),
|
||||
);
|
||||
let key = CheckpointKey {
|
||||
execution: 0,
|
||||
firing: 1,
|
||||
attempt: 1,
|
||||
};
|
||||
let bytes = Bytes::from_static(b"payload");
|
||||
assert!(
|
||||
hooks
|
||||
.store_artifact(key, "assets/report.bin", bytes.clone())
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert!(
|
||||
!hooks
|
||||
.store_artifact(key, "assets/report.bin", bytes)
|
||||
.await
|
||||
.unwrap()
|
||||
);
|
||||
assert_eq!(records.records.records(&run).len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_selection_keeps_the_smallest_files_within_the_budgets() {
|
||||
|
|
|
|||
|
|
@ -64,6 +64,7 @@
|
|||
//! workspace `Cargo.toml` for how they are tracked.
|
||||
|
||||
pub mod admission;
|
||||
pub mod artifacts;
|
||||
pub mod blobs;
|
||||
pub mod check;
|
||||
pub mod checkpoint;
|
||||
|
|
|
|||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -19,6 +19,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::{
|
||||
|
|
@ -32,9 +33,10 @@ use fabro_petri::providers::{DaytonaCredentials, SandboxProviderConfig};
|
|||
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;
|
||||
|
|
@ -80,39 +82,51 @@ fn admit(workflow: &str, settings: &str) -> AdmittedGraphs {
|
|||
|
||||
/// 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())),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -431,10 +445,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() {
|
||||
let mut harness = Harness::new();
|
||||
|
|
@ -465,14 +479,26 @@ 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_ne!(artifacts[0].digest, artifacts[1].digest);
|
||||
assert_ne!(artifacts[0].source.hash(), artifacts[1].source.hash());
|
||||
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.
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
@ -10,6 +10,7 @@ use object_store::buffered::BufWriter;
|
|||
use object_store::path::Path as ObjectPath;
|
||||
use percent_encoding::{AsciiSet, NON_ALPHANUMERIC, percent_decode_str, utf8_percent_encode};
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use tokio::task;
|
||||
|
||||
use crate::{Error, Result, StageId};
|
||||
|
||||
|
|
@ -74,13 +75,79 @@ 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?;
|
||||
self.put_at(
|
||||
&self.artifact_path(run_id, key)?,
|
||||
Bytes::copy_from_slice(data),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Publish a complete captured file under its content digest within the
|
||||
/// run, after checking that `data` hashes to `hash`. Repeating a put of
|
||||
/// the same content is safe; metadata is recorded separately only after
|
||||
/// this operation succeeds.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// [`Error::CaptureDigestMismatch`] when `data` does not hash to `hash`;
|
||||
/// nothing is written then.
|
||||
pub async fn put_capture(&self, run_id: &RunId, hash: &BlobHash, data: Bytes) -> Result<()> {
|
||||
let path = self.capture_path(run_id, hash)?;
|
||||
// Up to the capture size limit of bytes: hash off the async workers.
|
||||
let (actual, data) = task::spawn_blocking(move || (BlobHash::new(&data), data))
|
||||
.await
|
||||
.map_err(|err| Error::Other(format!("capture digest task failed: {err}")))?;
|
||||
if actual != *hash {
|
||||
return Err(Error::CaptureDigestMismatch { expected: *hash });
|
||||
}
|
||||
self.put_at(&path, data).await
|
||||
}
|
||||
|
||||
/// Remove one captured file's content. Removing absent content succeeds.
|
||||
pub async fn delete_capture(&self, run_id: &RunId, hash: &BlobHash) -> Result<()> {
|
||||
match self
|
||||
.object_store
|
||||
.delete(&self.capture_path(run_id, hash)?)
|
||||
.await
|
||||
{
|
||||
Ok(()) | Err(object_store::Error::NotFound { .. }) => Ok(()),
|
||||
Err(err) => Err(err.into()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Read run-owned content named by a capture record, without consulting
|
||||
/// 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
|
||||
}
|
||||
|
||||
/// Everything under a run's `captures/` belongs to recorded captures,
|
||||
/// never to historical stage keys, whatever digest namespace it uses.
|
||||
fn captures_root(&self, run_id: &RunId) -> Result<ObjectPath> {
|
||||
Ok(self.run_prefix(run_id)?.child("captures"))
|
||||
}
|
||||
|
||||
fn capture_prefix(&self, run_id: &RunId) -> Result<ObjectPath> {
|
||||
Ok(self.captures_root(run_id)?.child("sha256"))
|
||||
}
|
||||
|
||||
fn capture_path(&self, run_id: &RunId, hash: &BlobHash) -> Result<ObjectPath> {
|
||||
Ok(self.capture_prefix(run_id)?.child(hash.to_string()))
|
||||
}
|
||||
|
||||
async fn put_at(&self, path: &ObjectPath, data: Bytes) -> Result<()> {
|
||||
self.object_store.put(path, 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 +182,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 +214,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.captures_root(run_id)?;
|
||||
let mut stream = self.object_store.list(Some(&prefix));
|
||||
let mut artifacts = Vec::new();
|
||||
while let Some(meta) = stream.next().await.transpose()? {
|
||||
// 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 +495,108 @@ 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::from_static(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_is_never_listed_as_stage_artifacts() {
|
||||
let objects: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
|
||||
let store = ArtifactStore::new(objects.clone(), "artifacts");
|
||||
let run = RunId::new();
|
||||
let legacy = ArtifactKey::new(StageId::new("build", 1), 1, "report.txt");
|
||||
store.put(&run, &legacy, b"legacy").await.unwrap();
|
||||
for location in [
|
||||
format!("artifacts/{run}/captures/elsewhere/other"),
|
||||
format!("artifacts/{run}/captures/stray"),
|
||||
] {
|
||||
objects
|
||||
.put(
|
||||
&ObjectPath::from(location),
|
||||
Bytes::from_static(b"other").into(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
let listed = store.list_for_run(&run).await.unwrap();
|
||||
assert_eq!(listed.len(), 1);
|
||||
assert_eq!(listed[0].filename, "report.txt");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn put_capture_rejects_content_that_does_not_match_its_digest() {
|
||||
let objects: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
|
||||
let store = ArtifactStore::new(objects.clone(), "artifacts");
|
||||
let run = RunId::new();
|
||||
let hash = fabro_types::BlobHash::new(b"expected");
|
||||
let err = store
|
||||
.put_capture(&run, &hash, Bytes::from_static(b"other"))
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(err, Error::CaptureDigestMismatch { expected } if expected == hash));
|
||||
assert!(store.get_capture(&run, &hash).await.unwrap().is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_capture_removes_content_and_tolerates_absence() {
|
||||
let store = test_store();
|
||||
let run = RunId::new();
|
||||
let hash = fabro_types::BlobHash::new(b"content");
|
||||
store
|
||||
.put_capture(&run, &hash, Bytes::from_static(b"content"))
|
||||
.await
|
||||
.unwrap();
|
||||
store.delete_capture(&run, &hash).await.unwrap();
|
||||
assert!(store.get_capture(&run, &hash).await.unwrap().is_none());
|
||||
store.delete_capture(&run, &hash).await.unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn write_metadata_persists_store_marker() {
|
||||
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
|
||||
|
|
|
|||
|
|
@ -26,6 +26,8 @@ pub enum Error {
|
|||
BlobHashConflict { blob_hash: BlobHash },
|
||||
#[error("stored blob data does not match requested hash {blob_hash}")]
|
||||
BlobIntegrity { blob_hash: BlobHash },
|
||||
#[error("captured content does not match its digest {expected}")]
|
||||
CaptureDigestMismatch { expected: BlobHash },
|
||||
#[error("I/O error: {0}")]
|
||||
Io(#[from] std::io::Error),
|
||||
#[error("Invalid event payload: {0}")]
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -456,16 +457,61 @@ 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. Its hash is the SHA-256 of the
|
||||
/// bytes: with `path`, the identity a later capture of the same
|
||||
/// unchanged file is matched by. The wire also carries it as `digest`.
|
||||
#[serde(flatten, with = "source_with_digest")]
|
||||
pub source: ArtifactSource,
|
||||
pub bytes: u64,
|
||||
/// The SHA-256 of the bytes as lowercase hex: with `path`, the identity
|
||||
/// a later capture of the same unchanged file is matched by.
|
||||
pub digest: String,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub operation: Option<OperationKey>,
|
||||
}
|
||||
|
||||
/// The record's `digest` key repeats its source's hash as lowercase hex, as
|
||||
/// every earlier record wrote it. It is derived on write and checked on read.
|
||||
mod source_with_digest {
|
||||
use fabro_types::ArtifactSource;
|
||||
use serde::de::Error as _;
|
||||
use serde::{Deserialize, Deserializer, Serialize, Serializer};
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct Written<'a> {
|
||||
#[serde(flatten)]
|
||||
source: &'a ArtifactSource,
|
||||
digest: String,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct Read {
|
||||
#[serde(flatten)]
|
||||
source: ArtifactSource,
|
||||
digest: String,
|
||||
}
|
||||
|
||||
pub(super) fn serialize<S: Serializer>(
|
||||
source: &ArtifactSource,
|
||||
serializer: S,
|
||||
) -> Result<S::Ok, S::Error> {
|
||||
Written {
|
||||
source,
|
||||
digest: source.hash().to_string(),
|
||||
}
|
||||
.serialize(serializer)
|
||||
}
|
||||
|
||||
pub(super) fn deserialize<'de, D: Deserializer<'de>>(
|
||||
deserializer: D,
|
||||
) -> Result<ArtifactSource, D::Error> {
|
||||
let read = Read::deserialize(deserializer)?;
|
||||
if read.digest != read.source.hash().to_string() {
|
||||
return Err(D::Error::custom(
|
||||
"artifact checksum does not match its payload source",
|
||||
));
|
||||
}
|
||||
Ok(read.source)
|
||||
}
|
||||
}
|
||||
|
||||
/// 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 +790,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,9 +893,8 @@ 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 {
|
||||
execution: 0,
|
||||
decision: DecisionRef::AttemptStart {
|
||||
|
|
|
|||
|
|
@ -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", &[]),
|
||||
|
|
|
|||
|
|
@ -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::{
|
||||
|
|
|
|||
67
lib/foundation/fabro-api/tests/artifact_round_trip.rs
Normal file
67
lib/foundation/fabro-api/tests/artifact_round_trip.rs
Normal 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());
|
||||
}
|
||||
}
|
||||
|
|
@ -1842,6 +1842,28 @@ 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: Bytes,
|
||||
) -> Result<()> {
|
||||
let data = &data;
|
||||
self.send_api(|client| async move {
|
||||
client
|
||||
.write_run_artifact_content()
|
||||
.id(run_id.to_string())
|
||||
.digest(*digest)
|
||||
// A `Bytes` clone shares the buffer, including on a retry.
|
||||
.body(data.clone())
|
||||
.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 {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
55
lib/foundation/fabro-types/src/artifact_source.rs
Normal file
55
lib/foundation/fabro-types/src/artifact_source.rs
Normal 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;
|
||||
|
|
@ -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,17 @@ pub struct ServerSettings {
|
|||
impl ServerSettings {
|
||||
#[must_use]
|
||||
pub fn with_storage_override(mut self, path: &Path) -> Self {
|
||||
// The local artifact store always lives under the storage directory.
|
||||
// A configured `local.root` has never moved it, and existing objects
|
||||
// are only found there.
|
||||
if let ObjectStoreSettings::Local { root } = &mut self.server.artifacts.store {
|
||||
*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,
|
||||
|
|
|
|||
|
|
@ -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::{
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
@ -59,7 +59,9 @@ pub struct RunProjection {
|
|||
pub pending_interviews: BTreeMap<String, PendingInterviewRecord>,
|
||||
/// The files collected from the run's workspaces under
|
||||
/// `[run.artifacts] include`, one entry per capture, in the order they
|
||||
/// were recorded. The bytes are in the blob table under `blob`.
|
||||
/// were recorded. Each entry's `source` says where its bytes are: the
|
||||
/// configured artifact store for `object`, the SQLite blob table for
|
||||
/// earlier captures under `blob`.
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub artifacts: Vec<RunArtifact>,
|
||||
stages: HashMap<StageId, StageProjection>,
|
||||
|
|
@ -75,8 +77,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 +925,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;
|
||||
|
|
|
|||
|
|
@ -217,6 +217,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 {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
29
lib/packages/fabro-api-client/src/models/artifact-source.ts
generated
Normal file
29
lib/packages/fabro-api-client/src/models/artifact-source.ts
generated
Normal 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;
|
||||
}
|
||||
|
|
@ -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';
|
||||
|
|
|
|||
36
lib/packages/fabro-api-client/src/models/run-artifact.ts
generated
Normal file
36
lib/packages/fabro-api-client/src/models/run-artifact.ts
generated
Normal 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;
|
||||
}
|
||||
|
|
@ -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.
|
||||
*/
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue