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