mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
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 <n>`), 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 <noreply@anthropic.com>
This commit is contained in:
parent
55a145f5a4
commit
0820252bcf
17 changed files with 1384 additions and 47 deletions
|
|
@ -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 <n>` 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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
meta: Option<Box<Map<String, Value>>>,
|
||||
}
|
||||
|
||||
#[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<String>,
|
||||
/// 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<Box<Map<String, Value>>>,
|
||||
}
|
||||
|
||||
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<String>,
|
||||
code: impl Into<String>,
|
||||
meta: Map<String, Value>,
|
||||
) -> 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();
|
||||
|
|
|
|||
|
|
@ -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<types::PetriOpenResponse> {
|
||||
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<Vec<types::PetriRecord>> {
|
||||
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<BlobHash> {
|
||||
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<Option<Bytes>> {
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -6,6 +6,28 @@ use serde::de::DeserializeOwned;
|
|||
pub struct ApiFailure {
|
||||
pub status: fabro_http::StatusCode,
|
||||
pub code: Option<String>,
|
||||
/// The error entry's `meta`: structured details specific to `code`.
|
||||
pub meta: Option<serde_json::Value>,
|
||||
}
|
||||
|
||||
impl ApiFailure {
|
||||
#[must_use]
|
||||
pub fn new(status: fabro_http::StatusCode, code: Option<String>) -> 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<String>,
|
||||
pub code: Option<String>,
|
||||
pub meta: Option<serde_json::Value>,
|
||||
}
|
||||
|
||||
// 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<String>, Option<String>) {
|
||||
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<String>,
|
||||
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::<serde_json::Value>(&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::<serde_json::Value>(&body) {
|
||||
let (_, parsed_code) = parse_error_response_value(&value);
|
||||
code = parsed_code;
|
||||
}
|
||||
let entry = serde_json::from_str::<serde_json::Value>(&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::<serde_json::Value>::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(
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<RequestArgs> => {
|
||||
// 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<RequestArgs> => {
|
||||
// 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<RequestArgs> => {
|
||||
// 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<RequestArgs> => {
|
||||
// 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<RequestArgs> => {
|
||||
// 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<RequestArgs> => {
|
||||
// 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<void>> {
|
||||
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<PetriRecordList>> {
|
||||
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<PetriOpenResponse>> {
|
||||
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<File>> {
|
||||
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<void>> {
|
||||
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<WriteBlobResponse>> {
|
||||
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<void> {
|
||||
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<File> {
|
||||
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<PetriRecordList> {
|
||||
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<PaginatedEventList> {
|
||||
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<PetriOpenResponse> {
|
||||
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<void> {
|
||||
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<File> {
|
||||
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<File> {
|
||||
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<void> {
|
||||
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<WorkflowSettings> {
|
||||
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<WriteBlobResponse> {
|
||||
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
|
||||
|
|
|
|||
|
|
@ -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; };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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';
|
||||
|
|
|
|||
27
lib/packages/fabro-api-client/src/models/petri-access.ts
generated
Normal file
27
lib/packages/fabro-api-client/src/models/petri-access.ts
generated
Normal file
|
|
@ -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];
|
||||
29
lib/packages/fabro-api-client/src/models/petri-append-request.ts
generated
Normal file
29
lib/packages/fabro-api-client/src/models/petri-append-request.ts
generated
Normal file
|
|
@ -0,0 +1,29 @@
|
|||
/* tslint:disable */
|
||||
/* eslint-disable */
|
||||
/**
|
||||
* Fabro Run API
|
||||
* HTTP API for managing Fabro workflow run executions.
|
||||
*
|
||||
* The version of the OpenAPI document: 0.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
// 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<PetriRecord>;
|
||||
}
|
||||
29
lib/packages/fabro-api-client/src/models/petri-open-request.ts
generated
Normal file
29
lib/packages/fabro-api-client/src/models/petri-open-request.ts
generated
Normal file
|
|
@ -0,0 +1,29 @@
|
|||
/* tslint:disable */
|
||||
/* eslint-disable */
|
||||
/**
|
||||
* Fabro Run API
|
||||
* HTTP API for managing Fabro workflow run executions.
|
||||
*
|
||||
* The version of the OpenAPI document: 0.2.0
|
||||
*
|
||||
*
|
||||
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
|
||||
* https://openapi-generator.tech
|
||||
* Do not edit the class manually.
|
||||
*/
|
||||
|
||||
|
||||
// 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;
|
||||
}
|
||||
25
lib/packages/fabro-api-client/src/models/petri-open-response.ts
generated
Normal file
25
lib/packages/fabro-api-client/src/models/petri-open-response.ts
generated
Normal file
|
|
@ -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;
|
||||
}
|
||||
25
lib/packages/fabro-api-client/src/models/petri-record-list.ts
generated
Normal file
25
lib/packages/fabro-api-client/src/models/petri-record-list.ts
generated
Normal file
|
|
@ -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<PetriRecord>;
|
||||
}
|
||||
33
lib/packages/fabro-api-client/src/models/petri-record.ts
generated
Normal file
33
lib/packages/fabro-api-client/src/models/petri-record.ts
generated
Normal file
|
|
@ -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; };
|
||||
}
|
||||
25
lib/packages/fabro-api-client/src/models/petri-release-request.ts
generated
Normal file
25
lib/packages/fabro-api-client/src/models/petri-release-request.ts
generated
Normal file
|
|
@ -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;
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue