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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 20:12:20 -04:00
parent 25d47ebcd3
commit 2ca1e1cda0
No known key found for this signature in database
8 changed files with 720 additions and 1 deletions

5
Cargo.lock generated
View file

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

View file

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

View file

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

View file

@ -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<HashMap<(RunId, OwnerId), Arc<dyn RunLogs>>>,
}
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<Arc<dyn RunLogs>, 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<Arc<dyn RunLogs>, 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<Arc<dyn RunLogs>, 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::<Vec<_>>()
};
if !dropped.is_empty() {
debug!(
run_id = %run_id,
handles = dropped.len(),
"Petri run handles released at worker exit"
);
}
drop(dropped);
}
}
fn lock<T>(mutex: &Mutex<T>) -> 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<Notify>,
}
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<StartedWorker> {
self.running.store(true, Ordering::SeqCst);
let exit = Arc::clone(&self.exit);
let stderr: Pin<Box<dyn AsyncRead + Send + 'static>> = 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<dyn WorkerRuntime>)
.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);
}
}

View file

@ -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<Arc<AppState>> for RequireWorkerRunScoped {
}
}
impl FromRequestParts<Arc<AppState>> for RequireWorkerRunSegment {
type Rejection = Response;
async fn from_request_parts(
parts: &mut Parts,
state: &Arc<AppState>,
) -> Result<Self, Self::Rejection> {
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<Arc<AppState>> for RequireRunManagementTarget {
type Rejection = Response;

View file

@ -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<dyn WorkerControlBus>,
pub(crate) worker_runtime: Arc<dyn WorkerRuntime>,
/// 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<AuthCodeStore> {
&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<Arc<AppS
.context("load mcp servers")?,
);
let variables = Arc::new(VariableStore::new(db_pool.clone()));
let petri_runs = PetriRuns::new(db_pool.clone());
let session_records = Arc::new(RunSessionRecordStore::new(db_pool.clone()));
let secret_store = Arc::new(SecretStore::new(db_pool));
let vault = preloaded_vault;
@ -2583,6 +2601,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
max_concurrent_runs,
worker_control_bus,
worker_runtime,
petri_runs,
scheduler_notify: Notify::new(),
automation_scheduler_notify: Notify::new(),
pull_request_scheduler_notify: Notify::new(),
@ -4455,6 +4474,9 @@ async fn execute_run_subprocess(state: Arc<AppState>, 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 {

View file

@ -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<Arc<AppState>> {
.merge(lifecycle::routes())
.merge(steer::routes())
.merge(pair::routes())
.merge(petri::routes())
.merge(graph::manifest_routes())
.merge(graph::run_routes())
.merge(models::routes())

View file

@ -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<Arc<AppState>> {
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<Arc<AppState>>,
Json(request): Json<PetriOpenRequest>,
) -> 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<Arc<AppState>>,
Json(request): Json<PetriReleaseRequest>,
) -> 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<Arc<AppState>>,
) -> 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::<Result<Vec<_>, _>>()
{
Ok(records) => Json(PetriRecordList { records }).into_response(),
Err(err) => err.into_response(),
}
}
async fn append_records(
RequireWorkerRunSegment(id, log): RequireWorkerRunSegment,
State(state): State<Arc<AppState>>,
Json(request): Json<PetriAppendRequest>,
) -> Response {
let Some(log) = parse_log_id(&log) else {
return unknown_log(&log);
};
let records = match request
.records
.into_iter()
.map(stored_record)
.collect::<Result<Vec<_>, _>>()
{
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<Arc<AppState>>,
Query(query): Query<OwnerQuery>,
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::<BlobHash>() {
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<Arc<AppState>>,
) -> Response {
let digest = match blob_hash
.parse::<BlobHash>()
.map(|hash| hash.to_string().parse::<Digest>())
{
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 <n>`."
))
.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<Record, ApiError> {
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<PetriRecord, ApiError> {
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<const N: usize>(members: [(&str, Value); N]) -> Map<String, Value> {
members
.into_iter()
.map(|(name, value)| (name.to_string(), value))
.collect()
}