From 2ca1e1cda049847ed8821b781bd0da6418bca10e Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 20:12:20 -0400 Subject: [PATCH] Serve the Petri run store to workers The server answers the `/api/v1/runs/{id}/petri/*` endpoints from one `SqliteRunStore` over its pool. `PetriRuns` in `AppState` keeps the writer handle each worker opened, keyed by the run and the worker's owner id, so the lease semantics stay the store's: the handle drops on the worker's `release`, and every handle of a run drops when the server observes the run's worker exit, in the subprocess wait path. Never by timeout. A write from an owner with no held handle reopens only when the lease row still names that owner, so a server restart or a lost open reply recovers, and an owner the lease moved away from gets `petri_stale_owner`. Every endpoint is worker-scoped through the existing worker auth; a new `RequireWorkerRunSegment` extractor covers the two-segment routes. Store errors answer with a machine-readable code, the leased owner and the conflict position under `meta`, and a backend failure's cause goes to the server log rather than the worker. A test drives a held worker through the scheduler, opens the run over the API with its token, ends the worker, and sees the lease end. Co-Authored-By: Claude Fable 5.1 --- Cargo.lock | 5 + lib/apps/fabro-server/Cargo.toml | 2 + lib/apps/fabro-server/src/lib.rs | 1 + lib/apps/fabro-server/src/petri_runs.rs | 363 ++++++++++++++++++ .../fabro-server/src/principal_middleware.rs | 20 + lib/apps/fabro-server/src/server.rs | 24 +- .../fabro-server/src/server/handler/mod.rs | 2 + .../fabro-server/src/server/handler/petri.rs | 304 +++++++++++++++ 8 files changed, 720 insertions(+), 1 deletion(-) create mode 100644 lib/apps/fabro-server/src/petri_runs.rs create mode 100644 lib/apps/fabro-server/src/server/handler/petri.rs diff --git a/Cargo.lock b/Cargo.lock index 70342d49b..e885f192f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2886,8 +2886,12 @@ dependencies = [ name = "fabro-petri" version = "0.357.0-nightly.0" dependencies = [ + "anyhow", "async-trait", + "fabro-api", + "fabro-client", "fabro-db", + "fabro-http", "fabro-store", "fabro-types", "petri-attractor-steps", @@ -3005,6 +3009,7 @@ dependencies = [ "fabro-macros", "fabro-manifest", "fabro-mcp-store", + "fabro-petri", "fabro-proc", "fabro-redact", "fabro-sandbox", diff --git a/lib/apps/fabro-server/Cargo.toml b/lib/apps/fabro-server/Cargo.toml index 811dec41f..9135bc4c3 100644 --- a/lib/apps/fabro-server/Cargo.toml +++ b/lib/apps/fabro-server/Cargo.toml @@ -42,6 +42,7 @@ pebble-coding-agent.workspace = true fabro-llm = { path = "../../components/fabro-llm" } fabro-manifest = { path = "../../components/fabro-manifest" } fabro-mcp-store = { path = "../../components/fabro-mcp-store" } +fabro-petri = { path = "../../components/fabro-petri" } fabro-proc = { path = "../../foundation/fabro-proc" } fabro-template = { path = "../../foundation/fabro-template" } fabro-tool = { path = "../../components/fabro-tool" } @@ -115,6 +116,7 @@ chrono = { workspace = true } [dev-dependencies] fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"] } +fabro-petri = { path = "../../components/fabro-petri", features = ["test-support"] } fabro-llm = { path = "../../components/fabro-llm", features = ["test-support"] } git2.workspace = true tokio = { workspace = true, features = ["test-util", "macros"] } diff --git a/lib/apps/fabro-server/src/lib.rs b/lib/apps/fabro-server/src/lib.rs index 5b9ee5447..4d1be5fc3 100644 --- a/lib/apps/fabro-server/src/lib.rs +++ b/lib/apps/fabro-server/src/lib.rs @@ -32,6 +32,7 @@ mod interp; pub mod jwt_auth; pub mod manifest_validation; mod migrations; +mod petri_runs; mod principal_middleware; mod request_id; mod run_compiler; diff --git a/lib/apps/fabro-server/src/petri_runs.rs b/lib/apps/fabro-server/src/petri_runs.rs new file mode 100644 index 000000000..a3b22ae0a --- /dev/null +++ b/lib/apps/fabro-server/src/petri_runs.rs @@ -0,0 +1,363 @@ +//! The Petri runs the server holds open for its workers. +//! +//! A worker reaches its run's Petri records over the API +//! (`/api/v1/runs/{id}/petri/*`, `server::handler::petri`), and the server +//! answers from one `SqliteRunStore` over its pool. The store's lease is +//! held by a handle, so the server keeps the handle a worker opened, keyed +//! by the run and the worker's owner id, for as long as the worker's lease +//! should last: until the worker releases it, or until the server observes +//! the worker exit. That is the integration plan's rule for a lease: it ends +//! when the handle drops, when the server observes the worker exit, or by +//! operator release, never by timeout. +//! +//! The Petri run key of a Fabro run is the run id's text, as the plan sets +//! `RunOptions::run_key`. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; + +use fabro_db::DbPool; +use fabro_petri::SqliteRunStore; +use fabro_petri::petri::{Access, OwnerId, RunKey, RunLogs, RunStore as _, StoreError}; +use fabro_types::RunId; +use tracing::debug; + +pub(crate) struct PetriRuns { + store: SqliteRunStore, + /// The writer handle each worker holds open, by run and owner. + handles: Mutex>>, +} + +impl PetriRuns { + pub(crate) fn new(pool: DbPool) -> Self { + Self { + store: SqliteRunStore::new(pool), + handles: Mutex::default(), + } + } + + /// The Petri run key of a Fabro run. + pub(crate) fn key(run_id: &RunId) -> RunKey { + RunKey::new(run_id.to_string()) + } + + /// The store itself, for an operator's release and for inspection. + #[cfg(any(test, feature = "test-support"))] + pub(crate) fn store(&self) -> &SqliteRunStore { + &self.store + } + + /// Open the run as a worker asked. A writer handle is kept for the + /// owner until [`release`](Self::release) or + /// [`worker_exited`](Self::worker_exited); a reader handle is not kept. + pub(crate) async fn open( + &self, + run_id: RunId, + access: Access, + ) -> Result, StoreError> { + let handle = self.store.open(&Self::key(&run_id), access.clone()).await?; + if let Some(owner) = access.owner() { + lock(&self.handles).insert((run_id, owner.clone()), Arc::clone(&handle)); + } + Ok(handle) + } + + /// The writer handle `owner` holds on the run: the one kept from its + /// open, or a reopen when the store's lease row still names the owner + /// (the server restarted, or the open's reply was lost). An owner the + /// lease no longer names gets `StaleOwner`. + pub(crate) async fn writer( + &self, + run_id: RunId, + owner: &OwnerId, + ) -> Result, StoreError> { + if let Some(handle) = lock(&self.handles).get(&(run_id, owner.clone())) { + return Ok(Arc::clone(handle)); + } + let holder = self.store.owner(&Self::key(&run_id)).await?; + if holder.as_ref() != Some(owner) { + return Err(StoreError::StaleOwner); + } + debug!(run_id = %run_id, owner = %owner, "Petri run handle reopened for its lease holder"); + match self + .open(run_id, Access::Write { + owner: owner.clone(), + }) + .await + { + Err(StoreError::Leased { .. }) => Err(StoreError::StaleOwner), + opened => opened, + } + } + + /// A reader handle on the run: no lease, never kept. + pub(crate) async fn reader(&self, run_id: RunId) -> Result, StoreError> { + self.store.open(&Self::key(&run_id), Access::Read).await + } + + /// Drop the handle `owner` holds on the run: the worker's own release. + /// The store ends the lease when this was the owner's last handle. + pub(crate) fn release(&self, run_id: RunId, owner: &OwnerId) { + let handle = lock(&self.handles).remove(&(run_id, owner.clone())); + debug!( + run_id = %run_id, + owner = %owner, + held = handle.is_some(), + "Petri run handle released by its worker" + ); + drop(handle); + } + + /// Drop every handle held on the run: what the server does when it + /// observes the run's worker exit, so a worker that died without + /// releasing does not keep the lease. + pub(crate) fn worker_exited(&self, run_id: RunId) { + let dropped = { + let mut handles = lock(&self.handles); + let owners: Vec<_> = handles + .keys() + .filter(|(held, _)| *held == run_id) + .cloned() + .collect(); + owners + .into_iter() + .filter_map(|slot| handles.remove(&slot)) + .collect::>() + }; + if !dropped.is_empty() { + debug!( + run_id = %run_id, + handles = dropped.len(), + "Petri run handles released at worker exit" + ); + } + drop(dropped); + } +} + +fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { + mutex.lock().unwrap_or_else(PoisonError::into_inner) +} + +#[cfg(test)] +mod tests { + use std::pin::Pin; + use std::sync::Arc; + use std::sync::atomic::{AtomicBool, Ordering}; + use std::time::Duration; + + use axum::body::{Body, to_bytes}; + use axum::http::{Request, StatusCode, header}; + use fabro_config::Storage; + use fabro_config::bind::Bind; + use fabro_config::daemon::ServerDaemon; + use fabro_petri::petri::RunStore as _; + use fabro_static::EnvVars; + use fabro_types::{RunId, WorkflowPath, WorkflowVersion}; + use serde_json::json; + use tokio::io::AsyncRead; + use tokio::sync::Notify; + use tokio::time; + use tower::ServiceExt as _; + + use super::*; + use crate::server::{AppState, spawn_scheduler}; + use crate::test_support::{ + TestAppStateBuilder, build_test_router, test_register_workflow_version, + }; + use crate::worker_runtime::{ + StartedWorker, WorkerExit, WorkerLaunchSpec, WorkerRef, WorkerRuntime, + }; + + const MINIMAL_DOT: &str = r#"digraph Test { + graph [goal="Test"] + start [shape=Mdiamond] + exit [shape=Msquare] + start -> exit +}"#; + + /// A worker runtime whose one worker runs until the test ends it, so + /// the test can act while the server waits on the worker. + #[derive(Default)] + struct HeldWorkerRuntime { + started: Notify, + running: AtomicBool, + exit: Arc, + } + + impl HeldWorkerRuntime { + async fn wait_for_start(&self) { + time::timeout(Duration::from_secs(10), self.started.notified()) + .await + .expect("the scheduler starts the worker"); + } + + fn end_worker(&self) { + self.running.store(false, Ordering::SeqCst); + self.exit.notify_one(); + } + } + + #[async_trait::async_trait] + impl WorkerRuntime for HeldWorkerRuntime { + async fn start(&self, _spec: WorkerLaunchSpec) -> anyhow::Result { + self.running.store(true, Ordering::SeqCst); + let exit = Arc::clone(&self.exit); + let stderr: Pin> = Box::pin(tokio::io::empty()); + let started = StartedWorker { + worker_ref: WorkerRef::Local { pid: u32::MAX }, + stderr, + wait: Box::pin(async move { + exit.notified().await; + Ok(WorkerExit { + success: false, + detail: "test worker ended without a terminal event".to_string(), + }) + }), + }; + self.started.notify_one(); + Ok(started) + } + + async fn request_stop(&self, _worker_ref: &WorkerRef) { + self.end_worker(); + } + + async fn force_stop(&self, _worker_ref: &WorkerRef) { + self.end_worker(); + } + + async fn is_alive(&self, _worker_ref: &WorkerRef) -> bool { + self.running.load(Ordering::SeqCst) + } + } + + /// The server record the worker launch spec reads. + fn write_test_server_record(state: &AppState) { + let runtime_directory = Storage::new(state.server_storage_dir()).runtime_directory(); + ServerDaemon::new( + std::process::id(), + Bind::Tcp( + "127.0.0.1:32276" + .parse() + .expect("the test bind address parses"), + ), + runtime_directory.log_path(), + ) + .write(&runtime_directory) + .expect("the test server record writes"); + } + + /// A run created and started through the API, as a client would. + async fn create_and_start_run(app: &axum::Router) -> RunId { + let path = WorkflowPath::new("workflow.fabro").expect("a workflow path"); + let version = WorkflowVersion::new( + path.clone(), + std::collections::BTreeMap::from([(path, MINIMAL_DOT.to_string())]), + std::collections::BTreeMap::new(), + ) + .expect("a workflow version"); + let version_id = test_register_workflow_version(app, &version, None).await; + let intent = json!({ + "workflow_version_id": version_id, + "target": { "kind": "none" }, + "args": {}, + }); + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/runs") + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(intent.to_string())) + .expect("the create request builds"), + ) + .await + .expect("the create request completes"); + assert_eq!(response.status(), StatusCode::CREATED); + let body = to_bytes(response.into_body(), usize::MAX) + .await + .expect("the create body reads"); + let body: serde_json::Value = + serde_json::from_slice(&body).expect("the create body is JSON"); + let run_id: RunId = body["id"] + .as_str() + .expect("the created run has an id") + .parse() + .expect("the run id parses"); + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/v1/runs/{run_id}/start")) + .body(Body::empty()) + .expect("the start request builds"), + ) + .await + .expect("the start request completes"); + assert_eq!(response.status(), StatusCode::OK); + run_id + } + + /// The lease a worker took over the API ends when the server observes + /// the worker exit, with no release from the worker itself. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn the_lease_ends_when_the_server_observes_the_worker_exit() { + let runtime = Arc::new(HeldWorkerRuntime::default()); + let state = TestAppStateBuilder::new() + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .worker_runtime(Arc::clone(&runtime) as Arc) + .build(); + write_test_server_record(&state); + let app = build_test_router(Arc::clone(&state)); + let run_id = create_and_start_run(&app).await; + spawn_scheduler(Arc::clone(&state)); + runtime.wait_for_start().await; + + // The worker opens its run over the API and never releases it. + let token = state.test_issue_worker_token(&run_id); + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/v1/runs/{run_id}/petri/open")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from( + json!({ "access": "create", "owner": "worker-1" }).to_string(), + )) + .expect("the open request builds"), + ) + .await + .expect("the open request completes"); + assert_eq!(response.status(), StatusCode::OK); + let key = PetriRuns::key(&run_id); + let store = state.petri_runs.store(); + assert_eq!( + store.owner(&key).await.expect("reads the lease"), + Some(OwnerId::new("worker-1")) + ); + + runtime.end_worker(); + + let mut holder = store.owner(&key).await.expect("reads the lease"); + for _ in 0..500 { + if holder.is_none() { + break; + } + time::sleep(Duration::from_millis(10)).await; + holder = store.owner(&key).await.expect("reads the lease"); + } + assert_eq!(holder, None, "the lease ended when the worker exited"); + let resumed = store + .open(&key, Access::Write { + owner: OwnerId::new("resumer"), + }) + .await + .expect("the next owner takes the run"); + drop(resumed); + } +} diff --git a/lib/apps/fabro-server/src/principal_middleware.rs b/lib/apps/fabro-server/src/principal_middleware.rs index 3d9bb0542..ca3b16bc2 100644 --- a/lib/apps/fabro-server/src/principal_middleware.rs +++ b/lib/apps/fabro-server/src/principal_middleware.rs @@ -60,6 +60,9 @@ pub(crate) struct RequiredRunManagementActor(pub(crate) Principal); pub(crate) struct RequiredRunToolActor(pub(crate) Principal); pub(crate) struct RequireRunScoped(pub(crate) RunId); pub(crate) struct RequireWorkerRunScoped(pub(crate) RunId); +/// A worker-scoped route with one more path segment after the run id, handed +/// back as its text for the handler to parse. +pub(crate) struct RequireWorkerRunSegment(pub(crate) RunId, pub(crate) String); pub(crate) struct RequireRunManagementTarget(pub(crate) RunId, pub(crate) Principal); pub(crate) struct RequireRunBlob(pub(crate) RunId, pub(crate) BlobHash); pub(crate) struct RequireRunStageScoped(pub(crate) RunId, pub(crate) String); @@ -265,6 +268,23 @@ impl FromRequestParts> for RequireWorkerRunScoped { } } +impl FromRequestParts> for RequireWorkerRunSegment { + type Rejection = Response; + + async fn from_request_parts( + parts: &mut Parts, + state: &Arc, + ) -> Result { + let Path((id, segment)): Path<(String, String)> = Path::from_request_parts(parts, state) + .await + .map_err(IntoResponse::into_response)?; + let run_id = parse_run_id_path(&id)?; + require_worker_for_run(&auth_slot_from_parts(parts), &run_id) + .map_err(IntoResponse::into_response)?; + Ok(Self(run_id, segment)) + } +} + impl FromRequestParts> for RequireRunManagementTarget { type Rejection = Response; diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 7c3a7c83e..5c3e24b6f 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -146,10 +146,11 @@ use crate::github_webhooks::{ WEBHOOK_ROUTE, WEBHOOK_SECRET_ENV, parse_event_metadata, verify_signature, }; use crate::jwt_auth::{self, AuthMode}; +use crate::petri_runs::PetriRuns; use crate::principal_middleware::{ AuthContextSlot, RequestAuth, RequestAuthContext, RequireRunBlob, RequireRunManagementTarget, RequireRunScoped, RequireRunStageScoped, RequireStageArtifact, RequireWorkerRunScoped, - RequiredUser, principal_middleware, + RequireWorkerRunSegment, RequiredUser, principal_middleware, }; use crate::request_id::{self, RequestId}; use crate::run_files::{FilesInFlight, new_files_in_flight}; @@ -1112,6 +1113,8 @@ pub struct AppState { max_concurrent_runs: usize, pub(crate) worker_control_bus: Arc, pub(crate) worker_runtime: Arc, + /// The Petri runs held open for workers over the API. + pub(crate) petri_runs: PetriRuns, scheduler_notify: Notify, automation_scheduler_notify: Notify, pull_request_scheduler_notify: Notify, @@ -1172,6 +1175,20 @@ impl AppState { pub fn test_auth_code_store(&self) -> &Arc { &self.stores.auth_codes } + + /// The Petri run store the worker endpoints answer from, so a test can + /// release a lease as an operator would and read who holds one. + #[must_use] + pub fn test_petri_run_store(&self) -> &fabro_petri::SqliteRunStore { + self.petri_runs.store() + } + + /// A worker token for `run_id` with the plain `run:worker` scope, as the + /// server mints for the worker it launches. + pub fn test_issue_worker_token(&self, run_id: &RunId) -> String { + issue_worker_token_with_scopes(&self.worker_tokens, run_id, WorkerScopeSet::run_worker()) + .expect("a test worker token signs") + } } impl AppState { @@ -2467,6 +2484,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result anyhow::Result, run_id: RunId) { return; } + // The worker is gone: whatever Petri run handles it held open over the + // API drop here, so its lease never outlives it. + state.petri_runs.worker_exited(run_id); append_worker_exit_failure(&run_store, run_id, &worker_exit).await; let final_state = match run_store.state().await { diff --git a/lib/apps/fabro-server/src/server/handler/mod.rs b/lib/apps/fabro-server/src/server/handler/mod.rs index c5d1c2018..d22bb8d29 100644 --- a/lib/apps/fabro-server/src/server/handler/mod.rs +++ b/lib/apps/fabro-server/src/server/handler/mod.rs @@ -18,6 +18,7 @@ mod llm_sse; mod mcp_servers; mod models; mod pair; +mod petri; pub(in crate::server) mod pull_requests; pub(in crate::server) mod runs; mod sandbox; @@ -219,6 +220,7 @@ pub(super) fn real_routes() -> Router> { .merge(lifecycle::routes()) .merge(steer::routes()) .merge(pair::routes()) + .merge(petri::routes()) .merge(graph::manifest_routes()) .merge(graph::run_routes()) .merge(models::routes()) diff --git a/lib/apps/fabro-server/src/server/handler/petri.rs b/lib/apps/fabro-server/src/server/handler/petri.rs new file mode 100644 index 000000000..cc6824147 --- /dev/null +++ b/lib/apps/fabro-server/src/server/handler/petri.rs @@ -0,0 +1,304 @@ +//! The Petri run store over HTTP: the endpoints a run's worker uses to reach +//! the run's Petri records (`fabro_petri::HttpRunStore` is the client). Every +//! endpoint is worker-scoped, and the server answers from the handles +//! `crate::petri_runs::PetriRuns` holds, so the lease and the `(log, seq)` +//! rule are the store's own. +//! +//! Each store error answers with a machine-readable `code`: +//! `petri_run_exists` and `petri_run_leased` (with the holder under +//! `meta.owner`) on `open`, `petri_run_not_found` wherever the run is +//! missing, `petri_stale_owner` and `petri_record_conflict` (with the +//! position under `meta.log` and `meta.seq`) on a write, `petri_read_only` +//! should a reader ever be asked to write, `petri_blob_not_found` on a blob +//! read, and `petri_store_failed` for the backend itself, whose cause goes to +//! the server log and not to the worker. + +use std::sync::Arc; + +use axum::extract::DefaultBodyLimit; +use axum::routing::{get, post}; +use fabro_api::types::{ + PetriAccess, PetriAppendRequest, PetriOpenRequest, PetriOpenResponse, PetriRecord, + PetriRecordList, PetriReleaseRequest, WriteBlobResponse, +}; +use fabro_petri::petri::{Access, Digest, OwnerId, Record, StoreError}; +use fabro_petri::run_store::{log_id_text, parse_log_id}; +use fabro_types::BlobHash; +use fabro_util::error::collect_chain; +use serde_json::{Map, Value, json}; + +use super::super::{ + ApiError, AppState, Bytes, IntoResponse, Json, Query, RequireWorkerRunScoped, + RequireWorkerRunSegment, Response, Router, RunId, State, StatusCode, octet_stream_response, +}; + +/// The largest batch of records one append may carry. Petri batches an +/// execution's records per step, and a step's output can be large. +const RECORD_BATCH_BODY_LIMIT: usize = 256 * 1024 * 1024; + +pub(super) fn routes() -> Router> { + Router::new() + .route("/runs/{id}/petri/open", post(open_run)) + .route("/runs/{id}/petri/release", post(release_run)) + .route( + "/runs/{id}/petri/logs/{log}/records", + get(list_records) + .post(append_records) + .layer(DefaultBodyLimit::max(RECORD_BATCH_BODY_LIMIT)), + ) + .route( + "/runs/{id}/petri/blobs", + post(write_blob).layer(DefaultBodyLimit::disable()), + ) + .route("/runs/{id}/petri/blobs/{blobHash}", get(read_blob)) +} + +#[derive(serde::Deserialize)] +struct OwnerQuery { + owner: String, +} + +async fn open_run( + RequireWorkerRunScoped(id): RequireWorkerRunScoped, + State(state): State>, + Json(request): Json, +) -> Response { + let access = match (request.access, request.owner) { + (PetriAccess::Create, Some(owner)) => Access::Create { + owner: OwnerId::new(owner), + }, + (PetriAccess::Write, Some(owner)) => Access::Write { + owner: OwnerId::new(owner), + }, + (PetriAccess::Read, _) => Access::Read, + (PetriAccess::Create | PetriAccess::Write, None) => { + return ApiError::bad_request("`owner` is required to open a Petri run for writing.") + .into_response(); + } + }; + match state.petri_runs.open(id, access).await { + Ok(handle) => Json(PetriOpenResponse { + locator: handle.locator(), + }) + .into_response(), + Err(err) => store_error_response(id, &err), + } +} + +async fn release_run( + RequireWorkerRunScoped(id): RequireWorkerRunScoped, + State(state): State>, + Json(request): Json, +) -> Response { + state.petri_runs.release(id, &OwnerId::new(request.owner)); + StatusCode::NO_CONTENT.into_response() +} + +async fn list_records( + RequireWorkerRunSegment(id, log): RequireWorkerRunSegment, + State(state): State>, +) -> Response { + let Some(log) = parse_log_id(&log) else { + return unknown_log(&log); + }; + let reader = match state.petri_runs.reader(id).await { + Ok(reader) => reader, + Err(err) => return store_error_response(id, &err), + }; + let records = match reader.read(&log).await { + Ok(records) => records, + Err(err) => return store_error_response(id, &err), + }; + match records + .into_iter() + .map(wire_record) + .collect::, _>>() + { + Ok(records) => Json(PetriRecordList { records }).into_response(), + Err(err) => err.into_response(), + } +} + +async fn append_records( + RequireWorkerRunSegment(id, log): RequireWorkerRunSegment, + State(state): State>, + Json(request): Json, +) -> Response { + let Some(log) = parse_log_id(&log) else { + return unknown_log(&log); + }; + let records = match request + .records + .into_iter() + .map(stored_record) + .collect::, _>>() + { + Ok(records) => records, + Err(err) => return err.into_response(), + }; + let writer = match state + .petri_runs + .writer(id, &OwnerId::new(request.owner)) + .await + { + Ok(writer) => writer, + Err(err) => return store_error_response(id, &err), + }; + match writer.append(&log, &records).await { + Ok(()) => StatusCode::NO_CONTENT.into_response(), + Err(err) => store_error_response(id, &err), + } +} + +async fn write_blob( + RequireWorkerRunScoped(id): RequireWorkerRunScoped, + State(state): State>, + Query(query): Query, + body: Bytes, +) -> Response { + let writer = match state + .petri_runs + .writer(id, &OwnerId::new(query.owner)) + .await + { + Ok(writer) => writer, + Err(err) => return store_error_response(id, &err), + }; + let digest = match writer.put_blob(&body).await { + Ok(digest) => digest, + Err(err) => return store_error_response(id, &err), + }; + match digest.to_hex().parse::() { + Ok(hash) => Json(WriteBlobResponse { hash }).into_response(), + Err(err) => ApiError::with_code( + StatusCode::INTERNAL_SERVER_ERROR, + format!("The stored blob's digest is not a blob hash: {err}"), + "petri_store_failed", + ) + .into_response(), + } +} + +async fn read_blob( + RequireWorkerRunSegment(id, blob_hash): RequireWorkerRunSegment, + State(state): State>, +) -> Response { + let digest = match blob_hash + .parse::() + .map(|hash| hash.to_string().parse::()) + { + Ok(Ok(digest)) => digest, + Ok(Err(err)) => return ApiError::bad_request(err.to_string()).into_response(), + Err(err) => return ApiError::bad_request(err.to_string()).into_response(), + }; + let reader = match state.petri_runs.reader(id).await { + Ok(reader) => reader, + Err(err) => return store_error_response(id, &err), + }; + match reader.get_blob(digest).await { + Ok(Some(bytes)) => octet_stream_response(Bytes::from(bytes)), + Ok(None) => ApiError::with_code( + StatusCode::NOT_FOUND, + "The run holds no blob with this digest.", + "petri_blob_not_found", + ) + .into_response(), + Err(err) => store_error_response(id, &err), + } +} + +fn unknown_log(log: &str) -> Response { + ApiError::bad_request(format!( + "`{log}` is not a Petri log: expected `coordinator`, `resources` or `execution `." + )) + .into_response() +} + +/// A wire record into the record the store keeps: its JSON, with `seq` and +/// `recorded_at` lifted from it, which must agree with the ones sent +/// beside it. +fn stored_record(wire: PetriRecord) -> Result { + let record = Record::from_value(Value::Object(wire.record)) + .map_err(|err| ApiError::bad_request(format!("Invalid Petri record: {err}")))?; + if record.seq != wire.seq || record.recorded_at != wire.recorded_at { + return Err(ApiError::bad_request(format!( + "Invalid Petri record: it carries seq {} and recorded_at {}, but was sent as seq {} \ + and recorded_at {}.", + record.seq, record.recorded_at, wire.seq, wire.recorded_at + ))); + } + Ok(record) +} + +/// A stored record as the wire carries it. A stored record is always a JSON +/// object; one that is not is the backend's fault. +fn wire_record(record: Record) -> Result { + match record.record { + Value::Object(map) => Ok(PetriRecord { + seq: record.seq, + recorded_at: record.recorded_at, + record: map, + }), + _ => Err(ApiError::with_code( + StatusCode::INTERNAL_SERVER_ERROR, + format!( + "The stored record at seq {} is not a JSON object.", + record.seq + ), + "petri_store_failed", + )), + } +} + +/// The store's answer as the worker's client maps it back: a status, a +/// code, and the members the code documents under `meta`. +fn store_error_response(run_id: RunId, err: &StoreError) -> Response { + let error = match err { + StoreError::Exists { .. } => { + ApiError::with_code(StatusCode::CONFLICT, err.to_string(), "petri_run_exists") + } + StoreError::NotFound { .. } => ApiError::with_code( + StatusCode::NOT_FOUND, + err.to_string(), + "petri_run_not_found", + ), + StoreError::Leased { owner, .. } => ApiError::with_code_and_meta( + StatusCode::CONFLICT, + err.to_string(), + "petri_run_leased", + members([("owner", json!(owner.as_str()))]), + ), + StoreError::StaleOwner => { + ApiError::with_code(StatusCode::CONFLICT, err.to_string(), "petri_stale_owner") + } + StoreError::ReadOnly => { + ApiError::with_code(StatusCode::CONFLICT, err.to_string(), "petri_read_only") + } + StoreError::Conflict { log, seq } => ApiError::with_code_and_meta( + StatusCode::CONFLICT, + err.to_string(), + "petri_record_conflict", + members([("log", json!(log_id_text(log))), ("seq", json!(seq))]), + ), + StoreError::Backend { .. } => { + tracing::error!( + run_id = %run_id, + error = %collect_chain(err).join(": "), + "Petri run store failed" + ); + ApiError::with_code( + StatusCode::INTERNAL_SERVER_ERROR, + "The Petri run store failed; see the server log.", + "petri_store_failed", + ) + } + }; + error.into_response() +} + +fn members(members: [(&str, Value); N]) -> Map { + members + .into_iter() + .map(|(name, value)| (name.to_string(), value)) + .collect() +}