From 0820252bcfe383c84688f44e8d1fdad12f5f509c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 20:12:20 -0400 Subject: [PATCH 1/4] Add the Petri run store endpoints to the API The worker shape of the integration plan (F1.3) needs a run's worker to reach the run's Petri records over the server's API. This adds the contract: six worker-scoped endpoints under `/api/v1/runs/{id}/petri/` (open, release, list and append records of one log, write and read a blob), their request and response schemas, and the generated Rust and TypeScript clients. A store error needs more than a code: `petri_run_leased` names the holding owner and `petri_record_conflict` names the refused position. `ErrorResponseEntry` gains an optional `meta` object for such code-specific members, `ApiError` can carry it, and the client's `ApiFailure` parses it beside the code so a caller can act on it. Records travel as `{seq, recorded_at, record}`, the store's own unit, with `seq` and `recorded_at` as `uint64`. The log path segment is the log id's text (`coordinator`, `resources`, `execution `), which the generated client percent-encodes. The blob write reuses `WriteBlobResponse`, since Petri's digest is Fabro's blob hash. Co-Authored-By: Claude Fable 5.1 --- docs/public/api-reference/fabro-api.yaml | 353 ++++++++++++ lib/apps/fabro-cli/src/commands/parent/mod.rs | 4 +- lib/apps/fabro-server/src/error.rs | 59 +- lib/foundation/fabro-client/src/client.rs | 150 ++++- lib/foundation/fabro-client/src/error.rs | 132 +++-- lib/foundation/fabro-client/src/lib.rs | 6 +- .../src/.openapi-generator/FILES | 7 + .../src/api/run-internals-api.ts | 516 ++++++++++++++++++ .../src/models/error-response-entry.ts | 4 + .../fabro-api-client/src/models/index.ts | 7 + .../src/models/petri-access.ts | 27 + .../src/models/petri-append-request.ts | 29 + .../src/models/petri-open-request.ts | 29 + .../src/models/petri-open-response.ts | 25 + .../src/models/petri-record-list.ts | 25 + .../src/models/petri-record.ts | 33 ++ .../src/models/petri-release-request.ts | 25 + 17 files changed, 1384 insertions(+), 47 deletions(-) create mode 100644 lib/packages/fabro-api-client/src/models/petri-access.ts create mode 100644 lib/packages/fabro-api-client/src/models/petri-append-request.ts create mode 100644 lib/packages/fabro-api-client/src/models/petri-open-request.ts create mode 100644 lib/packages/fabro-api-client/src/models/petri-open-response.ts create mode 100644 lib/packages/fabro-api-client/src/models/petri-record-list.ts create mode 100644 lib/packages/fabro-api-client/src/models/petri-record.ts create mode 100644 lib/packages/fabro-api-client/src/models/petri-release-request.ts diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index e620dbd55..25ba828cf 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -3195,6 +3195,229 @@ paths: schema: $ref: "#/components/schemas/ErrorResponse" + # ── Petri run store (worker) ────────────────────────────────────────── + # + # The Petri run store over HTTP: what a run's worker process uses to reach + # the run's Petri records in the server's database. Every endpoint is + # worker-scoped: the worker token's run must be the path's run. `id` is + # the Petri run key, which is the Fabro run id. Errors carry a + # machine-readable `code`; `petri_run_leased` carries the holding owner + # under `meta.owner`, and `petri_record_conflict` carries the refused + # position under `meta.log` and `meta.seq`. + + /api/v1/runs/{id}/petri/open: + post: + operationId: openPetriRun + tags: [Run Internals] + summary: Open Petri Run + description: | + Opens the run in the Petri run store for the worker. `create` inserts + the run and takes its writer lease for `owner`; `write` takes the lease + of an existing run; `read` takes no lease. The lease is idempotent per + owner: a retry by the owner that holds it gets the same lease. Another + live owner is refused with `petri_run_leased`. The lease ends when the + worker releases it, when the server observes the worker exit, or by + operator release, never by timeout. + parameters: + - $ref: "#/components/parameters/RunId" + requestBody: + required: true + content: + application/json: + schema: + $ref: "#/components/schemas/PetriOpenRequest" + responses: + "200": + description: Run opened + content: + application/json: + schema: + $ref: "#/components/schemas/PetriOpenResponse" + "404": + description: Run not in the store (`petri_run_not_found`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "409": + description: | + The run exists (`petri_run_exists`, on `create`) or another live + owner holds its lease (`petri_run_leased`, with the holder under + `meta.owner`). + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + + /api/v1/runs/{id}/petri/release: + post: + operationId: releasePetriRun + tags: [Run Internals] + summary: Release Petri Run + description: | + Ends the worker's writer lease on the run when `owner` still holds + it: what a worker sends when it drops its store handle. A lease that + already moved to another owner is left alone. + parameters: + - $ref: "#/components/parameters/RunId" + requestBody: + required: true + content: + application/json: + schema: + $ref: "#/components/schemas/PetriReleaseRequest" + responses: + "204": + description: Lease released, or not held by this owner + + /api/v1/runs/{id}/petri/logs/{log}/records: + get: + operationId: listPetriRecords + tags: [Run Internals] + summary: List Petri Records + description: Every record of one log of the run, in `seq` order, unchanged. + parameters: + - $ref: "#/components/parameters/RunId" + - $ref: "#/components/parameters/PetriLog" + responses: + "200": + description: The log's records + content: + application/json: + schema: + $ref: "#/components/schemas/PetriRecordList" + "404": + description: Run not in the store (`petri_run_not_found`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + post: + operationId: appendPetriRecords + tags: [Run Internals] + summary: Append Petri Records + description: | + Appends one batch of records to one log at the sequences they carry, + durably, in one transaction. A record equal to the one already stored + at its `seq` is accepted without a second append, so a batch whose + reply was lost is safe to resend. A different record at a taken `seq`, + or a `seq` past the log's end, is refused with `petri_record_conflict` + and the batch stores nothing. + parameters: + - $ref: "#/components/parameters/RunId" + - $ref: "#/components/parameters/PetriLog" + requestBody: + required: true + content: + application/json: + schema: + $ref: "#/components/schemas/PetriAppendRequest" + responses: + "204": + description: Records durable + "404": + description: Run not in the store (`petri_run_not_found`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "409": + description: | + A record conflicts with the log (`petri_record_conflict`, with the + position under `meta.log` and `meta.seq`), or `owner` no longer + holds the run's lease (`petri_stale_owner`). + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + + /api/v1/runs/{id}/petri/blobs: + post: + operationId: writePetriBlob + tags: [Run Internals] + summary: Write Petri Blob + description: | + Stores a blob by content for the run's owner and returns its SHA-256 + digest, the same content address Fabro's blob store uses. Idempotent + by construction. + parameters: + - $ref: "#/components/parameters/RunId" + - $ref: "#/components/parameters/PetriOwner" + requestBody: + required: true + content: + application/octet-stream: + schema: + type: string + format: binary + responses: + "200": + description: Blob stored + content: + application/json: + schema: + $ref: "#/components/schemas/WriteBlobResponse" + "404": + description: Run not in the store (`petri_run_not_found`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + "409": + description: "`owner` no longer holds the run's lease (`petri_stale_owner`)" + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + + /api/v1/runs/{id}/petri/blobs/{blobHash}: + get: + operationId: readPetriBlob + tags: [Run Internals] + summary: Read Petri Blob + description: The blob with this digest, if the store holds one. + parameters: + - $ref: "#/components/parameters/RunId" + - $ref: "#/components/parameters/BlobHash" + responses: + "200": + description: Blob contents + content: + application/octet-stream: + schema: + type: string + format: binary + "404": + description: Run not in the store (`petri_run_not_found`) or no such blob (`petri_blob_not_found`) + headers: + x-request-id: + $ref: "#/components/headers/XRequestId" + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + /api/v1/runs/{id}/stages/{stageId}/logs/output: get: operationId: getRunStageCommandLog @@ -5978,6 +6201,27 @@ components: $ref: "#/components/schemas/BlobHash" example: 2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824 + PetriLog: + name: log + in: path + required: true + description: >- + A Petri log of the run, as its id renders: `coordinator`, `resources`, + or `execution ` for execution `n`. The space is percent-encoded on + the wire. + schema: + type: string + example: execution 0 + + PetriOwner: + name: owner + in: query + required: true + description: The owner id the worker opened the run's writer lease with. + schema: + type: string + example: 18f3c2a9e1b4-42017-0-9f3a1c7e2b5d + ArtifactFilename: name: filename in: query @@ -10184,6 +10428,12 @@ components: type: string format: uuid description: Server-generated request identifier; matches the x-request-id response header. + meta: + type: object + additionalProperties: true + description: >- + Optional structured details specific to the error `code`, for + clients that act on them. Each code documents the members it sets. ErrorResponse: description: Standard error response containing one or more error entries. @@ -10647,6 +10897,109 @@ components: hash: $ref: "#/components/schemas/BlobHash" + PetriAccess: + description: >- + How a worker opens a Petri run. `create` inserts the run and takes its + writer lease, `write` takes the lease of an existing run, `read` takes + no lease. + type: string + enum: + - create + - write + - read + + PetriOpenRequest: + description: An open of a Petri run for a worker. + type: object + required: + - access + properties: + access: + $ref: "#/components/schemas/PetriAccess" + owner: + type: string + nullable: true + description: >- + The owner id the writer lease is taken for. Required for `create` + and `write`, absent for `read`. + example: 18f3c2a9e1b4-42017-0-9f3a1c7e2b5d + + PetriOpenResponse: + description: A Petri run opened for a worker. + type: object + required: + - locator + properties: + locator: + type: string + description: Where the run lives, for messages. + example: "sqlite database /var/lib/fabro/db/fabro.sqlite3, run `01JNQVR7M0EJ5GKAT2SC4ERS1Z`" + + PetriReleaseRequest: + description: A release of a Petri run's writer lease by its owner. + type: object + required: + - owner + properties: + owner: + type: string + description: The owner id that holds the lease. + example: 18f3c2a9e1b4-42017-0-9f3a1c7e2b5d + + PetriRecord: + description: >- + One stored line of a Petri log: the record as JSON, with `seq` and + `recorded_at` lifted out of it so the store can key and index without + reading into the JSON. What the store hands back equals what it was + given as a JSON value. + type: object + required: + - seq + - recorded_at + - record + properties: + seq: + type: integer + format: uint64 + description: The record's position in its log, from 0. + example: 7 + recorded_at: + type: integer + format: uint64 + description: Milliseconds since the Unix epoch when Petri appended the record. + example: 1758067200123 + record: + type: object + additionalProperties: true + description: The record itself, stored and read back unchanged. + + PetriAppendRequest: + description: One batch of records for one Petri log. + type: object + required: + - owner + - records + properties: + owner: + type: string + description: The owner id the worker holds the run's writer lease with. + example: 18f3c2a9e1b4-42017-0-9f3a1c7e2b5d + records: + type: array + items: + $ref: "#/components/schemas/PetriRecord" + + PetriRecordList: + description: Every record of one Petri log, in `seq` order. + type: object + required: + - records + properties: + records: + type: array + items: + $ref: "#/components/schemas/PetriRecord" + CommandTermination: description: Terminal state for a command execution. type: string diff --git a/lib/apps/fabro-cli/src/commands/parent/mod.rs b/lib/apps/fabro-cli/src/commands/parent/mod.rs index aae5a326e..09f86bb7a 100644 --- a/lib/apps/fabro-cli/src/commands/parent/mod.rs +++ b/lib/apps/fabro-cli/src/commands/parent/mod.rs @@ -9,7 +9,9 @@ use crate::command_context::CommandContext; pub(crate) async fn dispatch(ns: ParentNamespace, base_ctx: &CommandContext) -> Result<()> { match ns.command { - ParentCommand::Link(args) => link::link_command(args, base_ctx).await, + // The link command's future carries several client calls and sits + // past clippy's stack budget; box it once at the call. + ParentCommand::Link(args) => Box::pin(link::link_command(args, base_ctx)).await, ParentCommand::Unlink(args) => unlink::unlink_command(args, base_ctx).await, } } diff --git a/lib/apps/fabro-server/src/error.rs b/lib/apps/fabro-server/src/error.rs index 910e9b2dc..4a31fc65b 100644 --- a/lib/apps/fabro-server/src/error.rs +++ b/lib/apps/fabro-server/src/error.rs @@ -4,6 +4,7 @@ use axum::response::{IntoResponse, Response}; use fabro_api::types::ErrorResponseEntry; use fabro_vault::Error as VaultError; use serde::Serialize; +use serde_json::{Map, Value}; #[derive(Debug, thiserror::Error)] pub enum Error { @@ -59,6 +60,8 @@ struct ErrorEntry { detail: String, #[serde(skip_serializing_if = "Option::is_none")] code: Option, + #[serde(skip_serializing_if = "Option::is_none")] + meta: Option>>, } #[derive(Serialize)] @@ -69,12 +72,15 @@ struct ErrorBody { /// Uniform API error response. /// /// Serializes to `{"errors": [{"status": "4xx", "title": "...", "detail": -/// "..."}]}`. +/// "..."}]}`, with `code` and `meta` when the error carries them. #[derive(Clone, Debug)] pub struct ApiError { status: StatusCode, detail: String, code: Option, + /// Structured details specific to `code`, for clients that act on them. + /// Boxed so the error stays small in every `Result` that carries it. + meta: Option>>, } impl ApiError { @@ -83,6 +89,7 @@ impl ApiError { status, detail: detail.into(), code: None, + meta: None, } } @@ -95,6 +102,22 @@ impl ApiError { status, detail: detail.into(), code: Some(code.into()), + meta: None, + } + } + + /// An error whose `code` documents the members of `meta`. + pub fn with_code_and_meta( + status: StatusCode, + detail: impl Into, + code: impl Into, + meta: Map, + ) -> Self { + Self { + status, + detail: detail.into(), + code: Some(code.into()), + meta: Some(Box::new(meta)), } } @@ -148,6 +171,7 @@ impl ApiError { .to_string(), detail: self.detail, code: self.code, + meta: self.meta.map(|meta| *meta).unwrap_or_default(), request_id: None, } } @@ -201,6 +225,7 @@ impl IntoResponse for ApiError { title, detail: self.detail, code: self.code, + meta: self.meta, }], }; (self.status, Json(body)).into_response() @@ -244,6 +269,38 @@ mod tests { ); } + #[tokio::test] + async fn api_error_with_meta_serializes_the_members_beside_the_code() { + let mut meta = serde_json::Map::new(); + meta.insert("owner".to_string(), json!("worker-1")); + let response = ApiError::with_code_and_meta( + StatusCode::CONFLICT, + "run is leased", + "petri_run_leased", + meta, + ) + .into_response(); + + let body = to_bytes(response.into_body(), usize::MAX) + .await + .expect("body should serialize"); + let body: serde_json::Value = + serde_json::from_slice(&body).expect("response body should be valid json"); + + assert_eq!( + body, + json!({ + "errors": [{ + "status": "409", + "title": "Conflict", + "detail": "run is leased", + "code": "petri_run_leased", + "meta": { "owner": "worker-1" } + }] + }) + ); + } + #[tokio::test] async fn unauthorized_without_code_omits_code_key() { let response = ApiError::unauthorized().into_response(); diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index 73ebce687..03f16e479 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1936,6 +1936,143 @@ impl Client { } } + // ── The Petri run store, as a worker reaches it ────────────────────── + // + // Each method is one request and answers with the server's reply as it + // is: an error carries the store's `code` and `meta` in its `ApiFailure` + // (see `api_failure_for`), and a transport failure carries none. Retry + // policy belongs to the store implementation over these calls, not here. + + /// Open the run in the Petri run store for the worker. + pub async fn open_petri_run( + &self, + run_id: &RunId, + body: types::PetriOpenRequest, + ) -> Result { + let response = self + .send_api(|client| async move { + client + .open_petri_run() + .id(run_id.to_string()) + .body(body.clone()) + .send() + .await + }) + .await?; + Ok(response.into_inner()) + } + + /// End the worker's writer lease on the run when `owner` still holds it. + pub async fn release_petri_run(&self, run_id: &RunId, owner: &str) -> Result<()> { + self.send_api(|client| async move { + client + .release_petri_run() + .id(run_id.to_string()) + .body(types::PetriReleaseRequest { + owner: owner.to_string(), + }) + .send() + .await + }) + .await?; + Ok(()) + } + + /// Append one batch of records to one log of the run. `log` is the log + /// id as text; the generated client percent-encodes it. + pub async fn append_petri_records( + &self, + run_id: &RunId, + log: &str, + body: types::PetriAppendRequest, + ) -> Result<()> { + self.send_api(|client| async move { + client + .append_petri_records() + .id(run_id.to_string()) + .log(log) + .body(body.clone()) + .send() + .await + }) + .await?; + Ok(()) + } + + /// Every record of one log of the run, in `seq` order. + pub async fn list_petri_records( + &self, + run_id: &RunId, + log: &str, + ) -> Result> { + let response = self + .send_api(|client| async move { + client + .list_petri_records() + .id(run_id.to_string()) + .log(log) + .send() + .await + }) + .await?; + Ok(response.into_inner().records) + } + + /// Store a blob by content for the run's `owner` and get its digest. + pub async fn write_petri_blob( + &self, + run_id: &RunId, + owner: &str, + data: &[u8], + ) -> Result { + let response = self + .send_api(|client| async move { + client + .write_petri_blob() + .id(run_id.to_string()) + .owner(owner) + .body(data.to_vec()) + .send() + .await + }) + .await?; + Ok(response.into_inner().hash) + } + + /// The blob with this digest, or `None` when the store holds no such + /// blob. A run the store does not hold is an error. + pub async fn read_petri_blob( + &self, + run_id: &RunId, + blob_hash: &BlobHash, + ) -> Result> { + let response = self + .current_state() + .client + .read_petri_blob() + .id(run_id.to_string()) + .blob_hash(*blob_hash) + .send() + .await; + match response { + Ok(response) => { + let mut stream = response.into_inner(); + let mut bytes = Vec::new(); + while let Some(chunk) = stream.next().await { + let chunk = chunk.map_err(anyhow::Error::new)?; + bytes.extend_from_slice(&chunk); + } + Ok(Some(Bytes::from(bytes))) + } + Err(err) => { + let err = classify_api_error(err).await.error; + let blob_missing = api_failure_for(&err) + .is_some_and(|failure| failure.code.as_deref() == Some("petri_blob_not_found")); + if blob_missing { Ok(None) } else { Err(err) } + } + } + } + #[expect( clippy::disallowed_types, reason = "Client builds raw server API request URLs for wire transit; logging redaction is handled at log boundaries." @@ -3261,10 +3398,7 @@ mod tests { fn add_pr_upgrade_hint_appends_on_unstructured_404() { let err = tag_with_failure( anyhow!("request failed with status 404 Not Found"), - ApiFailure { - status: fabro_http::StatusCode::NOT_FOUND, - code: None, - }, + ApiFailure::new(fabro_http::StatusCode::NOT_FOUND, None), ); let wrapped = super::add_pr_upgrade_hint(err); let message = wrapped.to_string(); @@ -3279,10 +3413,10 @@ mod tests { fn add_pr_upgrade_hint_does_not_touch_structured_404() { let err = tag_with_failure( anyhow!("No pull request found in store. Create one first with: fabro pr create abc"), - ApiFailure { - status: fabro_http::StatusCode::NOT_FOUND, - code: Some("no_stored_record".to_string()), - }, + ApiFailure::new( + fabro_http::StatusCode::NOT_FOUND, + Some("no_stored_record".to_string()), + ), ); let wrapped = super::add_pr_upgrade_hint(err); let message = wrapped.to_string(); diff --git a/lib/foundation/fabro-client/src/error.rs b/lib/foundation/fabro-client/src/error.rs index b5b2c9a07..03a03efc0 100644 --- a/lib/foundation/fabro-client/src/error.rs +++ b/lib/foundation/fabro-client/src/error.rs @@ -6,6 +6,28 @@ use serde::de::DeserializeOwned; pub struct ApiFailure { pub status: fabro_http::StatusCode, pub code: Option, + /// The error entry's `meta`: structured details specific to `code`. + pub meta: Option, +} + +impl ApiFailure { + #[must_use] + pub fn new(status: fabro_http::StatusCode, code: Option) -> Self { + Self { + status, + code, + meta: None, + } + } +} + +/// The first entry of an `ErrorResponse` body, as far as a client acts on +/// it. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct ParsedErrorEntry { + pub detail: Option, + pub code: Option, + pub meta: Option, } // Transparent wrapper that attaches an ApiFailure to an anyhow error while @@ -66,19 +88,31 @@ impl ApiError { } pub fn parse_error_response_value(value: &serde_json::Value) -> (Option, Option) { + let entry = parse_error_response_entry(value); + (entry.detail, entry.code) +} + +/// The first error entry of an `ErrorResponse` body: its detail, code and +/// `meta`, each absent when the body does not carry it. +pub fn parse_error_response_entry(value: &serde_json::Value) -> ParsedErrorEntry { let first = value .get("errors") .and_then(serde_json::Value::as_array) .and_then(|errors| errors.first()); - let detail = first - .and_then(|entry| entry.get("detail")) - .and_then(serde_json::Value::as_str) - .map(ToOwned::to_owned); - let code = first - .and_then(|entry| entry.get("code")) - .and_then(serde_json::Value::as_str) - .map(ToOwned::to_owned); - (detail, code) + let text = |field: &str| { + first + .and_then(|entry| entry.get(field)) + .and_then(serde_json::Value::as_str) + .map(ToOwned::to_owned) + }; + ParsedErrorEntry { + detail: text("detail"), + code: text("code"), + meta: first + .and_then(|entry| entry.get("meta")) + .filter(|meta| meta.is_object()) + .cloned(), + } } fn classify_from_status(err: anyhow::Error, status: fabro_http::StatusCode) -> anyhow::Error { @@ -92,9 +126,13 @@ fn classify_from_status(err: anyhow::Error, status: fabro_http::StatusCode) -> a fn build_structured_error( error: anyhow::Error, status: fabro_http::StatusCode, - code: Option, + entry: ParsedErrorEntry, ) -> StructuredApiError { - let failure = ApiFailure { status, code }; + let failure = ApiFailure { + status, + code: entry.code, + meta: entry.meta, + }; let tagged = tag_with_failure(error, failure.clone()); StructuredApiError { error: classify_from_status(tagged, status), @@ -110,12 +148,11 @@ where progenitor_client::Error::UnexpectedResponse(response) => { let status = response.status(); let body = response.text().await.unwrap_or_default(); - let mut code = None; + let mut entry = ParsedErrorEntry::default(); if let Ok(value) = serde_json::from_str::(&body) { - let (detail, parsed_code) = parse_error_response_value(&value); - code = parsed_code; - if let Some(detail) = detail { - return build_structured_error(anyhow!("{detail}"), status, code); + entry = parse_error_response_entry(&value); + if let Some(detail) = entry.detail.take() { + return build_structured_error(anyhow!("{detail}"), status, entry); } } let error = if body.is_empty() { @@ -123,7 +160,7 @@ where } else { anyhow!("request failed with status {status}: {body}") }; - build_structured_error(error, status, code) + build_structured_error(error, status, entry) } other => map_api_error_structured(other), } @@ -136,19 +173,26 @@ where match err { progenitor_client::Error::ErrorResponse(response) => { let status = response.status(); - let mut code = None; + let mut entry = ParsedErrorEntry::default(); if let Ok(value) = serde_json::to_value(response.into_inner()) { - let (detail, parsed_code) = parse_error_response_value(&value); - code = parsed_code; - if let Some(detail) = detail { - return build_structured_error(anyhow!("{detail}"), status, code); + entry = parse_error_response_entry(&value); + if let Some(detail) = entry.detail.take() { + return build_structured_error(anyhow!("{detail}"), status, entry); } } - build_structured_error(anyhow!("request failed with status {status}"), status, code) + build_structured_error( + anyhow!("request failed with status {status}"), + status, + entry, + ) } progenitor_client::Error::UnexpectedResponse(response) => { let status = response.status(); - build_structured_error(anyhow!("request failed with status {status}"), status, None) + build_structured_error( + anyhow!("request failed with status {status}"), + status, + ParsedErrorEntry::default(), + ) } other => StructuredApiError { error: anyhow::Error::new(other), @@ -202,17 +246,19 @@ pub async fn classify_http_response( let status = response.status(); let headers = response.headers().clone(); let body = response.text().await.unwrap_or_default(); - let mut code = None; - if let Ok(value) = serde_json::from_str::(&body) { - let (_, parsed_code) = parse_error_response_value(&value); - code = parsed_code; - } + let entry = serde_json::from_str::(&body) + .map(|value| parse_error_response_entry(&value)) + .unwrap_or_default(); Ok(Err(ApiError { status, headers, body, - failure: ApiFailure { status, code }, + failure: ApiFailure { + status, + code: entry.code, + meta: entry.meta, + }, })) } @@ -269,10 +315,7 @@ mod tests { }] })) .unwrap(), - failure: ApiFailure { - status, - code: Some(code.to_string()), - }, + failure: ApiFailure::new(status, Some(code.to_string())), } } @@ -313,6 +356,27 @@ mod tests { assert_eq!(failure.code.as_deref(), Some("invalid_manifest")); } + #[test] + fn map_api_error_carries_the_entry_meta() { + let response = progenitor_client::ResponseValue::new( + json!({ + "errors": [{ + "detail": "run is leased", + "code": "petri_run_leased", + "meta": { "owner": "worker-1" }, + }] + }), + fabro_http::StatusCode::CONFLICT, + fabro_http::HeaderMap::new(), + ); + let err = + map_api_error(progenitor_client::Error::::ErrorResponse(response)); + + let failure = api_failure_for(&err).expect("error should carry API failure metadata"); + assert_eq!(failure.code.as_deref(), Some("petri_run_leased")); + assert_eq!(failure.meta, Some(json!({ "owner": "worker-1" }))); + } + #[test] fn raw_response_failure_error_marks_401_as_auth_required() { let err = raw_response_failure_error(&api_error( diff --git a/lib/foundation/fabro-client/src/lib.rs b/lib/foundation/fabro-client/src/lib.rs index c5fd914bf..e65e02560 100644 --- a/lib/foundation/fabro-client/src/lib.rs +++ b/lib/foundation/fabro-client/src/lib.rs @@ -16,9 +16,9 @@ pub use client::{ }; pub use credential::{Credential, CredentialFallback}; pub use error::{ - ApiError, ApiFailure, StructuredApiError, classify_api_error, classify_http_response, - convert_type, is_not_found_error, map_api_error, parse_error_response_value, - raw_response_failure_error, + ApiError, ApiFailure, ParsedErrorEntry, StructuredApiError, api_failure_for, + classify_api_error, classify_http_response, convert_type, is_not_found_error, map_api_error, + parse_error_response_entry, parse_error_response_value, raw_response_failure_error, }; pub use session::OAuthSession; pub use target::ServerTarget; diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index fec89c29e..0d0b0f64d 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -295,6 +295,13 @@ models/parallel-branch-result.ts models/pending-interview-record.ts models/pending-reason.ts models/permission-level.ts +models/petri-access.ts +models/petri-append-request.ts +models/petri-open-request.ts +models/petri-open-response.ts +models/petri-record-list.ts +models/petri-record.ts +models/petri-release-request.ts models/preflight-check-detail.ts models/preflight-check-report.ts models/preflight-check-result.ts diff --git a/lib/packages/fabro-api-client/src/api/run-internals-api.ts b/lib/packages/fabro-api-client/src/api/run-internals-api.ts index 771501839..ba15825ed 100644 --- a/lib/packages/fabro-api-client/src/api/run-internals-api.ts +++ b/lib/packages/fabro-api-client/src/api/run-internals-api.ts @@ -34,6 +34,16 @@ import type { PaginatedEventList } from '../models'; // @ts-ignore import type { PaginatedRunStageList } from '../models'; // @ts-ignore +import type { PetriAppendRequest } from '../models'; +// @ts-ignore +import type { PetriOpenRequest } from '../models'; +// @ts-ignore +import type { PetriOpenResponse } from '../models'; +// @ts-ignore +import type { PetriRecordList } from '../models'; +// @ts-ignore +import type { PetriReleaseRequest } from '../models'; +// @ts-ignore import type { RunArtifactListResponse } from '../models'; // @ts-ignore import type { RunCheckpoint } from '../models'; @@ -56,6 +66,55 @@ import type { WriteRunBlobRequest } from '../models'; */ export const RunInternalsApiAxiosParamCreator = function (configuration?: Configuration) { return { + /** + * Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing. + * @summary Append Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {PetriAppendRequest} petriAppendRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + appendPetriRecords: async (id: string, log: string, petriAppendRequest: PetriAppendRequest, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('appendPetriRecords', 'id', id) + // verify required parameter 'log' is not null or undefined + assertParamExists('appendPetriRecords', 'log', log) + // verify required parameter 'petriAppendRequest' is not null or undefined + assertParamExists('appendPetriRecords', 'petriAppendRequest', petriAppendRequest) + const localVarPath = `/api/v1/runs/{id}/petri/logs/{log}/records` + .replace(`{${"id"}}`, encodeURIComponent(String(id))) + .replace(`{${"log"}}`, encodeURIComponent(String(log))); + // 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: 'POST', ...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/json'; + localVarHeaderParameter['Accept'] = 'application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + localVarRequestOptions.data = serializeDataIfNeeded(petriAppendRequest, localVarRequestOptions, configuration) + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Appends a validated event to the run event log. Intended for trusted internal callers. * @summary Append Run Event @@ -471,6 +530,50 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, + /** + * Every record of one log of the run, in `seq` order, unchanged. + * @summary List Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + listPetriRecords: async (id: string, log: string, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('listPetriRecords', 'id', id) + // verify required parameter 'log' is not null or undefined + assertParamExists('listPetriRecords', 'log', log) + const localVarPath = `/api/v1/runs/{id}/petri/logs/{log}/records` + .replace(`{${"id"}}`, encodeURIComponent(String(id))) + .replace(`{${"log"}}`, encodeURIComponent(String(log))); + // 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: 'GET', ...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['Accept'] = 'application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Lists captured artifact files for a run. * @summary List Run Artifacts @@ -719,6 +822,51 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, + /** + * Opens the run in the Petri run store for the worker. `create` inserts the run and takes its writer lease for `owner`; `write` takes the lease of an existing run; `read` takes no lease. The lease is idempotent per owner: a retry by the owner that holds it gets the same lease. Another live owner is refused with `petri_run_leased`. The lease ends when the worker releases it, when the server observes the worker exit, or by operator release, never by timeout. + * @summary Open Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriOpenRequest} petriOpenRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + openPetriRun: async (id: string, petriOpenRequest: PetriOpenRequest, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('openPetriRun', 'id', id) + // verify required parameter 'petriOpenRequest' is not null or undefined + assertParamExists('openPetriRun', 'petriOpenRequest', petriOpenRequest) + const localVarPath = `/api/v1/runs/{id}/petri/open` + .replace(`{${"id"}}`, encodeURIComponent(String(id))); + // 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: 'POST', ...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/json'; + localVarHeaderParameter['Accept'] = 'application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + localVarRequestOptions.data = serializeDataIfNeeded(petriOpenRequest, localVarRequestOptions, configuration) + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. * @summary Put Stage Artifact @@ -780,6 +928,50 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, + /** + * The blob with this digest, if the store holds one. + * @summary Read Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} blobHash Content-addressed blob hash. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + readPetriBlob: async (id: string, blobHash: string, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('readPetriBlob', 'id', id) + // verify required parameter 'blobHash' is not null or undefined + assertParamExists('readPetriBlob', 'blobHash', blobHash) + const localVarPath = `/api/v1/runs/{id}/petri/blobs/{blobHash}` + .replace(`{${"id"}}`, encodeURIComponent(String(id))) + .replace(`{${"blobHash"}}`, encodeURIComponent(String(blobHash))); + // 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: 'GET', ...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['Accept'] = 'application/octet-stream,application/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Reads a previously stored blob by hash. * @summary Read Run Blob @@ -824,6 +1016,50 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, + /** + * Ends the worker\'s writer lease on the run when `owner` still holds it: what a worker sends when it drops its store handle. A lease that already moved to another owner is left alone. + * @summary Release Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriReleaseRequest} petriReleaseRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + releasePetriRun: async (id: string, petriReleaseRequest: PetriReleaseRequest, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('releasePetriRun', 'id', id) + // verify required parameter 'petriReleaseRequest' is not null or undefined + assertParamExists('releasePetriRun', 'petriReleaseRequest', petriReleaseRequest) + const localVarPath = `/api/v1/runs/{id}/petri/release` + .replace(`{${"id"}}`, encodeURIComponent(String(id))); + // 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: 'POST', ...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/json'; + + setSearchParams(localVarUrlObj, localVarQueryParameter); + let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; + localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; + localVarRequestOptions.data = serializeDataIfNeeded(petriReleaseRequest, localVarRequestOptions, configuration) + + return { + url: toPathString(localVarUrlObj), + options: localVarRequestOptions, + }; + }, /** * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. * @summary Retrieve Run Checkpoint @@ -904,6 +1140,58 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, + /** + * Stores a blob by content for the run\'s owner and returns its SHA-256 digest, the same content address Fabro\'s blob store uses. Idempotent by construction. + * @summary Write Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} owner The owner id the worker opened the run\'s writer lease with. + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + writePetriBlob: async (id: string, owner: string, body: File, options: RawAxiosRequestConfig = {}): Promise => { + // verify required parameter 'id' is not null or undefined + assertParamExists('writePetriBlob', 'id', id) + // verify required parameter 'owner' is not null or undefined + assertParamExists('writePetriBlob', 'owner', owner) + // verify required parameter 'body' is not null or undefined + assertParamExists('writePetriBlob', 'body', body) + const localVarPath = `/api/v1/runs/{id}/petri/blobs` + .replace(`{${"id"}}`, encodeURIComponent(String(id))); + // 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: 'POST', ...baseOptions, ...options}; + const localVarHeaderParameter = {} as any; + const localVarQueryParameter = {} as any; + + // authentication SessionCookie required + + // authentication BearerAuth required + // http bearer authentication required + await setBearerAuthToObject(localVarHeaderParameter, configuration) + + if (owner !== undefined) { + localVarQueryParameter['owner'] = owner; + } + + 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 @@ -958,6 +1246,21 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarAxiosParamCreator = RunInternalsApiAxiosParamCreator(configuration) return { + /** + * Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing. + * @summary Append Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {PetriAppendRequest} petriAppendRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async appendPetriRecords(id: string, log: string, petriAppendRequest: PetriAppendRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.appendPetriRecords(id, log, petriAppendRequest, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.appendPetriRecords']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Appends a validated event to the run event log. Intended for trusted internal callers. * @summary Append Run Event @@ -1086,6 +1389,20 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.getStageArtifact']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, + /** + * Every record of one log of the run, in `seq` order, unchanged. + * @summary List Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async listPetriRecords(id: string, log: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.listPetriRecords(id, log, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.listPetriRecords']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Lists captured artifact files for a run. * @summary List Run Artifacts @@ -1161,6 +1478,20 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.listStageEvents']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, + /** + * Opens the run in the Petri run store for the worker. `create` inserts the run and takes its writer lease for `owner`; `write` takes the lease of an existing run; `read` takes no lease. The lease is idempotent per owner: a retry by the owner that holds it gets the same lease. Another live owner is refused with `petri_run_leased`. The lease ends when the worker releases it, when the server observes the worker exit, or by operator release, never by timeout. + * @summary Open Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriOpenRequest} petriOpenRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async openPetriRun(id: string, petriOpenRequest: PetriOpenRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.openPetriRun(id, petriOpenRequest, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.openPetriRun']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. * @summary Put Stage Artifact @@ -1178,6 +1509,20 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.putStageArtifact']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, + /** + * The blob with this digest, if the store holds one. + * @summary Read Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} blobHash Content-addressed blob hash. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async readPetriBlob(id: string, blobHash: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.readPetriBlob(id, blobHash, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.readPetriBlob']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Reads a previously stored blob by hash. * @summary Read Run Blob @@ -1192,6 +1537,20 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.readRunBlob']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, + /** + * Ends the worker\'s writer lease on the run when `owner` still holds it: what a worker sends when it drops its store handle. A lease that already moved to another owner is left alone. + * @summary Release Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriReleaseRequest} petriReleaseRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async releasePetriRun(id: string, petriReleaseRequest: PetriReleaseRequest, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.releasePetriRun(id, petriReleaseRequest, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.releasePetriRun']?.[localVarOperationServerIndex]?.url; + return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); + }, /** * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. * @summary Retrieve Run Checkpoint @@ -1218,6 +1577,21 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.retrieveRunSettings']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, + /** + * Stores a blob by content for the run\'s owner and returns its SHA-256 digest, the same content address Fabro\'s blob store uses. Idempotent by construction. + * @summary Write Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} owner The owner id the worker opened the run\'s writer lease with. + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + async writePetriBlob(id: string, owner: string, body: File, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.writePetriBlob(id, owner, body, options); + const localVarOperationServerIndex = configuration?.serverIndex ?? 0; + const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.writePetriBlob']?.[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 @@ -1241,6 +1615,18 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { export const RunInternalsApiFactory = function (configuration?: Configuration, basePath?: string, axios?: AxiosInstance) { const localVarFp = RunInternalsApiFp(configuration) return { + /** + * Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing. + * @summary Append Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {PetriAppendRequest} petriAppendRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + appendPetriRecords(id: string, log: string, petriAppendRequest: PetriAppendRequest, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.appendPetriRecords(id, log, petriAppendRequest, options).then((request) => request(axios, basePath)); + }, /** * Appends a validated event to the run event log. Intended for trusted internal callers. * @summary Append Run Event @@ -1342,6 +1728,17 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b getStageArtifact(id: string, stageId: string, filename: string, retry: number, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.getStageArtifact(id, stageId, filename, retry, options).then((request) => request(axios, basePath)); }, + /** + * Every record of one log of the run, in `seq` order, unchanged. + * @summary List Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + listPetriRecords(id: string, log: string, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.listPetriRecords(id, log, options).then((request) => request(axios, basePath)); + }, /** * Lists captured artifact files for a run. * @summary List Run Artifacts @@ -1402,6 +1799,17 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b listStageEvents(id: string, stageId: string, sinceSeq?: number, limit?: number, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.listStageEvents(id, stageId, sinceSeq, limit, options).then((request) => request(axios, basePath)); }, + /** + * Opens the run in the Petri run store for the worker. `create` inserts the run and takes its writer lease for `owner`; `write` takes the lease of an existing run; `read` takes no lease. The lease is idempotent per owner: a retry by the owner that holds it gets the same lease. Another live owner is refused with `petri_run_leased`. The lease ends when the worker releases it, when the server observes the worker exit, or by operator release, never by timeout. + * @summary Open Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriOpenRequest} petriOpenRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + openPetriRun(id: string, petriOpenRequest: PetriOpenRequest, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.openPetriRun(id, petriOpenRequest, options).then((request) => request(axios, basePath)); + }, /** * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. * @summary Put Stage Artifact @@ -1416,6 +1824,17 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b putStageArtifact(id: string, stageId: string, retry: number, body: File, filename?: string, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.putStageArtifact(id, stageId, retry, body, filename, options).then((request) => request(axios, basePath)); }, + /** + * The blob with this digest, if the store holds one. + * @summary Read Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} blobHash Content-addressed blob hash. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + readPetriBlob(id: string, blobHash: string, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.readPetriBlob(id, blobHash, options).then((request) => request(axios, basePath)); + }, /** * Reads a previously stored blob by hash. * @summary Read Run Blob @@ -1427,6 +1846,17 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b readRunBlob(id: string, blobHash: string, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.readRunBlob(id, blobHash, options).then((request) => request(axios, basePath)); }, + /** + * Ends the worker\'s writer lease on the run when `owner` still holds it: what a worker sends when it drops its store handle. A lease that already moved to another owner is left alone. + * @summary Release Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriReleaseRequest} petriReleaseRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + releasePetriRun(id: string, petriReleaseRequest: PetriReleaseRequest, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.releasePetriRun(id, petriReleaseRequest, options).then((request) => request(axios, basePath)); + }, /** * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. * @summary Retrieve Run Checkpoint @@ -1447,6 +1877,18 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b retrieveRunSettings(id: string, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.retrieveRunSettings(id, options).then((request) => request(axios, basePath)); }, + /** + * Stores a blob by content for the run\'s owner and returns its SHA-256 digest, the same content address Fabro\'s blob store uses. Idempotent by construction. + * @summary Write Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} owner The owner id the worker opened the run\'s writer lease with. + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + writePetriBlob(id: string, owner: string, body: File, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.writePetriBlob(id, owner, body, options).then((request) => request(axios, basePath)); + }, /** * Writes an opaque binary blob and returns its content-addressed blob hash. * @summary Write Run Blob @@ -1465,6 +1907,19 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b * RunInternalsApi - object-oriented interface */ export class RunInternalsApi extends BaseAPI { + /** + * Appends one batch of records to one log at the sequences they carry, durably, in one transaction. A record equal to the one already stored at its `seq` is accepted without a second append, so a batch whose reply was lost is safe to resend. A different record at a taken `seq`, or a `seq` past the log\'s end, is refused with `petri_record_conflict` and the batch stores nothing. + * @summary Append Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {PetriAppendRequest} petriAppendRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public appendPetriRecords(id: string, log: string, petriAppendRequest: PetriAppendRequest, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).appendPetriRecords(id, log, petriAppendRequest, options).then((request) => request(this.axios, this.basePath)); + } + /** * Appends a validated event to the run event log. Intended for trusted internal callers. * @summary Append Run Event @@ -1575,6 +2030,18 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).getStageArtifact(id, stageId, filename, retry, options).then((request) => request(this.axios, this.basePath)); } + /** + * Every record of one log of the run, in `seq` order, unchanged. + * @summary List Petri Records + * @param {string} id Unique run identifier (ULID). + * @param {string} log A Petri log of the run, as its id renders: `coordinator`, `resources`, or `execution <n>` for execution `n`. The space is percent-encoded on the wire. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public listPetriRecords(id: string, log: string, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).listPetriRecords(id, log, options).then((request) => request(this.axios, this.basePath)); + } + /** * Lists captured artifact files for a run. * @summary List Run Artifacts @@ -1640,6 +2107,18 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).listStageEvents(id, stageId, sinceSeq, limit, options).then((request) => request(this.axios, this.basePath)); } + /** + * Opens the run in the Petri run store for the worker. `create` inserts the run and takes its writer lease for `owner`; `write` takes the lease of an existing run; `read` takes no lease. The lease is idempotent per owner: a retry by the owner that holds it gets the same lease. Another live owner is refused with `petri_run_leased`. The lease ends when the worker releases it, when the server observes the worker exit, or by operator release, never by timeout. + * @summary Open Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriOpenRequest} petriOpenRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public openPetriRun(id: string, petriOpenRequest: PetriOpenRequest, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).openPetriRun(id, petriOpenRequest, options).then((request) => request(this.axios, this.basePath)); + } + /** * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. * @summary Put Stage Artifact @@ -1655,6 +2134,18 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).putStageArtifact(id, stageId, retry, body, filename, options).then((request) => request(this.axios, this.basePath)); } + /** + * The blob with this digest, if the store holds one. + * @summary Read Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} blobHash Content-addressed blob hash. + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public readPetriBlob(id: string, blobHash: string, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).readPetriBlob(id, blobHash, options).then((request) => request(this.axios, this.basePath)); + } + /** * Reads a previously stored blob by hash. * @summary Read Run Blob @@ -1667,6 +2158,18 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).readRunBlob(id, blobHash, options).then((request) => request(this.axios, this.basePath)); } + /** + * Ends the worker\'s writer lease on the run when `owner` still holds it: what a worker sends when it drops its store handle. A lease that already moved to another owner is left alone. + * @summary Release Petri Run + * @param {string} id Unique run identifier (ULID). + * @param {PetriReleaseRequest} petriReleaseRequest + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public releasePetriRun(id: string, petriReleaseRequest: PetriReleaseRequest, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).releasePetriRun(id, petriReleaseRequest, options).then((request) => request(this.axios, this.basePath)); + } + /** * Returns the latest checkpoint data for a run, or null if no checkpoint has been recorded yet. * @summary Retrieve Run Checkpoint @@ -1689,6 +2192,19 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).retrieveRunSettings(id, options).then((request) => request(this.axios, this.basePath)); } + /** + * Stores a blob by content for the run\'s owner and returns its SHA-256 digest, the same content address Fabro\'s blob store uses. Idempotent by construction. + * @summary Write Petri Blob + * @param {string} id Unique run identifier (ULID). + * @param {string} owner The owner id the worker opened the run\'s writer lease with. + * @param {File} body + * @param {*} [options] Override http request option. + * @throws {RequiredError} + */ + public writePetriBlob(id: string, owner: string, body: File, options?: RawAxiosRequestConfig) { + return RunInternalsApiFp(this.configuration).writePetriBlob(id, owner, body, options).then((request) => request(this.axios, this.basePath)); + } + /** * Writes an opaque binary blob and returns its content-addressed blob hash. * @summary Write Run Blob diff --git a/lib/packages/fabro-api-client/src/models/error-response-entry.ts b/lib/packages/fabro-api-client/src/models/error-response-entry.ts index 081416711..8ff4acd64 100644 --- a/lib/packages/fabro-api-client/src/models/error-response-entry.ts +++ b/lib/packages/fabro-api-client/src/models/error-response-entry.ts @@ -38,4 +38,8 @@ export interface ErrorResponseEntry { * Server-generated request identifier; matches the x-request-id response header. */ 'request_id'?: string; + /** + * Optional structured details specific to the error `code`, for clients that act on them. Each code documents the members it sets. + */ + 'meta'?: { [key: string]: any; }; } diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index b52c24aec..897027f9d 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -265,6 +265,13 @@ export * from './parallel-branch-result'; export * from './pending-interview-record'; export * from './pending-reason'; export * from './permission-level'; +export * from './petri-access'; +export * from './petri-append-request'; +export * from './petri-open-request'; +export * from './petri-open-response'; +export * from './petri-record'; +export * from './petri-record-list'; +export * from './petri-release-request'; export * from './preflight-check-detail'; export * from './preflight-check-report'; export * from './preflight-check-result'; diff --git a/lib/packages/fabro-api-client/src/models/petri-access.ts b/lib/packages/fabro-api-client/src/models/petri-access.ts new file mode 100644 index 000000000..fcc983f3d --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-access.ts @@ -0,0 +1,27 @@ +/* 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. + */ + + + +/** + * How a worker opens a Petri run. `create` inserts the run and takes its writer lease, `write` takes the lease of an existing run, `read` takes no lease. + */ + +export const PetriAccess = { + CREATE: 'create', + WRITE: 'write', + READ: 'read' +} as const; + +export type PetriAccess = typeof PetriAccess[keyof typeof PetriAccess]; diff --git a/lib/packages/fabro-api-client/src/models/petri-append-request.ts b/lib/packages/fabro-api-client/src/models/petri-append-request.ts new file mode 100644 index 000000000..7bb9f619c --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-append-request.ts @@ -0,0 +1,29 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { PetriRecord } from './petri-record'; + +/** + * One batch of records for one Petri log. + */ +export interface PetriAppendRequest { + /** + * The owner id the worker holds the run\'s writer lease with. + */ + 'owner': string; + 'records': Array; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-open-request.ts b/lib/packages/fabro-api-client/src/models/petri-open-request.ts new file mode 100644 index 000000000..89ceea4b7 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-open-request.ts @@ -0,0 +1,29 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.2.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { PetriAccess } from './petri-access'; + +/** + * An open of a Petri run for a worker. + */ +export interface PetriOpenRequest { + 'access': PetriAccess; + /** + * The owner id the writer lease is taken for. Required for `create` and `write`, absent for `read`. + */ + 'owner'?: string; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-open-response.ts b/lib/packages/fabro-api-client/src/models/petri-open-response.ts new file mode 100644 index 000000000..356dbc96c --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-open-response.ts @@ -0,0 +1,25 @@ +/* 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 Petri run opened for a worker. + */ +export interface PetriOpenResponse { + /** + * Where the run lives, for messages. + */ + 'locator': string; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-record-list.ts b/lib/packages/fabro-api-client/src/models/petri-record-list.ts new file mode 100644 index 000000000..9345875ea --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-record-list.ts @@ -0,0 +1,25 @@ +/* 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. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { PetriRecord } from './petri-record'; + +/** + * Every record of one Petri log, in `seq` order. + */ +export interface PetriRecordList { + 'records': Array; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-record.ts b/lib/packages/fabro-api-client/src/models/petri-record.ts new file mode 100644 index 000000000..26421fa45 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-record.ts @@ -0,0 +1,33 @@ +/* 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. + */ + + + +/** + * One stored line of a Petri log: the record as JSON, with `seq` and `recorded_at` lifted out of it so the store can key and index without reading into the JSON. What the store hands back equals what it was given as a JSON value. + */ +export interface PetriRecord { + /** + * The record\'s position in its log, from 0. + */ + 'seq': number; + /** + * Milliseconds since the Unix epoch when Petri appended the record. + */ + 'recorded_at': number; + /** + * The record itself, stored and read back unchanged. + */ + 'record': { [key: string]: any; }; +} diff --git a/lib/packages/fabro-api-client/src/models/petri-release-request.ts b/lib/packages/fabro-api-client/src/models/petri-release-request.ts new file mode 100644 index 000000000..05246bb44 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/petri-release-request.ts @@ -0,0 +1,25 @@ +/* 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 release of a Petri run\'s writer lease by its owner. + */ +export interface PetriReleaseRequest { + /** + * The owner id that holds the lease. + */ + 'owner': string; +} From 25d47ebcd364874ede8058ae296215222f83ab2a Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 20:12:20 -0400 Subject: [PATCH 2/4] Implement Petri's run store over the server's API `HttpRunStore` is Petri's `RunStore` and `RunLogs` as a worker process reaches them: over `fabro_client::Client` with the worker's token, against the server's SQLite store. A run key is a Fabro run id, the `{id}` of every request, which is the plan's rule that Petri's run key is Fabro's run id. The lease rules are the store's. A same-owner reopen shares the live handle in the process, and the server makes a same-owner reopen after a lost reply the same lease. Dropping the last handle of an owner sends `release` on the current Tokio runtime, and the store awaits every such release before its next open, so a drop followed by an open observes it. The server's worker-exit release is the backstop. A reply that never arrives, a transport error or the client's request timeout, is retried by resending the same request up to three times. Every request is idempotent on the server, so that is safe; a reply that did arrive is never retried. Each server error code maps back to its `StoreError` variant, with the leased owner and the conflict position read from `meta`. `fabro_petri::petri` re-exports the store vocabulary for the server, and the `test-support` feature re-exports Petri's test kit so the server's tests can run the conformance suite over the wire. Both keep this crate the one place that names a Petri package. Co-Authored-By: Claude Fable 5.1 --- lib/components/fabro-petri/Cargo.toml | 11 + lib/components/fabro-petri/README.md | 14 + lib/components/fabro-petri/src/http_store.rs | 621 ++++++++++++++++++ lib/components/fabro-petri/src/lib.rs | 7 + lib/components/fabro-petri/src/petri.rs | 8 + .../fabro-petri/src/test_support.rs | 5 + 6 files changed, 666 insertions(+) create mode 100644 lib/components/fabro-petri/src/http_store.rs create mode 100644 lib/components/fabro-petri/src/petri.rs create mode 100644 lib/components/fabro-petri/src/test_support.rs diff --git a/lib/components/fabro-petri/Cargo.toml b/lib/components/fabro-petri/Cargo.toml index 9398c5481..70d8187fa 100644 --- a/lib/components/fabro-petri/Cargo.toml +++ b/lib/components/fabro-petri/Cargo.toml @@ -12,7 +12,15 @@ doctest = false [lints] workspace = true +[features] +# Petri's test kit, re-exported for Fabro crates that check a store +# implementation against Petri's contract from their own tests. Never on in +# a normal build. +test-support = ["dep:petri_testkit"] + [dependencies] +fabro-api = { path = "../../foundation/fabro-api" } +fabro-client = { path = "../../foundation/fabro-client" } fabro-db = { path = "../../foundation/fabro-db" } fabro-store = { path = "../fabro-store" } fabro-types = { path = "../../foundation/fabro-types" } @@ -22,6 +30,8 @@ petri_store.workspace = true petri_attractor_steps.workspace = true petri_frontend_attractor.workspace = true petri_frontend_fabro.workspace = true +petri_testkit = { workspace = true, optional = true } +anyhow.workspace = true async-trait.workspace = true serde_json.workspace = true sqlx.workspace = true @@ -29,6 +39,7 @@ tokio.workspace = true tracing.workspace = true [dev-dependencies] +fabro-http.workspace = true fabro-store = { path = "../fabro-store", features = ["test-support"] } petri_testkit.workspace = true tempfile = "3" diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 727afb430..628dba833 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -19,6 +19,15 @@ Every adapter the integration plan describes lands here. run and its writer lease, `petri_records` for every record of every log, and the shared `blobs` table). The module docs state the lease and append rules. +- `HttpRunStore`: the same store as a run's worker process reaches it, over + the server's `/api/v1/runs/{id}/petri/*` endpoints with the worker's token. + The server answers from its `SqliteRunStore`, so the lease and the + `(log, seq)` rule are the store's; this layer carries requests, resends a + request whose reply was lost, and maps the server's error codes back to + `StoreError`. The module docs state the rules. +- `petri`: the Petri store vocabulary re-exported for the server, which + answers the worker endpoints from a `SqliteRunStore` without naming a Petri + package in its own manifest. - The platform adapters the plan adds after it: hooks, interviews, secrets, output storage, run tools, the event projection. @@ -37,6 +46,11 @@ Integration tests live under `tests/`: operator release, lease exclusivity, a crash between appends, and blob interoperation with Fabro's `BlobStore`. +The conformance suite over `HttpRunStore` needs a server to talk to, so it +lives with the server's integration tests +(`lib/apps/fabro-server/tests/it/api/petri_store.rs`), which reach the suite +through this crate's `test-support` feature (`fabro_petri::test_support`). + Run them with: ```sh diff --git a/lib/components/fabro-petri/src/http_store.rs b/lib/components/fabro-petri/src/http_store.rs new file mode 100644 index 000000000..d213706ed --- /dev/null +++ b/lib/components/fabro-petri/src/http_store.rs @@ -0,0 +1,621 @@ +//! Petri's run store as a run's worker process reaches it: the Petri store +//! crate's `RunStore` and `RunLogs` over the Fabro server's API, with the +//! worker's token. The server side is [`SqliteRunStore`](crate::SqliteRunStore) +//! behind the `/api/v1/runs/{id}/petri/*` endpoints, so the lease and the +//! `(log, seq)` rule are the store's: this layer carries requests and maps +//! replies. +//! +//! # Keys +//! +//! A run key over the API is a Fabro run id, the `{id}` of every endpoint. +//! The worker's token names the one run it may reach; any other key is +//! refused by the server. That is the integration plan's rule that Petri's +//! `run_key` is Fabro's run id. +//! +//! # The lease +//! +//! `Create` and `Write` take the run's writer lease for the handle's +//! `OwnerId` on the server, idempotently: a same-owner reopen in this +//! process shares the live handle, and a same-owner reopen after a lost +//! reply gets the same lease from the server. Another live owner is refused +//! with `Leased`. The lease ends when the last handle of the owner drops +//! (the drop sends `release`, best effort, on the current Tokio runtime), +//! when the server observes the worker exit, or by operator release. Never +//! by timeout. The store awaits every spawned release before its next +//! `open`, so a drop followed by an open observes the release; with no +//! runtime at drop, the server's worker-exit release is the backstop, and +//! the drop says so in the log. +//! +//! # Lost replies +//! +//! Every call is one request. A reply that never arrives (a transport error, +//! or the client's request timeout) is retried by resending the same +//! request a bounded number of times, with [`LOST_REPLY_RETRY_DELAYS`] +//! between attempts. That is safe because every request is idempotent on the +//! server: a repeated record at a taken seq is accepted, a blob write is +//! content-addressed, an open by the owner that holds the lease shares it, +//! and a release by an owner that no longer holds the lease is a no-op. A +//! `Create` whose reply was lost may find the run exists on the resend; the +//! store then takes the run with `Write` for the same owner, which the lease +//! rule makes the same lease. A reply that did arrive is never retried: the +//! server's answer, error or not, is the store's answer. +//! +//! # Errors +//! +//! The server names each store error with a machine-readable `code` +//! (`petri_run_exists`, `petri_run_not_found`, `petri_run_leased` with the +//! holder under `meta.owner`, `petri_stale_owner`, `petri_read_only`, +//! `petri_record_conflict` with the position under `meta.log` and +//! `meta.seq`), and the store maps each back to its `StoreError` variant. +//! Anything else, including a lost reply after the last retry, is +//! `StoreError::Backend`. + +use std::collections::HashMap; +use std::future::Future; +use std::sync::{Arc, Mutex, MutexGuard, PoisonError, Weak}; +use std::time::Duration; +use std::{fmt, mem, ptr}; + +use fabro_api::types::{PetriAccess, PetriAppendRequest, PetriOpenRequest, PetriRecord}; +use fabro_client::{Client, api_failure_for}; +use fabro_types::{BlobHash, RunId}; +use petri_store::{Access, Digest, LogId, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError}; +use serde_json::Value; +use tokio::runtime::Handle; +use tokio::task::JoinHandle; +use tokio::time; +use tracing::{debug, warn}; + +use crate::run_store::{log_id_text, parse_log_id}; + +/// The waits between attempts when a reply is lost: one request, then up +/// to three resends. +pub const LOST_REPLY_RETRY_DELAYS: [Duration; 3] = [ + Duration::from_millis(100), + Duration::from_millis(500), + Duration::from_secs(2), +]; + +/// Petri's run store over the Fabro server's API. +pub struct HttpRunStore { + shared: Arc, +} + +impl fmt::Debug for HttpRunStore { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("HttpRunStore") + .field("server", &self.shared.client.base_url()) + .finish_non_exhaustive() + } +} + +/// What the store and every handle it opens share. +struct Shared { + client: Client, + /// The writer handle alive in this process per run and owner, so a + /// same-owner reopen shares it and the lease lasts while any handle + /// does. + live: Mutex>>, + /// The releases dropped handles spawned, awaited before the next open. + releases: Mutex>>, +} + +impl HttpRunStore { + /// A store over a client that carries the worker's token. + #[must_use] + pub fn new(client: Client) -> Self { + Self { + shared: Arc::new(Shared { + client, + live: Mutex::default(), + releases: Mutex::default(), + }), + } + } + + /// The writer handle for `owner`, once the lease is taken: the live one + /// when this owner already holds a handle here, else a new one. + fn writer( + &self, + key: &RunKey, + run_id: RunId, + owner: OwnerId, + locator: String, + ) -> Arc { + let mut live = lock(&self.shared.live); + let slot = (key.clone(), owner.clone()); + if let Some(handle) = live.get(&slot).and_then(Weak::upgrade) { + return handle; + } + let handle = Arc::new(HttpRunLogs { + shared: self.shared.clone(), + run_id, + key: key.clone(), + owner: Some(owner), + locator, + }); + live.insert(slot, Arc::downgrade(&handle)); + handle + } +} + +impl Shared { + /// Where a run lives, for messages before the server has said. + fn locator(&self, key: &RunKey) -> String { + format!("Fabro server {}, run `{key}`", self.client.base_url()) + } + + /// The Fabro run id a key names, which is the `{id}` of every request. + fn run_id(&self, key: &RunKey) -> Result { + key.as_str().parse::().map_err(|cause| { + StoreError::backend( + self.locator(key), + "name the run", + format!("a Petri run key over Fabro's API is a Fabro run id: {cause}"), + ) + }) + } + + /// Await every release a dropped handle spawned, so what follows sees + /// the lease as the drops left it. + async fn drain_releases(&self) { + let pending = mem::take(&mut *lock(&self.releases)); + for release in pending { + // A release task never panics: it reports its own failure. + let _ = release.await; + } + } + + /// Run `request` until a reply arrives: a lost reply is resent after + /// each of [`LOST_REPLY_RETRY_DELAYS`], and the last loss is the error. + async fn until_replied( + &self, + key: &RunKey, + action: &'static str, + mut request: F, + ) -> Result + where + F: FnMut() -> Fut, + Fut: Future>, + { + let mut attempt = 0; + loop { + match request().await { + Ok(value) => return Ok(value), + Err(error) if reply_lost(&error) && attempt < LOST_REPLY_RETRY_DELAYS.len() => { + let delay = LOST_REPLY_RETRY_DELAYS[attempt]; + attempt += 1; + warn!( + run_id = %key, + action, + attempt, + delay_ms = delay.as_millis(), + error = %error, + "Petri store reply lost; resending the request" + ); + time::sleep(delay).await; + } + Err(error) => return Err(error), + } + } + } + + /// Map a request's failure to the store error the server named, or a + /// backend error for anything else. `log` is the request's log, the + /// fallback for a conflict whose `meta` does not name one. + fn store_error( + &self, + key: &RunKey, + action: &'static str, + log: Option<&LogId>, + error: anyhow::Error, + ) -> StoreError { + let Some(failure) = api_failure_for(&error) else { + return StoreError::backend(self.locator(key), action, error); + }; + let meta = |member: &str| { + failure + .meta + .as_ref() + .and_then(|meta| meta.get(member).cloned()) + }; + match failure.code.as_deref() { + Some("petri_run_exists") => StoreError::Exists { + key: key.clone(), + locator: self.locator(key), + }, + Some("petri_run_not_found") => StoreError::NotFound { + key: key.clone(), + locator: self.locator(key), + }, + Some("petri_run_leased") => StoreError::Leased { + locator: self.locator(key), + owner: OwnerId::new( + meta("owner") + .and_then(|owner| owner.as_str().map(ToOwned::to_owned)) + .unwrap_or_else(|| "".to_string()), + ), + }, + Some("petri_stale_owner") => StoreError::StaleOwner, + Some("petri_read_only") => StoreError::ReadOnly, + Some("petri_record_conflict") => { + let named = meta("log") + .and_then(|log| log.as_str().and_then(parse_log_id)) + .or_else(|| log.copied()); + let seq = meta("seq").and_then(|seq| seq.as_u64()); + match (named, seq) { + (Some(log), Some(seq)) => StoreError::Conflict { log, seq }, + _ => StoreError::backend(self.locator(key), action, error), + } + } + _ => StoreError::backend(self.locator(key), action, error), + } + } + + /// End the lease of `key` for `owner` on the server: what a dropped + /// handle does. A lease that already moved is left alone by the server. + async fn release(&self, key: &RunKey, run_id: RunId, owner: &OwnerId) { + let released = self + .until_replied(key, "release the run's lease", || { + self.client.release_petri_run(&run_id, owner.as_str()) + }) + .await; + match released { + Ok(()) => debug!(run_id = %key, owner = %owner, "Petri run lease released at drop"), + Err(error) => warn!( + run_id = %key, + owner = %owner, + error = %error, + "Petri run lease not released at drop; the server releases it when the worker exits" + ), + } + } +} + +/// Whether a request failed without a reply from the server: a transport +/// error or the client's request timeout. A reply, whatever its status, +/// carries an `ApiFailure`. +fn reply_lost(error: &anyhow::Error) -> bool { + api_failure_for(error).is_none() +} + +#[async_trait::async_trait] +impl RunStore for HttpRunStore { + async fn open(&self, key: &RunKey, access: Access) -> Result, StoreError> { + let shared = &self.shared; + shared.drain_releases().await; + let run_id = shared.run_id(key)?; + let (api_access, owner) = match &access { + Access::Create { owner } => (PetriAccess::Create, Some(owner)), + Access::Write { owner } => (PetriAccess::Write, Some(owner)), + Access::Read => (PetriAccess::Read, None), + }; + let request = PetriOpenRequest { + access: api_access, + owner: owner.map(|owner| owner.as_str().to_string()), + }; + let mut sent = 0_usize; + let opened = shared + .until_replied(key, "open the run", || { + sent += 1; + shared.client.open_petri_run(&run_id, request.clone()) + }) + .await; + let opened = match (opened, owner) { + // A `Create` whose first reply was lost may have created the + // run: the resend finds it and takes it as its owner. + (Err(error), Some(owner)) + if sent > 1 + && matches!(request.access, PetriAccess::Create) + && api_failure_for(&error).is_some_and(|failure| { + failure.code.as_deref() == Some("petri_run_exists") + }) => + { + debug!(run_id = %key, owner = %owner, "Petri run created by a resend; taking it"); + let retake = PetriOpenRequest { + access: PetriAccess::Write, + owner: Some(owner.as_str().to_string()), + }; + shared + .until_replied(key, "open the run", || { + shared.client.open_petri_run(&run_id, retake.clone()) + }) + .await + } + (opened, _) => opened, + }; + let opened = + opened.map_err(|error| shared.store_error(key, "open the run", None, error))?; + debug!(run_id = %key, access = ?request.access, "Petri run opened over the API"); + match owner { + Some(owner) => Ok(self.writer(key, run_id, owner.clone(), opened.locator)), + None => Ok(Arc::new(HttpRunLogs { + shared: shared.clone(), + run_id, + key: key.clone(), + owner: None, + locator: opened.locator, + })), + } + } +} + +/// One run on the server, opened. A writer handle carries the owner it was +/// opened with; a reader handle refuses every mutation. +struct HttpRunLogs { + shared: Arc, + run_id: RunId, + key: RunKey, + owner: Option, + locator: String, +} + +impl HttpRunLogs { + fn owner(&self) -> Result<&OwnerId, StoreError> { + self.owner.as_ref().ok_or(StoreError::ReadOnly) + } + + fn backend( + &self, + action: &'static str, + cause: impl Into>, + ) -> StoreError { + StoreError::backend(self.locator.clone(), action, cause) + } + + /// A record as the wire carries it: its JSON must be an object. + fn wire(&self, record: &Record) -> Result { + match &record.record { + Value::Object(map) => Ok(PetriRecord { + seq: record.seq, + recorded_at: record.recorded_at, + record: map.clone(), + }), + _ => Err(self.backend( + "encode a record", + format!("record at seq {} is not a JSON object", record.seq), + )), + } + } + + /// A wire record back into the record it was, checked against the seq + /// and recorded_at the server lifted beside it. + fn decode(&self, wire: PetriRecord) -> Result { + let record = Record::from_value(Value::Object(wire.record)) + .map_err(|cause| self.backend("decode a stored record", cause))?; + if record.seq != wire.seq || record.recorded_at != wire.recorded_at { + return Err(self.backend( + "decode a stored record", + format!( + "the server lifted seq {} and recorded_at {} beside a record carrying seq {} and recorded_at {}", + wire.seq, wire.recorded_at, record.seq, record.recorded_at + ), + )); + } + Ok(record) + } +} + +impl Drop for HttpRunLogs { + fn drop(&mut self) { + let Some(owner) = self.owner.clone() else { + return; + }; + { + let mut live = lock(&self.shared.live); + let slot = (self.key.clone(), owner.clone()); + let this: *const Self = self; + if live + .get(&slot) + .is_some_and(|weak| ptr::eq(weak.as_ptr(), this)) + { + live.remove(&slot); + } + } + match Handle::try_current() { + Ok(runtime) => { + let shared = self.shared.clone(); + let key = self.key.clone(); + let run_id = self.run_id; + let release = runtime.spawn(async move { + shared.release(&key, run_id, &owner).await; + }); + lock(&self.shared.releases).push(release); + } + Err(_) => { + warn!( + run_id = %self.key, + owner = %owner, + "Petri run lease not released at drop: no async runtime; the server releases it when the worker exits" + ); + } + } + } +} + +#[async_trait::async_trait] +impl RunLogs for HttpRunLogs { + fn locator(&self) -> String { + self.locator.clone() + } + + async fn append(&self, log: &LogId, records: &[Record]) -> Result<(), StoreError> { + let owner = self.owner()?; + let body = PetriAppendRequest { + owner: owner.as_str().to_string(), + records: records + .iter() + .map(|record| self.wire(record)) + .collect::, _>>()?, + }; + let log_text = log_id_text(log); + self.shared + .until_replied(&self.key, "append records", || { + self.shared + .client + .append_petri_records(&self.run_id, &log_text, body.clone()) + }) + .await + .map_err(|error| { + self.shared + .store_error(&self.key, "append records", Some(log), error) + }) + } + + async fn read(&self, log: &LogId) -> Result, StoreError> { + let log_text = log_id_text(log); + let records = self + .shared + .until_replied(&self.key, "read a log", || { + self.shared + .client + .list_petri_records(&self.run_id, &log_text) + }) + .await + .map_err(|error| { + self.shared + .store_error(&self.key, "read a log", Some(log), error) + })?; + records.into_iter().map(|wire| self.decode(wire)).collect() + } + + async fn put_blob(&self, bytes: &[u8]) -> Result { + let owner = self.owner()?; + let hash = self + .shared + .until_replied(&self.key, "store a blob", || { + self.shared + .client + .write_petri_blob(&self.run_id, owner.as_str(), bytes) + }) + .await + .map_err(|error| { + self.shared + .store_error(&self.key, "store a blob", None, error) + })?; + hash.to_string() + .parse() + .map_err(|cause| self.backend("store a blob", cause)) + } + + async fn get_blob(&self, digest: Digest) -> Result>, StoreError> { + let hash: BlobHash = digest + .to_hex() + .parse() + .map_err(|cause| self.backend("read a blob", cause))?; + let bytes = self + .shared + .until_replied(&self.key, "read a blob", || { + self.shared.client.read_petri_blob(&self.run_id, &hash) + }) + .await + .map_err(|error| { + self.shared + .store_error(&self.key, "read a blob", None, error) + })?; + Ok(bytes.map(|bytes| bytes.to_vec())) + } +} + +fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { + mutex.lock().unwrap_or_else(PoisonError::into_inner) +} + +#[cfg(test)] +mod tests { + use fabro_client::ApiFailure; + use fabro_client::error::tag_with_failure; + use petri_store::ExecutionId; + use serde_json::json; + + use super::*; + + fn store() -> HttpRunStore { + HttpRunStore::new(Client::new_no_proxy("http://127.0.0.1:1").expect("a client builds")) + } + + fn failure(status: u16, code: &str, meta: Option) -> anyhow::Error { + tag_with_failure(anyhow::anyhow!("the server said no"), ApiFailure { + status: fabro_http::StatusCode::from_u16(status).expect("a status"), + code: Some(code.to_string()), + meta, + }) + } + + #[test] + fn the_server_codes_map_back_to_the_store_errors() { + let store = store(); + let key = RunKey::new("01ARZ3NDEKTSV4RRFFQ69G5FAV"); + let log = LogId::Execution(ExecutionId::new(2)); + let map = |error| store.shared.store_error(&key, "act", Some(&log), error); + + assert!(matches!( + map(failure(409, "petri_run_exists", None)), + StoreError::Exists { .. } + )); + assert!(matches!( + map(failure(404, "petri_run_not_found", None)), + StoreError::NotFound { .. } + )); + assert!(matches!( + map(failure(409, "petri_run_leased", Some(json!({"owner": "first"})))), + StoreError::Leased { owner, .. } if owner.as_str() == "first" + )); + assert!(matches!( + map(failure(409, "petri_stale_owner", None)), + StoreError::StaleOwner + )); + assert!(matches!( + map(failure(409, "petri_read_only", None)), + StoreError::ReadOnly + )); + assert!(matches!( + map(failure( + 409, + "petri_record_conflict", + Some(json!({"log": "resources", "seq": 4})) + )), + StoreError::Conflict { + log: LogId::Resources, + seq: 4, + } + )); + assert!( + matches!( + map(failure(409, "petri_record_conflict", Some(json!({"seq": 1})))), + StoreError::Conflict { log: named, seq: 1 } if named == log + ), + "a conflict that names no log is on the request's log" + ); + assert!(matches!( + map(failure(500, "petri_store_failed", None)), + StoreError::Backend { .. } + )); + assert!(matches!( + map(anyhow::anyhow!("connection reset")), + StoreError::Backend { .. } + )); + } + + #[test] + fn a_reply_is_lost_only_without_an_api_failure() { + assert!(reply_lost(&anyhow::anyhow!("server request timed out"))); + assert!(!reply_lost(&failure(409, "petri_stale_owner", None))); + } + + #[test] + fn a_key_over_the_api_is_a_fabro_run_id() { + let store = store(); + let error = store + .shared + .run_id(&RunKey::new("lifecycle")) + .expect_err("a plain word is not a run id"); + assert!(matches!(error, StoreError::Backend { .. }), "{error}"); + assert!( + store + .shared + .run_id(&RunKey::new("01ARZ3NDEKTSV4RRFFQ69G5FAV")) + .is_ok() + ); + } +} diff --git a/lib/components/fabro-petri/src/lib.rs b/lib/components/fabro-petri/src/lib.rs index dd3f77139..677b50cf2 100644 --- a/lib/components/fabro-petri/src/lib.rs +++ b/lib/components/fabro-petri/src/lib.rs @@ -10,12 +10,19 @@ //! //! - [`SqliteRunStore`]: Petri's run store over Fabro's SQLite database, so a //! run's records are its source of truth in Fabro's tables; +//! - [`HttpRunStore`]: the same store as a run's worker process reaches it, +//! over the server's API with the worker's token; //! - the platform adapters: hooks, interviews, secrets, output storage, the run //! tools, the event projection. //! //! The Petri packages are pinned by revision in the workspace `Cargo.toml` //! under `petri_*` keys. +pub mod http_store; +pub mod petri; pub mod run_store; +#[cfg(feature = "test-support")] +pub mod test_support; +pub use http_store::HttpRunStore; pub use run_store::SqliteRunStore; diff --git a/lib/components/fabro-petri/src/petri.rs b/lib/components/fabro-petri/src/petri.rs new file mode 100644 index 000000000..726edac63 --- /dev/null +++ b/lib/components/fabro-petri/src/petri.rs @@ -0,0 +1,8 @@ +//! Petri's store vocabulary, re-exported for the Fabro crates that hold a +//! Petri run handle or answer for one (the server's worker endpoints) without +//! depending on the Petri packages themselves. Only this crate names them in +//! its `Cargo.toml`. + +pub use petri_store::{ + Access, Digest, ExecutionId, LogId, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError, +}; diff --git a/lib/components/fabro-petri/src/test_support.rs b/lib/components/fabro-petri/src/test_support.rs new file mode 100644 index 000000000..8f677c1d0 --- /dev/null +++ b/lib/components/fabro-petri/src/test_support.rs @@ -0,0 +1,5 @@ +//! Petri's test kit, for Fabro crates that check a store implementation +//! against Petri's contract from their own tests. Compiled only with the +//! `test-support` feature, which a dev-dependency turns on. + +pub use petri_testkit::run_store; From 2ca1e1cda049847ed8821b781bd0da6418bca10e Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 20:12:20 -0400 Subject: [PATCH 3/4] 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() +} From 6557a2f40460008f06ee616f7cc95387ca894974 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 20:12:20 -0400 Subject: [PATCH 4/4] Check the HTTP run store against a loopback server Petri's store conformance suite runs over `HttpRunStore` talking to an axum listener on a loopback port. The suite opens runs under keys of its own, while a key over the API is a Fabro run id the worker's token names, so an adapter gives each suite key a fresh run with a token minted for that run alone: the least a worker holds. Three more tests cover what the suite cannot: the operator release through the server's store turns the worker's handle stale; a middleware swallows the reply of one committed append and the store's resend leaves each record once; and two workers with owners of their own never hold one run's lease at the same time. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-server/tests/it/api/mod.rs | 1 + .../fabro-server/tests/it/api/petri_store.rs | 307 ++++++++++++++++++ 2 files changed, 308 insertions(+) create mode 100644 lib/apps/fabro-server/tests/it/api/petri_store.rs diff --git a/lib/apps/fabro-server/tests/it/api/mod.rs b/lib/apps/fabro-server/tests/it/api/mod.rs index 7fcdf254a..8833cfa21 100644 --- a/lib/apps/fabro-server/tests/it/api/mod.rs +++ b/lib/apps/fabro-server/tests/it/api/mod.rs @@ -8,6 +8,7 @@ mod events; mod install; mod install_openai_compatible; mod mcp_servers; +mod petri_store; mod routing; mod run_files; mod runs; diff --git a/lib/apps/fabro-server/tests/it/api/petri_store.rs b/lib/apps/fabro-server/tests/it/api/petri_store.rs new file mode 100644 index 000000000..9c8d6d157 --- /dev/null +++ b/lib/apps/fabro-server/tests/it/api/petri_store.rs @@ -0,0 +1,307 @@ +//! The Petri run store over HTTP: Petri's store conformance suite against +//! `HttpRunStore` talking to a loopback server, the operator release through +//! the server's store, a lost reply to an append, and two workers contending +//! for one run's lease. + +use std::collections::HashMap; +use std::future; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use axum::Router; +use axum::extract::{Request, State}; +use axum::http::{Method, StatusCode}; +use axum::middleware::{self, Next}; +use axum::response::Response; +use fabro_client::{Client, Credential, apply_bearer_token_auth}; +use fabro_petri::HttpRunStore; +use fabro_petri::petri::{Access, LogId, OwnerId, RunKey, RunLogs, RunStore, StoreError}; +use fabro_petri::test_support::run_store::{self, conformance, stale_owner_conformance}; +use fabro_server::server::{AppState, RouterOptions, build_router_with_options}; +use fabro_server::test_support::{test_app_state, test_auth_mode}; +use fabro_types::RunId; +use tokio::net::TcpListener; +use tokio::runtime::Handle; +use tokio::task; + +/// A loopback server over `state`, its router wrapped by `wrap`. +async fn serve(state: Arc, wrap: impl FnOnce(Router) -> Router) -> String { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("a loopback listener binds"); + let addr = listener.local_addr().expect("the listener has an address"); + let router = wrap(build_router_with_options( + state, + &test_auth_mode(), + RouterOptions::default(), + )); + tokio::spawn(async move { + let _ = axum::serve(listener, router).await; + }); + format!("http://{addr}") +} + +/// A worker's client: the run's worker token as its bearer, and a request +/// timeout after which a reply counts as lost. +async fn worker_client(base_url: &str, token: &str, timeout: Duration) -> Client { + let http = apply_bearer_token_auth(fabro_http::HttpClientBuilder::new().no_proxy(), token) + .expect("the bearer header builds") + .build() + .expect("the test HTTP client builds"); + Client::builder() + .transport(base_url, http) + .credential(Credential::Worker(token.to_string())) + .request_timeout(timeout) + .connect() + .await + .expect("the worker client connects") +} + +/// The Petri key of a Fabro run: its id. +fn petri_key(run_id: RunId) -> RunKey { + RunKey::new(run_id.to_string()) +} + +/// Wait until the server's store shows no lease holder on `key`: a dropped +/// handle's release travels to the server on its own, and only the store +/// that dropped it awaits that. Another worker starts after the first is +/// gone, which is what this waits for. +async fn wait_until_released(store: &fabro_petri::SqliteRunStore, key: &RunKey) { + for _ in 0..500 { + if store.owner(key).await.expect("reads the lease").is_none() { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + panic!("the lease on {key} was not released"); +} + +/// One worker's view of one run: a store over a client whose token names +/// the run. +async fn worker_store(state: &AppState, base_url: &str, run_id: RunId) -> HttpRunStore { + let token = state.test_issue_worker_token(&run_id); + HttpRunStore::new(worker_client(base_url, &token, Duration::from_secs(5)).await) +} + +/// The suite opens runs under keys of its own (`lifecycle`, `drop`), while a +/// key over the API is a Fabro run id that the worker's token names. This +/// adapter gives each suite key a Fabro run of its own, with a token minted +/// for that run alone: the least a worker holds. +struct WorkerRuns { + state: Arc, + base_url: String, + runs: Mutex>>, +} + +struct WorkerRun { + run_id: RunId, + store: HttpRunStore, +} + +impl WorkerRuns { + fn new(state: Arc, base_url: String) -> Self { + Self { + state, + base_url, + runs: Mutex::default(), + } + } + + async fn run(&self, key: &RunKey) -> Arc { + if let Some(run) = self.runs.lock().expect("runs lock").get(key) { + return Arc::clone(run); + } + let run_id = RunId::new(); + let store = worker_store(&self.state, &self.base_url, run_id).await; + let run = Arc::new(WorkerRun { run_id, store }); + self.runs + .lock() + .expect("runs lock") + .insert(key.clone(), Arc::clone(&run)); + run + } + + /// The Fabro run a suite key was given, once opened. + fn run_id(&self, key: &RunKey) -> RunId { + self.runs + .lock() + .expect("runs lock") + .get(key) + .expect("the suite opens a key before it releases it") + .run_id + } +} + +#[async_trait::async_trait] +impl RunStore for WorkerRuns { + async fn open(&self, key: &RunKey, access: Access) -> Result, StoreError> { + let run = self.run(key).await; + run.store.open(&petri_key(run.run_id), access).await + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn the_http_store_passes_the_conformance_suite() { + let state = test_app_state(); + let base_url = serve(Arc::clone(&state), |router| router).await; + let runs: Arc = Arc::new(WorkerRuns::new(state, base_url)); + conformance(|| Arc::clone(&runs)).await; +} + +/// An operator release through the server's own store ends the worker's +/// lease from outside: the worker's handle turns stale and the next writer +/// takes the run. The release is async, so the suite's synchronous closure +/// blocks on it in place. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_operator_release_through_the_server_makes_the_worker_stale() { + let state = test_app_state(); + let base_url = serve(Arc::clone(&state), |router| router).await; + let runs = WorkerRuns::new(Arc::clone(&state), base_url); + let release = |key: &RunKey| { + let key = petri_key(runs.run_id(key)); + let store = state.test_petri_run_store(); + task::block_in_place(|| Handle::current().block_on(store.release_lease(&key))) + .expect("the operator releases the lease"); + }; + stale_owner_conformance(&runs, release).await; +} + +/// Swallow the reply of the first append the server commits: the handler +/// runs, its transaction commits, and the response never leaves. That is +/// a lost reply as the worker sees it. +async fn swallow_one_append_reply( + State(swallowed): State>, + request: Request, + next: Next, +) -> Response { + let is_append = request.method() == Method::POST && request.uri().path().ends_with("/records"); + let response = next.run(request).await; + if is_append + && response.status() == StatusCode::NO_CONTENT + && swallowed + .compare_exchange(0, 1, Ordering::SeqCst, Ordering::SeqCst) + .is_ok() + { + future::pending::<()>().await; + } + response +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_lost_reply_to_an_append_is_safe_to_retry() { + let state = test_app_state(); + let swallowed = Arc::new(AtomicUsize::new(0)); + let base_url = serve(Arc::clone(&state), { + let swallowed = Arc::clone(&swallowed); + move |router| { + router.layer(middleware::from_fn_with_state( + swallowed, + swallow_one_append_reply, + )) + } + }) + .await; + let run_id = RunId::new(); + let token = state.test_issue_worker_token(&run_id); + let timeout = Duration::from_millis(500); + let store = HttpRunStore::new(worker_client(&base_url, &token, timeout).await); + let key = petri_key(run_id); + let logs = store + .open(&key, Access::Create { + owner: OwnerId::new("worker"), + }) + .await + .expect("the worker creates the run"); + + let batch = [ + run_store::record(0, "execution.started"), + run_store::record(1, "step.started"), + ]; + let started = Instant::now(); + logs.append(&LogId::Coordinator, &batch) + .await + .expect("the append succeeds once the resend is answered"); + assert_eq!( + swallowed.load(Ordering::SeqCst), + 1, + "the server committed one append whose reply was swallowed" + ); + assert!( + started.elapsed() >= timeout, + "the first attempt waited out the request timeout" + ); + assert_eq!( + logs.read(&LogId::Coordinator).await.expect("reads"), + batch.to_vec(), + "the log holds each record once" + ); +} + +/// Two workers with owners of their own never hold one run's lease at the +/// same time, in either order, and the server's store reports who holds it. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn two_workers_cannot_both_hold_a_run_lease() { + let state = test_app_state(); + let base_url = serve(Arc::clone(&state), |router| router).await; + let run_id = RunId::new(); + let key = petri_key(run_id); + let first_worker = worker_store(&state, &base_url, run_id).await; + let second_worker = worker_store(&state, &base_url, run_id).await; + let first = OwnerId::new("first"); + let second = OwnerId::new("second"); + let server_store = state.test_petri_run_store(); + + let held = first_worker + .open(&key, Access::Create { + owner: first.clone(), + }) + .await + .expect("the first worker creates"); + assert_eq!( + server_store.owner(&key).await.expect("reads"), + Some(first.clone()) + ); + let refused = second_worker + .open(&key, Access::Write { + owner: second.clone(), + }) + .await + .err() + .expect("the second worker is refused while the first holds the lease"); + assert!( + matches!(&refused, StoreError::Leased { owner, .. } if *owner == first), + "{refused}" + ); + assert!( + refused.to_string().contains(&run_id.to_string()), + "the message names the run: {refused}" + ); + + drop(held); + wait_until_released(server_store, &key).await; + let taken = second_worker + .open(&key, Access::Write { + owner: second.clone(), + }) + .await + .expect("the second worker takes the run once the first releases"); + assert_eq!( + server_store.owner(&key).await.expect("reads"), + Some(second.clone()) + ); + let refused = first_worker + .open(&key, Access::Write { + owner: first.clone(), + }) + .await + .err() + .expect("the first worker is refused in turn"); + assert!( + matches!(&refused, StoreError::Leased { owner, .. } if *owner == second), + "{refused}" + ); + + drop(taken); + wait_until_released(server_store, &key).await; +}