mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-09-30 01:53:45 +00:00
Serve a Petri run's events as one stream with a stream_seq cursor
`GET /runs/{id}/events` and `GET /runs/{id}/attach` serve a Petri run's
public events and Fabro's platform records as one ordered stream in a
Fabro envelope (`RunStreamItem`: `run_id`, `stream_seq`, `kind`, `id`,
`recorded_at`, `item`), read from the projector's `petri_stream` table.
The cursor is `stream_seq` (`?after=`); the item's own identity (the
Petri `EventId` as `<log>/<seq>/<index>`, or the platform record's seq)
travels beside it for deduplication. A legacy run keeps its envelope on
the same endpoints; the OpenAPI response is the union of the two lists,
and the stream list reports Petri's `EVENT_CONTRACT_VERSION`.
The attached stream follows the projector's commit signal (a wake-up,
with a poll as the fallback) and ends after the platform record of the
run's terminal lifecycle transition, the analog of the legacy stream's
`run.completed`, or a bounded grace after the projection went terminal.
`RunSpec.engine` (`RunEngine`, `PetriAdmission`, `PetriGraphRef`) is
named in the spec and reuses the Rust types. `fabro-client` matches the
union and adds `list_run_stream`, `list_run_stream_page` and
`attach_run_stream`.
A server test attaches to a two-branch parallel run, disconnects once
both branches started, records a platform notice while both branch
scripts run, reconnects from the last `stream_seq`, and checks the
union is the whole stream: every item once, in order, no gap, no
duplicate, the notice between the branch events, and the same as the
paged listing. The Petri scenarios capture their settled projection and
stream as JSON fixtures for the web app under
`FABRO_CAPTURE_PETRI_FIXTURES`.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
a6ac3f120e
commit
523f831a7a
32 changed files with 22243 additions and 65 deletions
3013
apps/fabro-web/app/test-fixtures/petri/command.json
Normal file
3013
apps/fabro-web/app/test-fixtures/petri/command.json
Normal file
File diff suppressed because it is too large
Load diff
4271
apps/fabro-web/app/test-fixtures/petri/gate.json
Normal file
4271
apps/fabro-web/app/test-fixtures/petri/gate.json
Normal file
File diff suppressed because it is too large
Load diff
4387
apps/fabro-web/app/test-fixtures/petri/hello.json
Normal file
4387
apps/fabro-web/app/test-fixtures/petri/hello.json
Normal file
File diff suppressed because it is too large
Load diff
8824
apps/fabro-web/app/test-fixtures/petri/parallel.json
Normal file
8824
apps/fabro-web/app/test-fixtures/petri/parallel.json
Normal file
File diff suppressed because it is too large
Load diff
|
|
@ -2970,23 +2970,40 @@ paths:
|
|||
tags: [Run Internals]
|
||||
summary: List Run Events
|
||||
description: |
|
||||
Returns a paginated JSON list of stored run events. Ascending order
|
||||
uses `since_seq` as an inclusive cursor. Descending order uses
|
||||
Returns a paginated JSON list of the run's events. The shape depends
|
||||
on the engine the run was created for (`RunSpec.engine`).
|
||||
|
||||
For a legacy run (`engine.kind = legacy`): stored run events in the
|
||||
legacy envelope (`PaginatedEventList`). Ascending order uses
|
||||
`since_seq` as an inclusive cursor. Descending order uses
|
||||
`before_seq` as an exclusive cursor and starts at the newest event
|
||||
when `before_seq` is omitted.
|
||||
|
||||
For a Petri run (`engine.kind = petri`): the run stream
|
||||
(`PaginatedRunStreamList`), one ordered delivery of Petri's own
|
||||
`RunEvent`s and Fabro's platform records in the `RunStreamItem`
|
||||
envelope, in `stream_seq` order. The cursor is `after`: the last
|
||||
`stream_seq` the client saw, exclusive; the first page is `after=0`.
|
||||
`since_seq`, `before_seq` and `order` are not accepted for a Petri
|
||||
run. A client that reconnects resumes from its last `stream_seq` and
|
||||
deduplicates by each item's `id`; every item is delivered once, in
|
||||
order, with no gap.
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/RunId"
|
||||
- $ref: "#/components/parameters/SinceSeq"
|
||||
- $ref: "#/components/parameters/EventLimit"
|
||||
- $ref: "#/components/parameters/BeforeSeq"
|
||||
- $ref: "#/components/parameters/EventOrder"
|
||||
- $ref: "#/components/parameters/StreamAfter"
|
||||
responses:
|
||||
"200":
|
||||
description: Paginated list of run events
|
||||
description: Paginated list of run events, in the run engine's envelope
|
||||
content:
|
||||
application/json:
|
||||
schema:
|
||||
$ref: "#/components/schemas/PaginatedEventList"
|
||||
oneOf:
|
||||
- $ref: "#/components/schemas/PaginatedEventList"
|
||||
- $ref: "#/components/schemas/PaginatedRunStreamList"
|
||||
"400":
|
||||
description: Invalid cursor and order combination
|
||||
headers:
|
||||
|
|
@ -3097,10 +3114,25 @@ paths:
|
|||
operationId: attachRunEvents
|
||||
tags: [Run Internals]
|
||||
summary: Attach Run Events
|
||||
description: Opens an ordered server-sent event stream starting at `since_seq`, replaying persisted events and continuing with live updates while the run remains active.
|
||||
description: |
|
||||
Opens an ordered server-sent event stream, replaying persisted items
|
||||
and continuing with live updates while the run remains active. Each
|
||||
`data:` frame is one JSON object in the run engine's envelope.
|
||||
|
||||
For a legacy run the frames are `EventEnvelope`s and the stream
|
||||
starts at `since_seq` (inclusive; the next unseen event when omitted).
|
||||
It ends after `run.completed` or `run.failed`.
|
||||
|
||||
For a Petri run the frames are `RunStreamItem`s and the stream
|
||||
starts after `after` (the last `stream_seq` the client saw; `0`
|
||||
replays the whole run; the next unseen item when omitted). It ends
|
||||
once the run is no longer active and every committed item has been
|
||||
sent. A reconnecting client passes its last `stream_seq` as `after`
|
||||
and deduplicates by `id`.
|
||||
parameters:
|
||||
- $ref: "#/components/parameters/RunId"
|
||||
- $ref: "#/components/parameters/SinceSeq"
|
||||
- $ref: "#/components/parameters/StreamAfter"
|
||||
responses:
|
||||
"200":
|
||||
description: Server-sent event stream
|
||||
|
|
@ -6290,6 +6322,20 @@ components:
|
|||
default: 100
|
||||
example: 100
|
||||
|
||||
StreamAfter:
|
||||
name: after
|
||||
in: query
|
||||
required: false
|
||||
description: |
|
||||
Run stream cursor for a Petri run: the last `stream_seq` the client
|
||||
saw, exclusive. `0` starts at the first item.
|
||||
schema:
|
||||
type: integer
|
||||
format: uint64
|
||||
minimum: 0
|
||||
default: 0
|
||||
example: 42
|
||||
|
||||
QuestionId:
|
||||
name: qid
|
||||
in: path
|
||||
|
|
@ -10869,6 +10915,99 @@ components:
|
|||
meta:
|
||||
$ref: "#/components/schemas/PaginationMeta"
|
||||
|
||||
RunStreamItemKind:
|
||||
description: Which item shape a run stream item carries.
|
||||
type: string
|
||||
enum: [petri, platform]
|
||||
|
||||
RunStreamItem:
|
||||
description: |
|
||||
One item of a Petri run's stream: a Petri `RunEvent` or a Fabro
|
||||
platform record in Fabro's envelope.
|
||||
|
||||
`stream_seq` is the durable per-run delivery sequence the projector
|
||||
assigned when the item's record was committed: dense, strictly
|
||||
increasing within the run, and the cursor for `after`. `id` is the
|
||||
item's own identity, kept beside the cursor so a client deduplicates
|
||||
by it: for a Petri event the `EventId` as `<log>/<seq>/<index>`
|
||||
(`coordinator/3/0`, `execution 1/23/0`); for a platform record its
|
||||
`seq`. Petri's `EventId` is per log and has no platform variant, so
|
||||
it is never the cursor.
|
||||
|
||||
A `petri` item is a Petri `RunEvent` passed through unchanged:
|
||||
`{id: {log, execution?, seq, index}, origin, recorded_at,
|
||||
observed_at?, context: {invocation, execution, parent?}, subject?,
|
||||
record?, derived?}`. Its vocabulary is Petri's public event contract
|
||||
(`crates/core/execution/EVENTS.md` in the Petri repository), not
|
||||
Fabro's: the recorded event's name is `record.body.event`
|
||||
(`<subject>.<verb>`, e.g. `visit.started`, `step.finished`,
|
||||
`run.finished`), a derived view event's is `derived.event`, and the
|
||||
stage a subject names is `(context.execution, subject.firing)` with
|
||||
`subject.node.name` and `subject.visit` as its display label. The
|
||||
server reports the contract version it serves in
|
||||
`PaginatedRunStreamList.event_contract_version`.
|
||||
|
||||
A `platform` item is a stored platform record: `{seq, recorded_at,
|
||||
record: {kind, ...}, position?: {execution, firing}}`. `record.kind`
|
||||
is one of `run.created`, `run.lifecycle`, `run.title`, `run.parent`,
|
||||
`run.archived`, `run.unarchived`, `run.superseded`, `run.notice`,
|
||||
`interview.answered`, `run.branch`, `git.identity`, `checkpoint`,
|
||||
`pull_request.created`, `notification.sent`, `run.paired`.
|
||||
type: object
|
||||
required:
|
||||
- run_id
|
||||
- stream_seq
|
||||
- kind
|
||||
- id
|
||||
- recorded_at
|
||||
- item
|
||||
properties:
|
||||
run_id:
|
||||
type: string
|
||||
stream_seq:
|
||||
type: integer
|
||||
format: uint64
|
||||
minimum: 0
|
||||
description: The delivery sequence; the cursor.
|
||||
kind:
|
||||
$ref: "#/components/schemas/RunStreamItemKind"
|
||||
id:
|
||||
type: string
|
||||
description: The item's own identity, for deduplication.
|
||||
recorded_at:
|
||||
type: integer
|
||||
format: uint64
|
||||
minimum: 0
|
||||
description: Milliseconds since the Unix epoch when the item's record was appended.
|
||||
item:
|
||||
type: object
|
||||
additionalProperties: true
|
||||
description: The Petri `RunEvent` or the stored platform record, unchanged.
|
||||
|
||||
PaginatedRunStreamList:
|
||||
description: |
|
||||
One page of a Petri run's stream, in `stream_seq` order.
|
||||
`event_contract_version` is Petri's `EVENT_CONTRACT_VERSION` the
|
||||
server was built against: the version of the event contract every
|
||||
`petri` item follows.
|
||||
type: object
|
||||
required:
|
||||
- data
|
||||
- meta
|
||||
- event_contract_version
|
||||
properties:
|
||||
data:
|
||||
type: array
|
||||
items:
|
||||
$ref: "#/components/schemas/RunStreamItem"
|
||||
meta:
|
||||
$ref: "#/components/schemas/PaginationMeta"
|
||||
event_contract_version:
|
||||
type: integer
|
||||
format: uint32
|
||||
minimum: 0
|
||||
example: 3
|
||||
|
||||
AppendEventResponse:
|
||||
description: Assigned sequence number for an appended event.
|
||||
type: object
|
||||
|
|
@ -12712,6 +12851,58 @@ components:
|
|||
oneOf:
|
||||
- $ref: "#/components/schemas/ForkSourceRef"
|
||||
- type: "null"
|
||||
engine:
|
||||
$ref: "#/components/schemas/RunEngine"
|
||||
description: |
|
||||
The engine the run was created for, with what it admitted.
|
||||
Absent in a spec written before the field existed, which means
|
||||
the legacy executor.
|
||||
|
||||
RunEngine:
|
||||
description: |
|
||||
The engine a run was created for. `legacy` is the in-process
|
||||
executor; `petri` names the Petri workflow engine and carries what
|
||||
Petri admitted at create time.
|
||||
oneOf:
|
||||
- type: object
|
||||
required: [kind]
|
||||
properties:
|
||||
kind:
|
||||
type: string
|
||||
enum: [legacy]
|
||||
- allOf:
|
||||
- type: object
|
||||
required: [kind]
|
||||
properties:
|
||||
kind:
|
||||
type: string
|
||||
enum: [petri]
|
||||
- $ref: "#/components/schemas/PetriAdmission"
|
||||
|
||||
PetriAdmission:
|
||||
description: |
|
||||
What Petri admitted for a run at create time: the lowered root graph
|
||||
and the pre-lowered child graphs, every one persisted in the blob
|
||||
store before the run exists.
|
||||
type: object
|
||||
required: [graph]
|
||||
properties:
|
||||
graph:
|
||||
$ref: "#/components/schemas/PetriGraphRef"
|
||||
children:
|
||||
type: array
|
||||
items:
|
||||
$ref: "#/components/schemas/PetriGraphRef"
|
||||
|
||||
PetriGraphRef:
|
||||
description: An admitted graph in the blob store, verified by digest on load.
|
||||
type: object
|
||||
required: [blob, digest]
|
||||
properties:
|
||||
blob:
|
||||
$ref: "#/components/schemas/BlobHash"
|
||||
digest:
|
||||
type: string
|
||||
|
||||
UpdateRunParentRequest:
|
||||
type: object
|
||||
|
|
|
|||
|
|
@ -1,12 +1,17 @@
|
|||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use axum::extract::DefaultBodyLimit;
|
||||
use fabro_api::types::PaginatedRunStreamList;
|
||||
use fabro_petri::petri::EVENT_CONTRACT_VERSION;
|
||||
use fabro_types::run_event::MAX_RUN_EVENT_BODY_BYTES;
|
||||
use fabro_types::{
|
||||
RunEventDetailContent, RunEventDetailContentKind, RunEventDetailEnvelope,
|
||||
RunEventDetailResponse,
|
||||
RunEventDetailResponse, RunStreamItem,
|
||||
};
|
||||
use fabro_workflow::event::build_redacted_event_payload;
|
||||
use tokio::sync::broadcast::error::RecvError;
|
||||
use tokio::time::{self, Instant};
|
||||
|
||||
use super::super::{
|
||||
ApiError, AppState, AppendEventResponse, BroadcastStream, Event, EventBody, EventEnvelope,
|
||||
|
|
@ -70,9 +75,12 @@ struct RunEventListParams {
|
|||
#[serde(default)]
|
||||
before_seq: Option<u32>,
|
||||
#[serde(default)]
|
||||
order: EventSequenceOrder,
|
||||
order: Option<EventSequenceOrder>,
|
||||
#[serde(default)]
|
||||
limit: Option<usize>,
|
||||
/// The run stream cursor of a Petri run: the last `stream_seq` seen.
|
||||
#[serde(default)]
|
||||
after: Option<u64>,
|
||||
}
|
||||
|
||||
impl RunEventListParams {
|
||||
|
|
@ -80,12 +88,21 @@ impl RunEventListParams {
|
|||
self.since_seq.unwrap_or(1).max(1)
|
||||
}
|
||||
|
||||
fn order(&self) -> EventSequenceOrder {
|
||||
self.order.unwrap_or_default()
|
||||
}
|
||||
|
||||
fn limit(&self) -> usize {
|
||||
self.limit.unwrap_or(100).clamp(1, 1000)
|
||||
}
|
||||
|
||||
fn cursor_error(&self) -> Option<&'static str> {
|
||||
match self.order {
|
||||
if self.after.is_some() && (self.since_seq.is_some() || self.before_seq.is_some()) {
|
||||
return Some(
|
||||
"after is the run stream cursor and cannot be combined with since_seq or before_seq.",
|
||||
);
|
||||
}
|
||||
match self.order() {
|
||||
EventSequenceOrder::Asc if self.before_seq.is_some() => {
|
||||
Some("before_seq requires order=desc.")
|
||||
}
|
||||
|
|
@ -95,12 +112,27 @@ impl RunEventListParams {
|
|||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Why the parameters do not address a Petri run's stream, if they do
|
||||
/// not: the legacy cursors have no meaning there.
|
||||
fn stream_cursor_error(&self) -> Option<&'static str> {
|
||||
if self.since_seq.is_some() || self.before_seq.is_some() || self.order.is_some() {
|
||||
return Some(
|
||||
"this run executes on Petri; its events are a run stream addressed by `after` \
|
||||
(the last stream_seq seen), not by since_seq, before_seq or order.",
|
||||
);
|
||||
}
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct AttachParams {
|
||||
#[serde(default)]
|
||||
since_seq: Option<u32>,
|
||||
/// The run stream cursor of a Petri run: the last `stream_seq` seen.
|
||||
#[serde(default)]
|
||||
after: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
|
|
@ -256,9 +288,25 @@ async fn list_run_events(
|
|||
}
|
||||
|
||||
let limit = params.limit();
|
||||
match run_is_petri(&state, &id).await {
|
||||
Ok(true) => {
|
||||
if let Some(detail) = params.stream_cursor_error() {
|
||||
return ApiError::bad_request(detail).into_response();
|
||||
}
|
||||
return list_run_stream(&state, id, params.after.unwrap_or(0), limit).await;
|
||||
}
|
||||
Ok(false) => {}
|
||||
Err(response) => return response,
|
||||
}
|
||||
if params.after.is_some() {
|
||||
return ApiError::bad_request(
|
||||
"after is the run stream cursor of a Petri run; this run's events use since_seq.",
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
match state.stores.runs.open_run_reader(&id).await {
|
||||
Ok(run_store) => {
|
||||
let events = match params.order {
|
||||
let events = match params.order() {
|
||||
EventSequenceOrder::Asc => {
|
||||
run_store
|
||||
.list_events_from_with_limit(params.since_seq(), limit)
|
||||
|
|
@ -291,6 +339,171 @@ async fn list_run_events(
|
|||
}
|
||||
}
|
||||
|
||||
/// Whether the run executes on Petri, from its stored spec; the canonical
|
||||
/// 404 when there is no such run.
|
||||
async fn run_is_petri(state: &AppState, id: &RunId) -> Result<bool, Response> {
|
||||
let projection = state
|
||||
.load_run_projection(id)
|
||||
.await
|
||||
.map_err(IntoResponse::into_response)?;
|
||||
Ok(projection.spec.engine.is_petri())
|
||||
}
|
||||
|
||||
/// One page of a Petri run's stream past `after`.
|
||||
async fn list_run_stream(state: &AppState, id: RunId, after: u64, limit: usize) -> Response {
|
||||
match state
|
||||
.petri_projector
|
||||
.stream_after(id, after, limit.saturating_add(1))
|
||||
.await
|
||||
{
|
||||
Ok(mut items) => {
|
||||
let has_more = items.len() > limit;
|
||||
items.truncate(limit);
|
||||
Json(PaginatedRunStreamList {
|
||||
data: items,
|
||||
meta: PaginationMeta {
|
||||
has_more,
|
||||
total: None,
|
||||
},
|
||||
event_contract_version: EVENT_CONTRACT_VERSION,
|
||||
})
|
||||
.into_response()
|
||||
}
|
||||
Err(err) => {
|
||||
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn sse_event_from_stream_item(item: &RunStreamItem) -> Option<Event> {
|
||||
let data = serde_json::to_string(item).ok()?;
|
||||
let data = redact_jsonl_line(&data);
|
||||
Some(Event::default().data(data))
|
||||
}
|
||||
|
||||
/// How many stream items one read takes while attached.
|
||||
const STREAM_ATTACH_BATCH_LIMIT: usize = 256;
|
||||
|
||||
/// How long an attached reader waits for a commit signal before it re-reads
|
||||
/// its cursor anyway: a signal is a wake-up, never the source of facts.
|
||||
const STREAM_ATTACH_POLL: Duration = Duration::from_secs(1);
|
||||
|
||||
/// How long an attached reader keeps following a run whose projection is
|
||||
/// already terminal, waiting for the platform record of the terminal
|
||||
/// lifecycle transition that ends the stream; after that it ends anyway.
|
||||
const STREAM_ATTACH_TERMINAL_GRACE: Duration = Duration::from_secs(15);
|
||||
|
||||
/// Whether the item ends an attached stream: the platform record of the
|
||||
/// run's terminal lifecycle transition, which Fabro writes after the engine
|
||||
/// recorded the run's finish. The analog of the legacy stream's
|
||||
/// `run.completed` and `run.failed`.
|
||||
fn stream_item_is_terminal(item: &RunStreamItem) -> bool {
|
||||
if item.kind != fabro_types::RunStreamItemKind::Platform {
|
||||
return false;
|
||||
}
|
||||
let record = &item.item["record"];
|
||||
record["kind"].as_str() == Some("run.lifecycle")
|
||||
&& matches!(
|
||||
record["transition"].as_str(),
|
||||
Some("succeeded" | "failed" | "dead")
|
||||
)
|
||||
}
|
||||
|
||||
/// The live stream of a Petri run from `after` (the last `stream_seq` the
|
||||
/// client saw; `None` starts at the next unseen item), as server-sent
|
||||
/// events. Every committed item past the cursor is sent once, in order,
|
||||
/// and the stream ends once the run is no longer active and every
|
||||
/// committed item is out.
|
||||
async fn attach_run_stream(state: Arc<AppState>, id: RunId, after: Option<u64>) -> Response {
|
||||
let cursor = match after {
|
||||
Some(after) => after,
|
||||
None => match state.petri_projector.stream_head(id).await {
|
||||
Ok(head) => head.unwrap_or(0),
|
||||
Err(err) => {
|
||||
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
||||
.into_response();
|
||||
}
|
||||
},
|
||||
};
|
||||
let (sender, receiver) = mpsc::unbounded_channel();
|
||||
let shutdown = state.shutdown_token();
|
||||
tokio::spawn(async move {
|
||||
// Subscribed before the first read, so a pass that commits between
|
||||
// the read and the wait is not missed.
|
||||
let mut committed = state.petri_projector.subscribe();
|
||||
let mut cursor = cursor;
|
||||
// Set once the projection is terminal: the stream then ends at the
|
||||
// terminal lifecycle record, or when the grace runs out.
|
||||
let mut terminal_deadline: Option<Instant> = None;
|
||||
loop {
|
||||
// Drain everything committed past the cursor.
|
||||
let mut drained = false;
|
||||
while !drained {
|
||||
let Ok(items) = state
|
||||
.petri_projector
|
||||
.stream_after(id, cursor, STREAM_ATTACH_BATCH_LIMIT)
|
||||
.await
|
||||
else {
|
||||
return;
|
||||
};
|
||||
drained = items.len() < STREAM_ATTACH_BATCH_LIMIT;
|
||||
for item in items {
|
||||
cursor = item.stream_seq;
|
||||
let terminal = stream_item_is_terminal(&item);
|
||||
if let Some(sse_event) = sse_event_from_stream_item(&item) {
|
||||
if sender
|
||||
.send(Ok::<Event, std::convert::Infallible>(sse_event))
|
||||
.is_err()
|
||||
{
|
||||
return;
|
||||
}
|
||||
}
|
||||
if terminal {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The run's status is read after the drain, so an item
|
||||
// committed with the finish is already out. Once terminal, the
|
||||
// stream keeps following for the terminal lifecycle record,
|
||||
// which Fabro writes after the engine's finish, for a bounded
|
||||
// time.
|
||||
if terminal_deadline.is_none() {
|
||||
let active = match state.stores.runs.load_run_projection(&id).await {
|
||||
Ok(Some(projection)) => run_projection_is_active(&projection),
|
||||
Ok(None) | Err(_) => false,
|
||||
};
|
||||
if !active {
|
||||
terminal_deadline = Some(Instant::now() + STREAM_ATTACH_TERMINAL_GRACE);
|
||||
}
|
||||
}
|
||||
if terminal_deadline.is_some_and(|deadline| Instant::now() >= deadline) {
|
||||
return;
|
||||
}
|
||||
|
||||
// Wait for the projector to commit more of this run, or poll.
|
||||
loop {
|
||||
tokio::select! {
|
||||
biased;
|
||||
() = shutdown.cancelled() => return,
|
||||
signal = committed.recv() => match signal {
|
||||
Ok(run_id) if run_id == id => break,
|
||||
Ok(_) => {}
|
||||
Err(RecvError::Lagged(_)) => break,
|
||||
Err(RecvError::Closed) => return,
|
||||
},
|
||||
() = time::sleep(STREAM_ATTACH_POLL) => break,
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
Sse::new(UnboundedReceiverStream::new(receiver))
|
||||
.keep_alive(KeepAlive::default())
|
||||
.into_response()
|
||||
}
|
||||
|
||||
async fn list_run_stage_events(
|
||||
RequireRunStageScoped(id, stage_id): RequireRunStageScoped,
|
||||
State(state): State<Arc<AppState>>,
|
||||
|
|
@ -446,6 +659,11 @@ async fn attach_run_events(
|
|||
Ok(id) => id,
|
||||
Err(response) => return response,
|
||||
};
|
||||
match run_is_petri(&state, &id).await {
|
||||
Ok(true) => return attach_run_stream(state, id, params.after).await,
|
||||
Ok(false) => {}
|
||||
Err(response) => return response,
|
||||
}
|
||||
let Ok(run_store) = state.stores.runs.open_run_reader(&id).await else {
|
||||
return ApiError::not_found("Run not found.").into_response();
|
||||
};
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ mod archive;
|
|||
mod dry_run;
|
||||
mod lifecycle;
|
||||
mod petri;
|
||||
mod petri_stream;
|
||||
mod run_completion;
|
||||
mod sse;
|
||||
mod usage;
|
||||
|
|
|
|||
|
|
@ -105,14 +105,14 @@ const PARALLEL_DOT: &str = r#"digraph Parallel {
|
|||
merge -> exit
|
||||
}"#;
|
||||
|
||||
const PLAIN_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
|
||||
pub(super) const PLAIN_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
|
||||
const PETRI_SETTINGS: &str =
|
||||
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n";
|
||||
|
||||
/// The host plugin as Petri's lookup finds it: the override variable, else
|
||||
/// the executable on `PATH`. `None`, after saying so, when the test should
|
||||
/// skip; a panic when the environment forbids a skip.
|
||||
fn host_plugin() -> Option<PathBuf> {
|
||||
pub(super) fn host_plugin() -> Option<PathBuf> {
|
||||
let found = env::var_os(HOST_PLUGIN_OVERRIDE)
|
||||
.map(PathBuf::from)
|
||||
.or_else(|| {
|
||||
|
|
@ -132,7 +132,7 @@ fn host_plugin() -> Option<PathBuf> {
|
|||
|
||||
/// Register a version whose entrypoint is `workflow.fabro`, with the given
|
||||
/// files beside it.
|
||||
async fn register_version(app: &axum::Router, files: &[(&str, &str)]) -> String {
|
||||
pub(super) async fn register_version(app: &axum::Router, files: &[(&str, &str)]) -> String {
|
||||
let entrypoint = WorkflowPath::new("workflow.fabro").expect("entrypoint path is valid");
|
||||
let files = files
|
||||
.iter()
|
||||
|
|
@ -150,7 +150,7 @@ async fn register_version(app: &axum::Router, files: &[(&str, &str)]) -> String
|
|||
.to_string()
|
||||
}
|
||||
|
||||
fn intent(version_id: &str, workspace: &std::path::Path) -> serde_json::Value {
|
||||
pub(super) fn intent(version_id: &str, workspace: &std::path::Path) -> serde_json::Value {
|
||||
serde_json::json!({
|
||||
"workflow_version_id": version_id,
|
||||
"target": {"kind": "folder", "path": workspace},
|
||||
|
|
@ -187,7 +187,11 @@ async fn petri_outcome(state: &AppState, run_id: &str) -> engine::RunOutcome {
|
|||
}
|
||||
|
||||
/// The run's projected state once its projector settled.
|
||||
async fn settled_state(state: &AppState, app: &axum::Router, run_id: &str) -> serde_json::Value {
|
||||
pub(super) async fn settled_state(
|
||||
state: &AppState,
|
||||
app: &axum::Router,
|
||||
run_id: &str,
|
||||
) -> serde_json::Value {
|
||||
let id: RunId = run_id.parse().expect("the run id parses");
|
||||
state.test_petri_projector().settle(id).await;
|
||||
let req = Request::builder()
|
||||
|
|
@ -335,6 +339,7 @@ async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() {
|
|||
.any(|request| request["model"] == OPENAI_MODEL),
|
||||
"the prompt stage should have called the twin, got {logs}"
|
||||
);
|
||||
super::petri_stream::capture_settled(&state, &app, &run_id, "hello").await;
|
||||
}
|
||||
|
||||
/// A command-only bundle runs on Petri when the server's setting names the
|
||||
|
|
@ -379,6 +384,7 @@ async fn a_command_bundle_runs_on_petri_under_the_server_setting() {
|
|||
);
|
||||
let stream = petri_stream_len(&state, &run_id).await;
|
||||
assert!(stream > 0, "the run's stream holds its events");
|
||||
super::petri_stream::capture_settled(&state, &app, &run_id, "command").await;
|
||||
}
|
||||
|
||||
/// A parallel bundle with two command branches runs on Petri through the
|
||||
|
|
@ -685,4 +691,5 @@ async fn a_human_gate_is_answered_through_the_questions_api() {
|
|||
record.principal.is_some(),
|
||||
"the answering principal: {record:?}"
|
||||
);
|
||||
super::petri_stream::capture_settled(&state, &app, &run_id, "gate").await;
|
||||
}
|
||||
|
|
|
|||
421
lib/apps/fabro-server/tests/it/scenario/petri_stream.rs
Normal file
421
lib/apps/fabro-server/tests/it/scenario/petri_stream.rs
Normal file
|
|
@ -0,0 +1,421 @@
|
|||
//! The run stream of a Petri run through the server: `GET /runs/{id}/events`
|
||||
//! pages it by `after`, `GET /runs/{id}/attach` follows it live, and a
|
||||
//! client that disconnects mid-run and reconnects from its last
|
||||
//! `stream_seq` receives every item once, in order, with no gap and no
|
||||
//! duplicate, including a platform record Fabro recorded between two
|
||||
//! concurrent child executions' events.
|
||||
//!
|
||||
//! The runs execute in the server process under the handler-registry test
|
||||
//! override and take their host scope through the sandbox-driver host
|
||||
//! plugin, so the tests skip, and say why, when the executable is not
|
||||
//! found (see `petri.rs`).
|
||||
//!
|
||||
//! With `FABRO_CAPTURE_PETRI_FIXTURES` set, a scenario also writes its
|
||||
//! settled projection and full stream as JSON under the web app's test
|
||||
//! fixtures (`apps/fabro-web/app/test-fixtures/petri/`), which the web
|
||||
//! app's rendering tests read.
|
||||
|
||||
#![expect(
|
||||
clippy::disallowed_methods,
|
||||
reason = "the tests locate the plugin executable and the capture switch through the process environment"
|
||||
)]
|
||||
#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")]
|
||||
|
||||
use std::collections::BTreeSet;
|
||||
use std::env;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use axum::body::Body;
|
||||
use axum::http::{Request, StatusCode};
|
||||
use fabro_server::server::AppState;
|
||||
use fabro_store::platform_records::{PlatformRecord, PlatformRecordStore, RunNoticeRecord};
|
||||
use fabro_types::run_event::RunNoticeLevel;
|
||||
use fabro_types::{RunId, RunStreamItem, RunStreamItemKind};
|
||||
use http_body_util::BodyExt;
|
||||
use tokio::time::timeout;
|
||||
use tower::ServiceExt;
|
||||
|
||||
use super::petri::{PLAIN_SETTINGS, host_plugin, intent, register_version, settled_state};
|
||||
use crate::helpers::{
|
||||
api, create_and_start_run_from_intent, repo_root, response_json, run_json, settings_from_toml,
|
||||
test_app_state_with_options, test_app_with_scheduler, wait_for_run_status,
|
||||
};
|
||||
|
||||
const CAPTURE_ENV: &str = "FABRO_CAPTURE_PETRI_FIXTURES";
|
||||
const FRAME_TIMEOUT: Duration = Duration::from_secs(20);
|
||||
|
||||
/// Two command branches that announce they started and wait for a release
|
||||
/// marker, so a test can act between their events.
|
||||
fn gated_parallel_dot(markers: &std::path::Path) -> String {
|
||||
let dir = markers.display();
|
||||
format!(
|
||||
r#"digraph Parallel {{
|
||||
graph [goal="Run two branches"]
|
||||
start [shape=Mdiamond]
|
||||
exit [shape=Msquare]
|
||||
fork [shape=component]
|
||||
a [shape=parallelogram, script="touch {dir}/a.started; while [ ! -f {dir}/go ]; do sleep 0.05; done; echo a"]
|
||||
b [shape=parallelogram, script="touch {dir}/b.started; while [ ! -f {dir}/go ]; do sleep 0.05; done; echo b"]
|
||||
merge [shape=tripleoctagon]
|
||||
start -> fork
|
||||
fork -> a
|
||||
fork -> b
|
||||
a -> merge
|
||||
b -> merge
|
||||
merge -> exit
|
||||
}}"#
|
||||
)
|
||||
}
|
||||
|
||||
/// One page of the run's stream past `after`.
|
||||
async fn stream_page(
|
||||
app: &axum::Router,
|
||||
run_id: &str,
|
||||
after: u64,
|
||||
limit: usize,
|
||||
) -> serde_json::Value {
|
||||
let req = Request::builder()
|
||||
.method("GET")
|
||||
.uri(api(&format!(
|
||||
"/runs/{run_id}/events?after={after}&limit={limit}"
|
||||
)))
|
||||
.body(Body::empty())
|
||||
.expect("events request should build");
|
||||
let response = app
|
||||
.clone()
|
||||
.oneshot(req)
|
||||
.await
|
||||
.expect("events request routes");
|
||||
response_json(
|
||||
response,
|
||||
StatusCode::OK,
|
||||
format!("GET /api/v1/runs/{run_id}/events?after={after}&limit={limit}"),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Every item of the run's stream, paged through the listing endpoint.
|
||||
pub(super) async fn list_stream(
|
||||
app: &axum::Router,
|
||||
run_id: &str,
|
||||
page_limit: usize,
|
||||
) -> Vec<RunStreamItem> {
|
||||
let mut after = 0;
|
||||
let mut items = Vec::new();
|
||||
loop {
|
||||
let page = stream_page(app, run_id, after, page_limit).await;
|
||||
let data: Vec<RunStreamItem> =
|
||||
serde_json::from_value(page["data"].clone()).expect("stream items decode");
|
||||
let has_more = page["meta"]["has_more"]
|
||||
.as_bool()
|
||||
.expect("has_more is a bool");
|
||||
assert_eq!(
|
||||
page["event_contract_version"].as_u64(),
|
||||
Some(3),
|
||||
"the server reports Petri's contract version: {page}"
|
||||
);
|
||||
let Some(last) = data.last() else {
|
||||
assert!(!has_more, "an empty page is the last");
|
||||
break;
|
||||
};
|
||||
after = last.stream_seq;
|
||||
items.extend(data);
|
||||
if !has_more {
|
||||
break;
|
||||
}
|
||||
}
|
||||
items
|
||||
}
|
||||
|
||||
/// An attached reader of the run's stream that stops reading when `until`
|
||||
/// says so, as a client that lost its connection would: the frames it saw
|
||||
/// so far come back.
|
||||
struct Attached {
|
||||
body: Body,
|
||||
pending: String,
|
||||
}
|
||||
|
||||
impl Attached {
|
||||
async fn open(app: &axum::Router, run_id: &str, after: u64) -> Self {
|
||||
let req = Request::builder()
|
||||
.method("GET")
|
||||
.uri(api(&format!("/runs/{run_id}/attach?after={after}")))
|
||||
.body(Body::empty())
|
||||
.expect("attach request should build");
|
||||
let response = app
|
||||
.clone()
|
||||
.oneshot(req)
|
||||
.await
|
||||
.expect("attach request routes");
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
assert!(
|
||||
response
|
||||
.headers()
|
||||
.get("content-type")
|
||||
.and_then(|value| value.to_str().ok())
|
||||
.is_some_and(|value| value.contains("text/event-stream")),
|
||||
"an SSE response"
|
||||
);
|
||||
Self {
|
||||
body: response.into_body(),
|
||||
pending: String::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// The next item on the stream, or `None` once the server ended it.
|
||||
async fn next(&mut self) -> Option<RunStreamItem> {
|
||||
loop {
|
||||
if let Some(end) = self.pending.find("\n\n") {
|
||||
let frame = self.pending[..end].to_string();
|
||||
self.pending.drain(..end + 2);
|
||||
let data = frame
|
||||
.lines()
|
||||
.filter_map(|line| line.strip_prefix("data:"))
|
||||
.map(str::trim)
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n");
|
||||
if data.is_empty() {
|
||||
continue;
|
||||
}
|
||||
return Some(serde_json::from_str(&data).expect("a stream item frame decodes"));
|
||||
}
|
||||
let frame = timeout(FRAME_TIMEOUT, self.body.frame())
|
||||
.await
|
||||
.expect("the attached stream keeps sending or ends");
|
||||
match frame {
|
||||
Some(Ok(frame)) => {
|
||||
if let Some(data) = frame.data_ref() {
|
||||
self.pending.push_str(&String::from_utf8_lossy(data));
|
||||
}
|
||||
}
|
||||
Some(Err(err)) => panic!("the attached stream failed: {err}"),
|
||||
None => return None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn petri_name(item: &RunStreamItem) -> Option<&str> {
|
||||
(item.kind == RunStreamItemKind::Petri)
|
||||
.then(|| item.name())
|
||||
.flatten()
|
||||
}
|
||||
|
||||
/// The subject's node name, for a node that is a stage of its own: the
|
||||
/// `parallel.branch` delegate the fork's execution holds for each branch
|
||||
/// shares the branch's name and is not one.
|
||||
fn subject_node(item: &RunStreamItem) -> Option<&str> {
|
||||
let node = &item.item["subject"]["node"];
|
||||
if node["meta"]["kind"].as_str() == Some("parallel.branch") {
|
||||
return None;
|
||||
}
|
||||
node["name"].as_str()
|
||||
}
|
||||
|
||||
fn wait_for_marker(path: &std::path::Path) {
|
||||
let deadline = std::time::Instant::now() + Duration::from_secs(20);
|
||||
while !path.exists() {
|
||||
assert!(
|
||||
std::time::Instant::now() < deadline,
|
||||
"{} never appeared",
|
||||
path.display()
|
||||
);
|
||||
std::thread::sleep(Duration::from_millis(20));
|
||||
}
|
||||
}
|
||||
|
||||
/// A client attached to a two-branch parallel run disconnects once both
|
||||
/// branches have started, Fabro records a platform notice while they run,
|
||||
/// the client reconnects from its last `stream_seq`, and the union of what
|
||||
/// it saw is the whole stream: every item once, in `stream_seq` order, no
|
||||
/// gap, no duplicate, with the notice between the branches' events.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn a_reconnecting_client_receives_every_stream_item_once_in_order() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let workspace = tempfile::tempdir().expect("workspace tempdir");
|
||||
let markers = tempfile::tempdir().expect("marker tempdir");
|
||||
let settings = settings_from_toml(
|
||||
"_version = 1\n\n[run.environment]\nid = \"local\"\n\n[server.execution]\nengine = \
|
||||
\"petri\"\n",
|
||||
);
|
||||
let state = test_app_state_with_options(settings, 5);
|
||||
let app = test_app_with_scheduler(Arc::clone(&state));
|
||||
|
||||
let dot = gated_parallel_dot(markers.path());
|
||||
let version_id = register_version(&app, &[
|
||||
("workflow.fabro", &dot),
|
||||
("workflow.toml", PLAIN_SETTINGS),
|
||||
])
|
||||
.await;
|
||||
let run_id =
|
||||
create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await;
|
||||
let id: RunId = run_id.parse().expect("the run id parses");
|
||||
|
||||
// First connection: from the start of the run until both branches
|
||||
// have a `visit.started` on the stream, then drop it.
|
||||
let mut first = Attached::open(&app, &run_id, 0).await;
|
||||
let mut seen_first: Vec<RunStreamItem> = Vec::new();
|
||||
let mut started: BTreeSet<String> = BTreeSet::new();
|
||||
while started.len() < 2 {
|
||||
let item = first
|
||||
.next()
|
||||
.await
|
||||
.expect("the stream runs until both branches started");
|
||||
if petri_name(&item) == Some("visit.started") {
|
||||
if let Some(node @ ("a" | "b")) = subject_node(&item) {
|
||||
started.insert(node.to_string());
|
||||
}
|
||||
}
|
||||
seen_first.push(item);
|
||||
}
|
||||
let last_seen = seen_first
|
||||
.last()
|
||||
.map(|item| item.stream_seq)
|
||||
.expect("something was seen");
|
||||
drop(first);
|
||||
|
||||
// Both branch scripts are running: record a platform fact between
|
||||
// their events, as a checkpoint or a notice would be, then let them go.
|
||||
wait_for_marker(&markers.path().join("a.started"));
|
||||
wait_for_marker(&markers.path().join("b.started"));
|
||||
let platform = PlatformRecordStore::new(state.test_petri_view_pool());
|
||||
let notice = platform
|
||||
.append(
|
||||
&id,
|
||||
&PlatformRecord::RunNotice(RunNoticeRecord {
|
||||
level: RunNoticeLevel::Info,
|
||||
code: "test.between_branches".to_string(),
|
||||
message: "recorded while both branches ran".to_string(),
|
||||
}),
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.expect("the notice appends");
|
||||
state.test_petri_projector().signal(id);
|
||||
state.test_petri_projector().settle(id).await;
|
||||
std::fs::write(markers.path().join("go"), b"").expect("the release marker writes");
|
||||
|
||||
// Second connection: resume from the last stream_seq seen and read to
|
||||
// the end of the stream, which the server closes once the run is done.
|
||||
let mut second = Attached::open(&app, &run_id, last_seen).await;
|
||||
let mut seen_second: Vec<RunStreamItem> = Vec::new();
|
||||
while let Some(item) = second.next().await {
|
||||
seen_second.push(item);
|
||||
}
|
||||
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
|
||||
let run = run_json(&app, &run_id).await;
|
||||
assert_eq!(status, "succeeded", "run: {run}");
|
||||
let projection = settled_state(&state, &app, &run_id).await;
|
||||
assert_eq!(projection["status"]["kind"], "succeeded", "{projection}");
|
||||
|
||||
// The union is the whole stream, once each, in order, with no gap.
|
||||
let mut union = seen_first;
|
||||
union.extend(seen_second);
|
||||
let seqs: Vec<u64> = union.iter().map(|item| item.stream_seq).collect();
|
||||
let expected: Vec<u64> = (1..=seqs.len() as u64).collect();
|
||||
assert_eq!(
|
||||
seqs, expected,
|
||||
"stream_seq is dense and strictly increasing"
|
||||
);
|
||||
let ids: BTreeSet<&str> = union.iter().map(|item| item.id.as_str()).collect();
|
||||
assert_eq!(ids.len(), union.len(), "every item identity appears once");
|
||||
for item in &union {
|
||||
assert_eq!(item.run_id, id);
|
||||
assert!(item.recorded_at > 0, "{item:?}");
|
||||
}
|
||||
|
||||
// The same stream, paged through the listing endpoint with small
|
||||
// pages, is item for item what the attached client saw.
|
||||
let listed = list_stream(&app, &run_id, 7).await;
|
||||
assert_eq!(listed, union, "the listing pages the same stream");
|
||||
|
||||
// The notice sits between the branches' events.
|
||||
let notice_seq = union
|
||||
.iter()
|
||||
.find(|item| item.kind == RunStreamItemKind::Platform && item.id == notice.seq.to_string())
|
||||
.map(|item| item.stream_seq)
|
||||
.expect("the notice is on the stream");
|
||||
let branch_seqs = |name: &str| -> Vec<u64> {
|
||||
union
|
||||
.iter()
|
||||
.filter(|item| {
|
||||
petri_name(item) == Some(name) && matches!(subject_node(item), Some("a" | "b"))
|
||||
})
|
||||
.map(|item| item.stream_seq)
|
||||
.collect()
|
||||
};
|
||||
let starts = branch_seqs("visit.started");
|
||||
let ends = branch_seqs("visit.completed");
|
||||
assert_eq!(starts.len(), 2, "{starts:?}");
|
||||
assert_eq!(ends.len(), 2, "{ends:?}");
|
||||
assert!(
|
||||
starts.iter().all(|seq| *seq < notice_seq) && ends.iter().all(|seq| *seq > notice_seq),
|
||||
"the notice ({notice_seq}) is between the branch starts {starts:?} and ends {ends:?}"
|
||||
);
|
||||
let finished = union
|
||||
.iter()
|
||||
.filter(|item| petri_name(item) == Some("run.finished"))
|
||||
.count();
|
||||
assert_eq!(finished, 1, "the stream ends with the run's finish");
|
||||
|
||||
// The legacy cursors are refused for a Petri run; the stream cursor is
|
||||
// refused for nothing else.
|
||||
let req = Request::builder()
|
||||
.method("GET")
|
||||
.uri(api(&format!("/runs/{run_id}/events?since_seq=1")))
|
||||
.body(Body::empty())
|
||||
.expect("events request should build");
|
||||
let response = app
|
||||
.clone()
|
||||
.oneshot(req)
|
||||
.await
|
||||
.expect("events request routes");
|
||||
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
|
||||
|
||||
capture_fixture(&app, &run_id, "parallel", &projection).await;
|
||||
}
|
||||
|
||||
/// Write the run's settled projection and its whole stream under the web
|
||||
/// app's test fixtures, when the capture switch is set.
|
||||
pub(super) async fn capture_fixture(
|
||||
app: &axum::Router,
|
||||
run_id: &str,
|
||||
name: &str,
|
||||
projection: &serde_json::Value,
|
||||
) {
|
||||
if env::var_os(CAPTURE_ENV).is_none() {
|
||||
return;
|
||||
}
|
||||
let stream = list_stream(app, run_id, 1000).await;
|
||||
let fixture = serde_json::json!({
|
||||
"run_id": run_id,
|
||||
"projection": projection,
|
||||
"stream": stream,
|
||||
});
|
||||
let dir = repo_root().join("apps/fabro-web/app/test-fixtures/petri");
|
||||
std::fs::create_dir_all(&dir).expect("the fixture directory creates");
|
||||
let path = dir.join(format!("{name}.json"));
|
||||
std::fs::write(
|
||||
&path,
|
||||
serde_json::to_string_pretty(&fixture).expect("the fixture serializes"),
|
||||
)
|
||||
.expect("the fixture writes");
|
||||
eprintln!("captured {}", path.display());
|
||||
}
|
||||
|
||||
/// Helpers the other Petri scenarios use to capture their fixtures.
|
||||
pub(super) async fn capture_settled(
|
||||
state: &AppState,
|
||||
app: &axum::Router,
|
||||
run_id: &str,
|
||||
name: &str,
|
||||
) {
|
||||
if env::var_os(CAPTURE_ENV).is_none() {
|
||||
return;
|
||||
}
|
||||
let projection = settled_state(state, app, run_id).await;
|
||||
capture_fixture(app, run_id, name, &projection).await;
|
||||
}
|
||||
|
|
@ -3,6 +3,11 @@
|
|||
//! depending on the Petri packages themselves. Only this crate names them in
|
||||
//! its `Cargo.toml`.
|
||||
|
||||
use petri_execution::events;
|
||||
pub use petri_store::{
|
||||
Access, Digest, ExecutionId, LogId, OwnerId, Record, RunKey, RunLogs, RunStore, StoreError,
|
||||
};
|
||||
|
||||
/// The version of Petri's public event contract this build serves on the
|
||||
/// run stream: every `petri` item of `GET /runs/{id}/events` follows it.
|
||||
pub const EVENT_CONTRACT_VERSION: u32 = events::EVENT_CONTRACT_VERSION;
|
||||
|
|
|
|||
|
|
@ -46,13 +46,13 @@ use std::time::Duration;
|
|||
use fabro_db::DbPool;
|
||||
use fabro_store::platform_records::{PlatformRecordStore, StoredPlatformRecord, now_ms};
|
||||
use fabro_store::{RunProjection, RunSummaryStore};
|
||||
use fabro_types::RunId;
|
||||
use fabro_types::{RunId, RunStreamItem, RunStreamItemKind};
|
||||
use fabro_util::error::collect_chain;
|
||||
use petri_execution::events::{self, EventId, EventSource, RunEvent};
|
||||
use petri_execution::{Access, RunKey, RunStore as _, inspect};
|
||||
use petri_store::StoreError;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::sync::Mutex as AsyncMutex;
|
||||
use tokio::sync::{Mutex as AsyncMutex, broadcast};
|
||||
use tokio::time;
|
||||
use tracing::{debug, info, warn};
|
||||
|
||||
|
|
@ -147,17 +147,21 @@ struct Slot {
|
|||
/// the server both are the one database; a test may hand it the run
|
||||
/// summary store's own pool for the views.
|
||||
pub struct Projector {
|
||||
records: DbPool,
|
||||
pool: DbPool,
|
||||
store: SqliteRunStore,
|
||||
platform: PlatformRecordStore,
|
||||
slots: Mutex<HashMap<RunId, Slot>>,
|
||||
records: DbPool,
|
||||
pool: DbPool,
|
||||
store: SqliteRunStore,
|
||||
platform: PlatformRecordStore,
|
||||
slots: Mutex<HashMap<RunId, Slot>>,
|
||||
/// One pass at a time per run: a signalled pass and the startup pass
|
||||
/// over the same run never interleave their reads and writes.
|
||||
passes: Mutex<HashMap<RunId, Arc<AsyncMutex<()>>>>,
|
||||
passes: Mutex<HashMap<RunId, Arc<AsyncMutex<()>>>>,
|
||||
/// Test-only: stop the next pass after its reads, before its view
|
||||
/// transaction, as a crash there would.
|
||||
fault: AtomicBool,
|
||||
fault: AtomicBool,
|
||||
/// Sent after each committed pass that wrote stream rows: the run whose
|
||||
/// stream grew. A wake-up for the stream's readers, never a source of
|
||||
/// facts; a reader that lags re-reads from its cursor.
|
||||
committed: broadcast::Sender<RunId>,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for Projector {
|
||||
|
|
@ -180,9 +184,42 @@ impl Projector {
|
|||
slots: Mutex::default(),
|
||||
passes: Mutex::default(),
|
||||
fault: AtomicBool::new(false),
|
||||
committed: broadcast::channel(COMMIT_SIGNAL_CAPACITY).0,
|
||||
})
|
||||
}
|
||||
|
||||
/// A receiver that learns which run's stream grew after each committed
|
||||
/// pass. A receiver that falls behind gets `Lagged` and treats it as a
|
||||
/// wake-up for every run it follows.
|
||||
#[must_use]
|
||||
pub fn subscribe(&self) -> broadcast::Receiver<RunId> {
|
||||
self.committed.subscribe()
|
||||
}
|
||||
|
||||
/// The run's stream past the cursor: up to `limit` items with
|
||||
/// `stream_seq > after`, in `stream_seq` order, each in Fabro's
|
||||
/// envelope. `after = 0` reads from the first item.
|
||||
pub async fn stream_after(
|
||||
&self,
|
||||
run_id: RunId,
|
||||
after: u64,
|
||||
limit: usize,
|
||||
) -> Result<Vec<RunStreamItem>, ProjectError> {
|
||||
stream_after(&self.pool, run_id, after, limit).await
|
||||
}
|
||||
|
||||
/// The last delivery sequence the run's view holds, or `None` when no
|
||||
/// pass has committed a view for it.
|
||||
pub async fn stream_head(&self, run_id: RunId) -> Result<Option<u64>, ProjectError> {
|
||||
let head: Option<i64> =
|
||||
sqlx::query_scalar("SELECT stream_seq FROM petri_projection WHERE run_id = ?")
|
||||
.bind(run_id.to_string())
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(ProjectError::Database)?;
|
||||
Ok(head.map(|head| u64::try_from(head).unwrap_or(0)))
|
||||
}
|
||||
|
||||
/// Schedule a pass for the run. A pass already running for it runs once
|
||||
/// more when it ends; any number of signals in between coalesce.
|
||||
pub fn signal(self: &Arc<Self>, run_id: RunId) {
|
||||
|
|
@ -471,6 +508,10 @@ impl Projector {
|
|||
stream_seq,
|
||||
"Petri projection pass committed"
|
||||
);
|
||||
if !rows.is_empty() {
|
||||
// No receiver is not an error: nobody follows the stream.
|
||||
let _ = self.committed.send(run_id);
|
||||
}
|
||||
Ok(PassReport {
|
||||
run_id,
|
||||
skipped: false,
|
||||
|
|
@ -658,6 +699,52 @@ struct StreamRow {
|
|||
event_json: String,
|
||||
}
|
||||
|
||||
/// How many commit signals a slow reader may fall behind before it is told
|
||||
/// it lagged and re-reads from its cursor.
|
||||
const COMMIT_SIGNAL_CAPACITY: usize = 1024;
|
||||
|
||||
/// The run's stream past the cursor, read from the view tables: up to
|
||||
/// `limit` rows with `stream_seq > after`, in order, in Fabro's envelope.
|
||||
pub async fn stream_after(
|
||||
views: &DbPool,
|
||||
run_id: RunId,
|
||||
after: u64,
|
||||
limit: usize,
|
||||
) -> Result<Vec<RunStreamItem>, ProjectError> {
|
||||
let rows: Vec<(i64, String, String, String)> = sqlx::query_as(
|
||||
"SELECT stream_seq, item_kind, item_id, event_json FROM petri_stream WHERE run_id = ? AND \
|
||||
stream_seq > ? ORDER BY stream_seq LIMIT ?",
|
||||
)
|
||||
.bind(run_id.to_string())
|
||||
.bind(column(after))
|
||||
.bind(i64::try_from(limit).unwrap_or(i64::MAX))
|
||||
.fetch_all(views)
|
||||
.await
|
||||
.map_err(ProjectError::Database)?;
|
||||
rows.into_iter()
|
||||
.map(|(stream_seq, item_kind, item_id, event_json)| {
|
||||
let item: serde_json::Value =
|
||||
serde_json::from_str(&event_json).map_err(ProjectError::Encode)?;
|
||||
let kind = match item_kind.as_str() {
|
||||
"platform" => RunStreamItemKind::Platform,
|
||||
_ => RunStreamItemKind::Petri,
|
||||
};
|
||||
let recorded_at = item
|
||||
.get("recorded_at")
|
||||
.and_then(serde_json::Value::as_u64)
|
||||
.unwrap_or(0);
|
||||
Ok(RunStreamItem {
|
||||
run_id,
|
||||
stream_seq: u64::try_from(stream_seq).unwrap_or(0),
|
||||
kind,
|
||||
id: item_id,
|
||||
recorded_at,
|
||||
item,
|
||||
})
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// A Petri event id as the stream names it: `<log>/<seq>/<index>`.
|
||||
#[must_use]
|
||||
pub fn event_id_text(id: &EventId) -> String {
|
||||
|
|
|
|||
|
|
@ -658,6 +658,11 @@ fn main() {
|
|||
&[],
|
||||
),
|
||||
("EventEnvelope", "fabro_types::EventEnvelope", &[]),
|
||||
("RunStreamItem", "fabro_types::RunStreamItem", &[]),
|
||||
("RunStreamItemKind", "fabro_types::RunStreamItemKind", &[]),
|
||||
("RunEngine", "fabro_types::RunEngine", &[]),
|
||||
("PetriAdmission", "fabro_types::PetriAdmission", &[]),
|
||||
("PetriGraphRef", "fabro_types::PetriGraphRef", &[]),
|
||||
("PullRequest", "fabro_types::PullRequest", &[]),
|
||||
("PullRequestLink", "fabro_types::PullRequestLink", &[]),
|
||||
(
|
||||
|
|
|
|||
|
|
@ -52,25 +52,25 @@ pub mod types {
|
|||
ModelTestMode, ModelUsage, PairId, PairMessageId, PairMessageRecord, PairMessageRequest,
|
||||
PairRecord, PairStartRequest, PairStatus, PairTarget, PairTranscriptEntry,
|
||||
PairTranscriptResponse, ParallelBranchId, ParallelBranchResult, PendingInterviewRecord,
|
||||
PermissionLevel, Principal, Provider, PullRequest, PullRequestCreation,
|
||||
PullRequestCreationId, PullRequestCreationStatus, PullRequestDetails,
|
||||
PermissionLevel, PetriAdmission, PetriGraphRef, Principal, Provider, PullRequest,
|
||||
PullRequestCreation, PullRequestCreationId, PullRequestCreationStatus, PullRequestDetails,
|
||||
PullRequestDetailsStatus, PullRequestDetailsUnavailableReason, PullRequestLink,
|
||||
PullRequestMeta, PullRequestResponse, QuestionType, RepositoryRef, ReviewTarget,
|
||||
ReviewTargetKind, Run, RunApproval, RunApprovalState, RunClientProvenance, RunEvent,
|
||||
RunEventDetailContentKind, RunEventDetailResponse, RunFailure, RunIntent, RunIntentArgs,
|
||||
RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource, RunSandbox,
|
||||
RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan, RunSandboxRuntime,
|
||||
RunServerProvenance, RunSessionMetadata, RunSize, RunTarget, SandboxDetails, SandboxInfo,
|
||||
SandboxListMeta, SandboxListResponse, SandboxProviderKind, SandboxProviderLookupError,
|
||||
SandboxService, SandboxServiceListResponse, SecretMetadata, SecretType, ServerSettings,
|
||||
SessionDetail, SessionId, SessionStatus, SessionSummary, SessionTurn,
|
||||
SkillActivationSource, SkillSummary, StageCompletion, StageContextWindow,
|
||||
StageContextWindowUnavailableReason, StageHandler, StageId, StageInferenceProjection,
|
||||
StageModelUsage, StageOutcome, StageProjection, StageState, StageToolBatchProjection,
|
||||
SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection,
|
||||
ToolCategory, ToolSource, ToolSummary, TurnId, UpdateVariableRequest, UserPrincipal,
|
||||
Variable, VariableListResponse, WorkflowPath, WorkflowSettings, WorkflowVersion,
|
||||
WorkflowVersionId,
|
||||
ReviewTargetKind, Run, RunApproval, RunApprovalState, RunClientProvenance, RunEngine,
|
||||
RunEvent, RunEventDetailContentKind, RunEventDetailResponse, RunFailure, RunIntent,
|
||||
RunIntentArgs, RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource,
|
||||
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,
|
||||
RunSandboxRuntime, RunServerProvenance, RunSessionMetadata, RunSize, RunStreamItem,
|
||||
RunStreamItemKind, RunTarget, SandboxDetails, SandboxInfo, SandboxListMeta,
|
||||
SandboxListResponse, SandboxProviderKind, SandboxProviderLookupError, SandboxService,
|
||||
SandboxServiceListResponse, SecretMetadata, SecretType, ServerSettings, SessionDetail,
|
||||
SessionId, SessionStatus, SessionSummary, SessionTurn, SkillActivationSource, SkillSummary,
|
||||
StageCompletion, StageContextWindow, StageContextWindowUnavailableReason, StageHandler,
|
||||
StageId, StageInferenceProjection, StageModelUsage, StageOutcome, StageProjection,
|
||||
StageState, StageToolBatchProjection, SystemActorKind, SystemIntegrationStatus,
|
||||
SystemIntegrationsResponse, TodoListProjection, ToolCategory, ToolSource, ToolSummary,
|
||||
TurnId, UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowPath,
|
||||
WorkflowSettings, WorkflowVersion, WorkflowVersionId,
|
||||
};
|
||||
pub use lithos_llm::catalog::{ModelHandle, ProviderId};
|
||||
pub use lithos_llm::types::{
|
||||
|
|
|
|||
57
lib/foundation/fabro-api/tests/run_engine_round_trip.rs
Normal file
57
lib/foundation/fabro-api/tests/run_engine_round_trip.rs
Normal file
|
|
@ -0,0 +1,57 @@
|
|||
use std::any::{TypeId, type_name};
|
||||
|
||||
use fabro_api::types::{
|
||||
PetriAdmission as ApiPetriAdmission, PetriGraphRef as ApiPetriGraphRef,
|
||||
RunEngine as ApiRunEngine,
|
||||
};
|
||||
use fabro_types::{PetriAdmission, PetriGraphRef, RunEngine};
|
||||
use serde_json::json;
|
||||
|
||||
#[test]
|
||||
fn run_engine_reuses_canonical_types() {
|
||||
assert_same_type::<ApiRunEngine, RunEngine>();
|
||||
assert_same_type::<ApiPetriAdmission, PetriAdmission>();
|
||||
assert_same_type::<ApiPetriGraphRef, PetriGraphRef>();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_legacy_engine_round_trips_as_its_kind_alone() {
|
||||
let value = json!({ "kind": "legacy" });
|
||||
let engine: RunEngine = serde_json::from_value(value.clone()).unwrap();
|
||||
assert!(engine.is_legacy());
|
||||
assert_eq!(serde_json::to_value(&engine).unwrap(), value);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_petri_engine_round_trips_with_its_admission_flattened() {
|
||||
let value = json!({
|
||||
"kind": "petri",
|
||||
"graph": {
|
||||
"blob": "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824",
|
||||
"digest": "sha256:root"
|
||||
},
|
||||
"children": [
|
||||
{
|
||||
"blob": "3cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824",
|
||||
"digest": "sha256:child"
|
||||
}
|
||||
]
|
||||
});
|
||||
let engine: RunEngine = serde_json::from_value(value.clone()).unwrap();
|
||||
let admission = engine
|
||||
.petri()
|
||||
.expect("a Petri engine carries its admission");
|
||||
assert_eq!(admission.graph.digest, "sha256:root");
|
||||
assert_eq!(admission.children.len(), 1);
|
||||
assert_eq!(serde_json::to_value(&engine).unwrap(), value);
|
||||
}
|
||||
|
||||
fn assert_same_type<T: 'static, U: 'static>() {
|
||||
assert_eq!(
|
||||
TypeId::of::<T>(),
|
||||
TypeId::of::<U>(),
|
||||
"{} should be the same type as {}",
|
||||
type_name::<T>(),
|
||||
type_name::<U>()
|
||||
);
|
||||
}
|
||||
74
lib/foundation/fabro-api/tests/run_stream_item_round_trip.rs
Normal file
74
lib/foundation/fabro-api/tests/run_stream_item_round_trip.rs
Normal file
|
|
@ -0,0 +1,74 @@
|
|||
use std::any::{TypeId, type_name};
|
||||
|
||||
use fabro_api::types::{
|
||||
RunStreamItem as ApiRunStreamItem, RunStreamItemKind as ApiRunStreamItemKind,
|
||||
};
|
||||
use fabro_types::{RunStreamItem, RunStreamItemKind, fixtures};
|
||||
use serde_json::json;
|
||||
|
||||
#[test]
|
||||
fn run_stream_item_reuses_canonical_types() {
|
||||
assert_same_type::<ApiRunStreamItem, RunStreamItem>();
|
||||
assert_same_type::<ApiRunStreamItemKind, RunStreamItemKind>();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_petri_item_round_trips_with_its_event_unchanged() {
|
||||
let value = json!({
|
||||
"run_id": fixtures::RUN_1.to_string(),
|
||||
"stream_seq": 12,
|
||||
"kind": "petri",
|
||||
"id": "execution 1/23/0",
|
||||
"recorded_at": 1_789_323_217_459_u64,
|
||||
"item": {
|
||||
"id": { "log": "execution", "execution": 1, "seq": 23, "index": 0 },
|
||||
"origin": "core",
|
||||
"recorded_at": 1_789_323_217_459_u64,
|
||||
"context": { "invocation": 0, "execution": 1 },
|
||||
"subject": {
|
||||
"node": { "id": 4, "name": "review", "kind": "attractor/agent", "meta": { "kind": "agent" } },
|
||||
"firing": 3, "visit": 1, "attempt": 1, "generation": 0, "branch": { "role": "none" }
|
||||
},
|
||||
"record": {
|
||||
"seq": 23, "origin": "core", "recorded_at": 1_789_323_217_459_u64,
|
||||
"body": { "event": "route.applied", "kind": "jump", "firing": 3, "target": 7 }
|
||||
},
|
||||
"derived": { "target": { "id": 7, "name": "finalize", "kind": "attractor/command", "meta": { "kind": "command" } } }
|
||||
}
|
||||
});
|
||||
let item: RunStreamItem = serde_json::from_value(value.clone()).unwrap();
|
||||
assert_eq!(item.kind, RunStreamItemKind::Petri);
|
||||
assert_eq!(item.name(), Some("route.applied"));
|
||||
assert_eq!(serde_json::to_value(&item).unwrap(), value);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_platform_item_round_trips_with_its_record_unchanged() {
|
||||
let value = json!({
|
||||
"run_id": fixtures::RUN_1.to_string(),
|
||||
"stream_seq": 13,
|
||||
"kind": "platform",
|
||||
"id": "4",
|
||||
"recorded_at": 1_789_323_217_500_u64,
|
||||
"item": {
|
||||
"seq": 4,
|
||||
"recorded_at": 1_789_323_217_500_u64,
|
||||
"record": { "kind": "checkpoint", "execution": 1, "firing": 3, "git_commit_sha": "abc123" },
|
||||
"position": { "execution": 1, "firing": 3 }
|
||||
}
|
||||
});
|
||||
let item: RunStreamItem = serde_json::from_value(value.clone()).unwrap();
|
||||
assert_eq!(item.kind, RunStreamItemKind::Platform);
|
||||
assert_eq!(item.name(), Some("checkpoint"));
|
||||
assert_eq!(serde_json::to_value(&item).unwrap(), value);
|
||||
}
|
||||
|
||||
fn assert_same_type<T: 'static, U: 'static>() {
|
||||
assert_eq!(
|
||||
TypeId::of::<T>(),
|
||||
TypeId::of::<U>(),
|
||||
"{} should be the same type as {}",
|
||||
type_name::<T>(),
|
||||
type_name::<U>()
|
||||
);
|
||||
}
|
||||
|
|
@ -15,7 +15,7 @@ use fabro_types::{
|
|||
ArtifactUpload, BlobHash, EventEnvelope, Model, ModelTestMode, PairId, PairMessageRecord,
|
||||
PairMessageRequest, PairRecord, PairStartRequest, PairTranscriptResponse, Run, RunEvent,
|
||||
RunEventDetailResponse, RunId, RunPairStatusResponse, RunProjection, RunSessionMetadata,
|
||||
SessionId, StageId, WorkflowVersion, WorkflowVersionId,
|
||||
RunStreamItem, SessionId, StageId, WorkflowVersion, WorkflowVersionId,
|
||||
};
|
||||
use fabro_util::exit::{ErrorExt, ExitClass};
|
||||
use futures::future::BoxFuture;
|
||||
|
|
@ -56,6 +56,24 @@ pub struct RunEventStream {
|
|||
buffered_events: VecDeque<EventEnvelope>,
|
||||
}
|
||||
|
||||
/// The live stream of a Petri run, as `GET /runs/{id}/attach` serves it:
|
||||
/// one `RunStreamItem` per `data:` frame, in `stream_seq` order.
|
||||
pub struct RunStreamItemStream {
|
||||
stream: progenitor_client::ByteStream,
|
||||
pending_bytes: Vec<u8>,
|
||||
buffered_items: VecDeque<RunStreamItem>,
|
||||
}
|
||||
|
||||
/// One page of a Petri run's stream.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct RunStreamPage {
|
||||
pub items: Vec<RunStreamItem>,
|
||||
pub has_more: bool,
|
||||
/// Petri's `EVENT_CONTRACT_VERSION` the server serves; `None` when the
|
||||
/// page was empty and the server reported no version beside it.
|
||||
pub event_contract_version: Option<u32>,
|
||||
}
|
||||
|
||||
type HttpByteStream = Pin<Box<dyn Stream<Item = Result<Bytes>> + Send>>;
|
||||
|
||||
pub struct SessionEventStream {
|
||||
|
|
@ -186,6 +204,42 @@ impl RunEventStream {
|
|||
}
|
||||
}
|
||||
|
||||
impl RunStreamItemStream {
|
||||
#[must_use]
|
||||
pub fn new(stream: progenitor_client::ByteStream) -> Self {
|
||||
Self {
|
||||
stream,
|
||||
pending_bytes: Vec::new(),
|
||||
buffered_items: VecDeque::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn next_item(&mut self) -> Result<Option<RunStreamItem>> {
|
||||
loop {
|
||||
if let Some(item) = self.buffered_items.pop_front() {
|
||||
return Ok(Some(item));
|
||||
}
|
||||
|
||||
if let Some(chunk) = self.stream.next().await {
|
||||
let chunk = chunk.map_err(anyhow::Error::new)?;
|
||||
self.pending_bytes.extend_from_slice(&chunk);
|
||||
self.buffer_sse_items(false)?;
|
||||
} else {
|
||||
self.buffer_sse_items(true)?;
|
||||
return Ok(self.buffered_items.pop_front());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn buffer_sse_items(&mut self, finalize: bool) -> Result<()> {
|
||||
for payload in sse::drain_sse_payloads(&mut self.pending_bytes, finalize) {
|
||||
self.buffered_items
|
||||
.push_back(serde_json::from_str(&payload)?);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl SessionEventStream {
|
||||
#[must_use]
|
||||
pub fn new(stream: HttpByteStream) -> Self {
|
||||
|
|
@ -1813,7 +1867,15 @@ impl Client {
|
|||
request.send().await
|
||||
})
|
||||
.await?;
|
||||
let parsed = response.into_inner();
|
||||
let parsed = match response.into_inner() {
|
||||
types::ListRunEventsResponse::EventList(page) => page,
|
||||
types::ListRunEventsResponse::RunStreamList(_) => {
|
||||
bail!(
|
||||
"run {run_id} executes on Petri; its events are served as a run stream \
|
||||
(list_run_stream)"
|
||||
);
|
||||
}
|
||||
};
|
||||
let events = parsed
|
||||
.data
|
||||
.into_iter()
|
||||
|
|
@ -1822,6 +1884,84 @@ impl Client {
|
|||
Ok((events, parsed.meta.has_more))
|
||||
}
|
||||
|
||||
/// One page of a Petri run's stream: up to `limit` items with
|
||||
/// `stream_seq > after`, in order.
|
||||
pub async fn list_run_stream_page(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
after: u64,
|
||||
limit: Option<usize>,
|
||||
) -> Result<RunStreamPage> {
|
||||
let response = self
|
||||
.send_api(|client| async move {
|
||||
let mut request = client.list_run_events().id(run_id.to_string()).after(after);
|
||||
let page_limit = limit.map(|limit| limit.min(1000));
|
||||
if let Some(limit) = page_limit.and_then(non_zero_u64_from_usize) {
|
||||
request = request.limit(limit);
|
||||
}
|
||||
request.send().await
|
||||
})
|
||||
.await?;
|
||||
match response.into_inner() {
|
||||
types::ListRunEventsResponse::RunStreamList(page) => Ok(RunStreamPage {
|
||||
items: page.data,
|
||||
has_more: page.meta.has_more,
|
||||
event_contract_version: Some(page.event_contract_version),
|
||||
}),
|
||||
// An empty page decodes as either list; a legacy page with
|
||||
// items is a run that does not execute on Petri.
|
||||
types::ListRunEventsResponse::EventList(page) if page.data.is_empty() => {
|
||||
Ok(RunStreamPage {
|
||||
items: Vec::new(),
|
||||
has_more: page.meta.has_more,
|
||||
event_contract_version: None,
|
||||
})
|
||||
}
|
||||
types::ListRunEventsResponse::EventList(_) => {
|
||||
bail!("run {run_id} executes on the legacy engine; its events are not a run stream")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Every item of a Petri run's stream past `after`, page by page.
|
||||
pub async fn list_run_stream(&self, run_id: &RunId, after: u64) -> Result<Vec<RunStreamItem>> {
|
||||
let mut cursor = after;
|
||||
let mut all = Vec::new();
|
||||
loop {
|
||||
let page = self.list_run_stream_page(run_id, cursor, None).await?;
|
||||
let Some(last) = page.items.last() else {
|
||||
break;
|
||||
};
|
||||
cursor = last.stream_seq;
|
||||
let has_more = page.has_more;
|
||||
all.extend(page.items);
|
||||
if !has_more {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Ok(all)
|
||||
}
|
||||
|
||||
/// The live stream of a Petri run from `after` (the last `stream_seq`
|
||||
/// seen; `Some(0)` replays the whole run; `None` starts at the next
|
||||
/// unseen item).
|
||||
pub async fn attach_run_stream(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
after: Option<u64>,
|
||||
) -> Result<RunStreamItemStream> {
|
||||
let response = self
|
||||
.send_api(|client| async move {
|
||||
let mut request = client.attach_run_events().id(run_id.to_string());
|
||||
if let Some(after) = after {
|
||||
request = request.after(after);
|
||||
}
|
||||
request.send().await
|
||||
})
|
||||
.await?;
|
||||
Ok(RunStreamItemStream::new(response.into_inner()))
|
||||
}
|
||||
|
||||
pub async fn attach_run_events(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
|
|
|
|||
|
|
@ -12,7 +12,8 @@ pub use auth_store::{
|
|||
AuthEntry, AuthStore, AuthStoreError, DevTokenEntry, LockError, OAuthEntry, StoredSubject,
|
||||
};
|
||||
pub use client::{
|
||||
Client, RunEventStream, SessionEventStream, TransportConnector, apply_bearer_token_auth,
|
||||
Client, RunEventStream, RunStreamItemStream, RunStreamPage, SessionEventStream,
|
||||
TransportConnector, apply_bearer_token_auth,
|
||||
};
|
||||
pub use credential::{Credential, CredentialFallback};
|
||||
pub use error::{
|
||||
|
|
|
|||
|
|
@ -35,6 +35,7 @@ pub mod run_id;
|
|||
pub mod run_intent;
|
||||
pub mod run_projection;
|
||||
pub mod run_sandbox;
|
||||
pub mod run_stream;
|
||||
pub mod run_summary;
|
||||
pub mod run_title;
|
||||
pub mod sandbox_details;
|
||||
|
|
@ -152,6 +153,7 @@ pub use run_sandbox::{
|
|||
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,
|
||||
RunSandboxRuntime,
|
||||
};
|
||||
pub use run_stream::{RunStreamItem, RunStreamItemKind, petri_event_name};
|
||||
pub use run_summary::{
|
||||
AskFabro, AskFabroUnavailableReason, AutomationRef, ResolvedAutomationGitWorkflowSource, Run,
|
||||
RunApproval, RunApprovalState, RunError, RunLifecycle, RunLinks, RunModel, RunOrigin,
|
||||
|
|
|
|||
168
lib/foundation/fabro-types/src/run_stream.rs
Normal file
168
lib/foundation/fabro-types/src/run_stream.rs
Normal file
|
|
@ -0,0 +1,168 @@
|
|||
//! The run stream: one ordered delivery of a Petri run's public events and
|
||||
//! Fabro's platform records, as `GET /runs/{id}/events` and the attach
|
||||
//! stream serve them for a run that executes on Petri.
|
||||
//!
|
||||
//! Each item is a Petri `RunEvent` (Petri's event contract, passed through
|
||||
//! as JSON) or a stored platform record (Fabro's own fact about the run:
|
||||
//! its lifecycle before and after the engine, a checkpoint with its commit,
|
||||
//! a pull request), in one Fabro envelope. The envelope carries:
|
||||
//!
|
||||
//! - `stream_seq`, the durable per-run delivery sequence the projector assigned
|
||||
//! when the item's record was committed. It is the cursor: a client resumes
|
||||
//! from the last `stream_seq` it saw. It is dense and strictly increasing
|
||||
//! within a run.
|
||||
//! - `id`, the item's own identity, kept beside the cursor so a client
|
||||
//! deduplicates by it: for a Petri event the `EventId` as
|
||||
//! `<log>/<seq>/<index>` (`coordinator/3/0`, `execution 1/23/0`), for a
|
||||
//! platform record its `seq`. Petri's `EventId` is per log and has no
|
||||
//! platform variant, so it is never the cursor.
|
||||
//! - `kind`, which of the two the item is.
|
||||
//! - `recorded_at`, when the item's record was appended, in milliseconds since
|
||||
//! the Unix epoch; the same field both item shapes carry.
|
||||
//! - `item`, the Petri `RunEvent` or the stored platform record, unchanged.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::RunId;
|
||||
|
||||
/// Which of the two item shapes a stream item carries.
|
||||
#[derive(
|
||||
Debug,
|
||||
Clone,
|
||||
Copy,
|
||||
PartialEq,
|
||||
Eq,
|
||||
Hash,
|
||||
Serialize,
|
||||
Deserialize,
|
||||
strum::Display,
|
||||
strum::EnumString,
|
||||
strum::IntoStaticStr,
|
||||
)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
#[strum(serialize_all = "lowercase")]
|
||||
pub enum RunStreamItemKind {
|
||||
/// A Petri `RunEvent`: `{id, origin, recorded_at, context, subject,
|
||||
/// record, derived}` under Petri's event contract.
|
||||
Petri,
|
||||
/// A stored platform record: `{seq, recorded_at, record: {kind, ...},
|
||||
/// position?}`.
|
||||
Platform,
|
||||
}
|
||||
|
||||
/// One item of a Petri run's stream, in Fabro's envelope.
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct RunStreamItem {
|
||||
pub run_id: RunId,
|
||||
/// The delivery sequence: the cursor.
|
||||
pub stream_seq: u64,
|
||||
pub kind: RunStreamItemKind,
|
||||
/// The item's own identity, for deduplication.
|
||||
pub id: String,
|
||||
/// Milliseconds since the Unix epoch when the item's record was
|
||||
/// appended.
|
||||
pub recorded_at: u64,
|
||||
/// The Petri `RunEvent` or the stored platform record, as JSON.
|
||||
pub item: serde_json::Value,
|
||||
}
|
||||
|
||||
impl RunStreamItem {
|
||||
/// The `<subject>.<verb>` name of a Petri event, or the `kind` of a
|
||||
/// platform record: what a listing shows and a filter matches on.
|
||||
#[must_use]
|
||||
pub fn name(&self) -> Option<&str> {
|
||||
match self.kind {
|
||||
RunStreamItemKind::Petri => petri_event_name(&self.item),
|
||||
RunStreamItemKind::Platform => self
|
||||
.item
|
||||
.get("record")
|
||||
.and_then(|record| record.get("kind"))
|
||||
.and_then(serde_json::Value::as_str),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The `<subject>.<verb>` name of a Petri `RunEvent` value: the recorded
|
||||
/// body's `event` tag, or a view event's tag under `derived`.
|
||||
#[must_use]
|
||||
pub fn petri_event_name(event: &serde_json::Value) -> Option<&str> {
|
||||
event
|
||||
.get("record")
|
||||
.and_then(|record| record.get("body"))
|
||||
.and_then(|body| body.get("event"))
|
||||
.and_then(serde_json::Value::as_str)
|
||||
.or_else(|| {
|
||||
event
|
||||
.get("derived")
|
||||
.and_then(|derived| derived.get("event"))
|
||||
.and_then(serde_json::Value::as_str)
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use serde_json::json;
|
||||
|
||||
use super::{RunStreamItem, RunStreamItemKind, petri_event_name};
|
||||
use crate::fixtures;
|
||||
|
||||
#[test]
|
||||
fn kind_names_are_lowercase_in_both_directions() {
|
||||
assert_eq!(RunStreamItemKind::Petri.to_string(), "petri");
|
||||
assert_eq!(
|
||||
"platform".parse::<RunStreamItemKind>(),
|
||||
Ok(RunStreamItemKind::Platform)
|
||||
);
|
||||
assert_eq!(
|
||||
serde_json::to_value(RunStreamItemKind::Platform).expect("kind serializes"),
|
||||
json!("platform")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_petri_item_is_named_by_its_recorded_event_tag() {
|
||||
let item = RunStreamItem {
|
||||
run_id: fixtures::RUN_1,
|
||||
stream_seq: 4,
|
||||
kind: RunStreamItemKind::Petri,
|
||||
id: "coordinator/3/0".to_string(),
|
||||
recorded_at: 1_789_323_217_366,
|
||||
item: json!({
|
||||
"id": {"log": "coordinator", "seq": 3, "index": 0},
|
||||
"record": {"seq": 3, "body": {"event": "execution.declared"}}
|
||||
}),
|
||||
};
|
||||
assert_eq!(item.name(), Some("execution.declared"));
|
||||
assert_eq!(
|
||||
petri_event_name(&json!({"origin": "derived", "derived": {"event": "visit.started"}})),
|
||||
Some("visit.started")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_platform_item_is_named_by_its_record_kind() {
|
||||
let item = RunStreamItem {
|
||||
run_id: fixtures::RUN_1,
|
||||
stream_seq: 5,
|
||||
kind: RunStreamItemKind::Platform,
|
||||
id: "2".to_string(),
|
||||
recorded_at: 1_789_323_217_400,
|
||||
item: json!({"seq": 2, "record": {"kind": "run.notice", "level": "info"}}),
|
||||
};
|
||||
assert_eq!(item.name(), Some("run.notice"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_envelope_round_trips_as_json() {
|
||||
let value = json!({
|
||||
"run_id": fixtures::RUN_1.to_string(),
|
||||
"stream_seq": 7,
|
||||
"kind": "platform",
|
||||
"id": "3",
|
||||
"recorded_at": 1_789_323_217_400_u64,
|
||||
"item": {"seq": 3, "recorded_at": 1_789_323_217_400_u64, "record": {"kind": "checkpoint", "execution": 1, "firing": 2}}
|
||||
});
|
||||
let item: RunStreamItem = serde_json::from_value(value.clone()).expect("item decodes");
|
||||
assert_eq!(serde_json::to_value(&item).expect("item encodes"), value);
|
||||
}
|
||||
}
|
||||
|
|
@ -216,6 +216,7 @@ models/interview-option.ts
|
|||
models/interview-provider-settings.ts
|
||||
models/interview-question-record.ts
|
||||
models/link-run-pull-request-request.ts
|
||||
models/list-run-events200-response.ts
|
||||
models/llm-output-kind.ts
|
||||
models/llm-retry-classification-after.ts
|
||||
models/llm-retry-classification-never.ts
|
||||
|
|
@ -271,6 +272,7 @@ models/paginated-run-commit-list.ts
|
|||
models/paginated-run-file-list.ts
|
||||
models/paginated-run-list.ts
|
||||
models/paginated-run-stage-list.ts
|
||||
models/paginated-run-stream-list.ts
|
||||
models/paginated-saved-query-list.ts
|
||||
models/paginated-session-list.ts
|
||||
models/paginated-workflow-list-response.ts
|
||||
|
|
@ -296,7 +298,9 @@ models/pending-interview-record.ts
|
|||
models/pending-reason.ts
|
||||
models/permission-level.ts
|
||||
models/petri-access.ts
|
||||
models/petri-admission.ts
|
||||
models/petri-append-request.ts
|
||||
models/petri-graph-ref.ts
|
||||
models/petri-open-request.ts
|
||||
models/petri-open-response.ts
|
||||
models/petri-record-list.ts
|
||||
|
|
@ -383,6 +387,9 @@ models/run-commit.ts
|
|||
models/run-commits-meta.ts
|
||||
models/run-control-action.ts
|
||||
models/run-diff.ts
|
||||
models/run-engine-one-of.ts
|
||||
models/run-engine-one-of1.ts
|
||||
models/run-engine.ts
|
||||
models/run-environment-settings.ts
|
||||
models/run-error.ts
|
||||
models/run-event-detail-response-content.ts
|
||||
|
|
@ -442,6 +449,8 @@ models/run-status-starting.ts
|
|||
models/run-status-submitted.ts
|
||||
models/run-status-succeeded.ts
|
||||
models/run-status.ts
|
||||
models/run-stream-item-kind.ts
|
||||
models/run-stream-item.ts
|
||||
models/run-superseded-by-props.ts
|
||||
models/run-target.ts
|
||||
models/run-timestamps.ts
|
||||
|
|
|
|||
|
|
@ -30,6 +30,8 @@ import type { CommandLogResponse } from '../models';
|
|||
// @ts-ignore
|
||||
import type { ErrorResponse } from '../models';
|
||||
// @ts-ignore
|
||||
import type { ListRunEvents200Response } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PaginatedEventList } from '../models';
|
||||
// @ts-ignore
|
||||
import type { PaginatedRunStageList } from '../models';
|
||||
|
|
@ -161,14 +163,15 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config
|
|||
};
|
||||
},
|
||||
/**
|
||||
* Opens an ordered server-sent event stream starting at `since_seq`, replaying persisted events and continuing with live updates while the run remains active.
|
||||
* Opens an ordered server-sent event stream, replaying persisted items and continuing with live updates while the run remains active. Each `data:` frame is one JSON object in the run engine\'s envelope. For a legacy run the frames are `EventEnvelope`s and the stream starts at `since_seq` (inclusive; the next unseen event when omitted). It ends after `run.completed` or `run.failed`. For a Petri run the frames are `RunStreamItem`s and the stream starts after `after` (the last `stream_seq` the client saw; `0` replays the whole run; the next unseen item when omitted). It ends once the run is no longer active and every committed item has been sent. A reconnecting client passes its last `stream_seq` as `after` and deduplicates by `id`.
|
||||
* @summary Attach Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
attachRunEvents: async (id: string, sinceSeq?: number, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
attachRunEvents: async (id: string, sinceSeq?: number, after?: number, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
// verify required parameter 'id' is not null or undefined
|
||||
assertParamExists('attachRunEvents', 'id', id)
|
||||
const localVarPath = `/api/v1/runs/{id}/attach`
|
||||
|
|
@ -194,6 +197,10 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config
|
|||
localVarQueryParameter['since_seq'] = sinceSeq;
|
||||
}
|
||||
|
||||
if (after !== undefined) {
|
||||
localVarQueryParameter['after'] = after;
|
||||
}
|
||||
|
||||
localVarHeaderParameter['Accept'] = 'text/event-stream,application/json';
|
||||
|
||||
setSearchParams(localVarUrlObj, localVarQueryParameter);
|
||||
|
|
@ -615,17 +622,18 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config
|
|||
};
|
||||
},
|
||||
/**
|
||||
* Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted.
|
||||
* Returns a paginated JSON list of the run\'s events. The shape depends on the engine the run was created for (`RunSpec.engine`). For a legacy run (`engine.kind = legacy`): stored run events in the legacy envelope (`PaginatedEventList`). Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. For a Petri run (`engine.kind = petri`): the run stream (`PaginatedRunStreamList`), one ordered delivery of Petri\'s own `RunEvent`s and Fabro\'s platform records in the `RunStreamItem` envelope, in `stream_seq` order. The cursor is `after`: the last `stream_seq` the client saw, exclusive; the first page is `after=0`. `since_seq`, `before_seq` and `order` are not accepted for a Petri run. A client that reconnects resumes from its last `stream_seq` and deduplicates by each item\'s `id`; every item is delivered once, in order, with no gap.
|
||||
* @summary List Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [limit] Maximum number of events to return.
|
||||
* @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event.
|
||||
* @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
listRunEvents: async (id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
listRunEvents: async (id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, after?: number, options: RawAxiosRequestConfig = {}): Promise<RequestArgs> => {
|
||||
// verify required parameter 'id' is not null or undefined
|
||||
assertParamExists('listRunEvents', 'id', id)
|
||||
const localVarPath = `/api/v1/runs/{id}/events`
|
||||
|
|
@ -663,6 +671,10 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config
|
|||
localVarQueryParameter['order'] = order;
|
||||
}
|
||||
|
||||
if (after !== undefined) {
|
||||
localVarQueryParameter['after'] = after;
|
||||
}
|
||||
|
||||
localVarHeaderParameter['Accept'] = 'application/json';
|
||||
|
||||
setSearchParams(localVarUrlObj, localVarQueryParameter);
|
||||
|
|
@ -1276,15 +1288,16 @@ export const RunInternalsApiFp = function(configuration?: Configuration) {
|
|||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Opens an ordered server-sent event stream starting at `since_seq`, replaying persisted events and continuing with live updates while the run remains active.
|
||||
* Opens an ordered server-sent event stream, replaying persisted items and continuing with live updates while the run remains active. Each `data:` frame is one JSON object in the run engine\'s envelope. For a legacy run the frames are `EventEnvelope`s and the stream starts at `since_seq` (inclusive; the next unseen event when omitted). It ends after `run.completed` or `run.failed`. For a Petri run the frames are `RunStreamItem`s and the stream starts after `after` (the last `stream_seq` the client saw; `0` replays the whole run; the next unseen item when omitted). It ends once the run is no longer active and every committed item has been sent. A reconnecting client passes its last `stream_seq` as `after` and deduplicates by `id`.
|
||||
* @summary Attach Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
async attachRunEvents(id: string, sinceSeq?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<string>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.attachRunEvents(id, sinceSeq, options);
|
||||
async attachRunEvents(id: string, sinceSeq?: number, after?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<string>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.attachRunEvents(id, sinceSeq, after, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.attachRunEvents']?.[localVarOperationServerIndex]?.url;
|
||||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
|
|
@ -1417,18 +1430,19 @@ export const RunInternalsApiFp = function(configuration?: Configuration) {
|
|||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
},
|
||||
/**
|
||||
* Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted.
|
||||
* Returns a paginated JSON list of the run\'s events. The shape depends on the engine the run was created for (`RunSpec.engine`). For a legacy run (`engine.kind = legacy`): stored run events in the legacy envelope (`PaginatedEventList`). Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. For a Petri run (`engine.kind = petri`): the run stream (`PaginatedRunStreamList`), one ordered delivery of Petri\'s own `RunEvent`s and Fabro\'s platform records in the `RunStreamItem` envelope, in `stream_seq` order. The cursor is `after`: the last `stream_seq` the client saw, exclusive; the first page is `after=0`. `since_seq`, `before_seq` and `order` are not accepted for a Petri run. A client that reconnects resumes from its last `stream_seq` and deduplicates by each item\'s `id`; every item is delivered once, in order, with no gap.
|
||||
* @summary List Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [limit] Maximum number of events to return.
|
||||
* @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event.
|
||||
* @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
async listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<PaginatedEventList>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.listRunEvents(id, sinceSeq, limit, beforeSeq, order, options);
|
||||
async listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, after?: number, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise<ListRunEvents200Response>> {
|
||||
const localVarAxiosArgs = await localVarAxiosParamCreator.listRunEvents(id, sinceSeq, limit, beforeSeq, order, after, options);
|
||||
const localVarOperationServerIndex = configuration?.serverIndex ?? 0;
|
||||
const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.listRunEvents']?.[localVarOperationServerIndex]?.url;
|
||||
return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath);
|
||||
|
|
@ -1639,15 +1653,16 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b
|
|||
return localVarFp.appendRunEvent(id, runEvent, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Opens an ordered server-sent event stream starting at `since_seq`, replaying persisted events and continuing with live updates while the run remains active.
|
||||
* Opens an ordered server-sent event stream, replaying persisted items and continuing with live updates while the run remains active. Each `data:` frame is one JSON object in the run engine\'s envelope. For a legacy run the frames are `EventEnvelope`s and the stream starts at `since_seq` (inclusive; the next unseen event when omitted). It ends after `run.completed` or `run.failed`. For a Petri run the frames are `RunStreamItem`s and the stream starts after `after` (the last `stream_seq` the client saw; `0` replays the whole run; the next unseen item when omitted). It ends once the run is no longer active and every committed item has been sent. A reconnecting client passes its last `stream_seq` as `after` and deduplicates by `id`.
|
||||
* @summary Attach Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
attachRunEvents(id: string, sinceSeq?: number, options?: RawAxiosRequestConfig): AxiosPromise<string> {
|
||||
return localVarFp.attachRunEvents(id, sinceSeq, options).then((request) => request(axios, basePath));
|
||||
attachRunEvents(id: string, sinceSeq?: number, after?: number, options?: RawAxiosRequestConfig): AxiosPromise<string> {
|
||||
return localVarFp.attachRunEvents(id, sinceSeq, after, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Streams a ZIP archive with the latest captured version of each artifact path. Stage order, retry number, and then stage ID determine the latest version, matching the artifacts page. Captures from the graph\'s boundary nodes are excluded, identified by their `start` and `exit` handler type rather than by node name. The archive streams, so the response status is sent before the first artifact is read. A failure after that point aborts the transfer rather than returning `500`. The ZIP central directory is written last, so a truncated download does not open as a valid archive.
|
||||
|
|
@ -1750,18 +1765,19 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b
|
|||
return localVarFp.listRunArtifacts(id, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted.
|
||||
* Returns a paginated JSON list of the run\'s events. The shape depends on the engine the run was created for (`RunSpec.engine`). For a legacy run (`engine.kind = legacy`): stored run events in the legacy envelope (`PaginatedEventList`). Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. For a Petri run (`engine.kind = petri`): the run stream (`PaginatedRunStreamList`), one ordered delivery of Petri\'s own `RunEvent`s and Fabro\'s platform records in the `RunStreamItem` envelope, in `stream_seq` order. The cursor is `after`: the last `stream_seq` the client saw, exclusive; the first page is `after=0`. `since_seq`, `before_seq` and `order` are not accepted for a Petri run. A client that reconnects resumes from its last `stream_seq` and deduplicates by each item\'s `id`; every item is delivered once, in order, with no gap.
|
||||
* @summary List Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [limit] Maximum number of events to return.
|
||||
* @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event.
|
||||
* @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options?: RawAxiosRequestConfig): AxiosPromise<PaginatedEventList> {
|
||||
return localVarFp.listRunEvents(id, sinceSeq, limit, beforeSeq, order, options).then((request) => request(axios, basePath));
|
||||
listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, after?: number, options?: RawAxiosRequestConfig): AxiosPromise<ListRunEvents200Response> {
|
||||
return localVarFp.listRunEvents(id, sinceSeq, limit, beforeSeq, order, after, options).then((request) => request(axios, basePath));
|
||||
},
|
||||
/**
|
||||
* Returns the ordered list of stages in a run\'s workflow graph with their current status and timing. Stages are bounded by the workflow graph size, typically fewer than 20.
|
||||
|
|
@ -1933,15 +1949,16 @@ export class RunInternalsApi extends BaseAPI {
|
|||
}
|
||||
|
||||
/**
|
||||
* Opens an ordered server-sent event stream starting at `since_seq`, replaying persisted events and continuing with live updates while the run remains active.
|
||||
* Opens an ordered server-sent event stream, replaying persisted items and continuing with live updates while the run remains active. Each `data:` frame is one JSON object in the run engine\'s envelope. For a legacy run the frames are `EventEnvelope`s and the stream starts at `since_seq` (inclusive; the next unseen event when omitted). It ends after `run.completed` or `run.failed`. For a Petri run the frames are `RunStreamItem`s and the stream starts after `after` (the last `stream_seq` the client saw; `0` replays the whole run; the next unseen item when omitted). It ends once the run is no longer active and every committed item has been sent. A reconnecting client passes its last `stream_seq` as `after` and deduplicates by `id`.
|
||||
* @summary Attach Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
public attachRunEvents(id: string, sinceSeq?: number, options?: RawAxiosRequestConfig) {
|
||||
return RunInternalsApiFp(this.configuration).attachRunEvents(id, sinceSeq, options).then((request) => request(this.axios, this.basePath));
|
||||
public attachRunEvents(id: string, sinceSeq?: number, after?: number, options?: RawAxiosRequestConfig) {
|
||||
return RunInternalsApiFp(this.configuration).attachRunEvents(id, sinceSeq, after, options).then((request) => request(this.axios, this.basePath));
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -2054,18 +2071,19 @@ export class RunInternalsApi extends BaseAPI {
|
|||
}
|
||||
|
||||
/**
|
||||
* Returns a paginated JSON list of stored run events. Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted.
|
||||
* Returns a paginated JSON list of the run\'s events. The shape depends on the engine the run was created for (`RunSpec.engine`). For a legacy run (`engine.kind = legacy`): stored run events in the legacy envelope (`PaginatedEventList`). Ascending order uses `since_seq` as an inclusive cursor. Descending order uses `before_seq` as an exclusive cursor and starts at the newest event when `before_seq` is omitted. For a Petri run (`engine.kind = petri`): the run stream (`PaginatedRunStreamList`), one ordered delivery of Petri\'s own `RunEvent`s and Fabro\'s platform records in the `RunStreamItem` envelope, in `stream_seq` order. The cursor is `after`: the last `stream_seq` the client saw, exclusive; the first page is `after=0`. `since_seq`, `before_seq` and `order` are not accepted for a Petri run. A client that reconnects resumes from its last `stream_seq` and deduplicates by each item\'s `id`; every item is delivered once, in order, with no gap.
|
||||
* @summary List Run Events
|
||||
* @param {string} id Unique run identifier (ULID).
|
||||
* @param {number} [sinceSeq] First event sequence number to include.
|
||||
* @param {number} [limit] Maximum number of events to return.
|
||||
* @param {number} [beforeSeq] Exclusive upper event sequence cursor for descending order. Omit on the first descending request to start from the newest event.
|
||||
* @param {ListRunEventsOrderEnum} [order] Event sequence order. `since_seq` is valid only with `asc`; `before_seq` is valid only with `desc`.
|
||||
* @param {number} [after] Run stream cursor for a Petri run: the last `stream_seq` the client saw, exclusive. `0` starts at the first item.
|
||||
* @param {*} [options] Override http request option.
|
||||
* @throws {RequiredError}
|
||||
*/
|
||||
public listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, options?: RawAxiosRequestConfig) {
|
||||
return RunInternalsApiFp(this.configuration).listRunEvents(id, sinceSeq, limit, beforeSeq, order, options).then((request) => request(this.axios, this.basePath));
|
||||
public listRunEvents(id: string, sinceSeq?: number, limit?: number, beforeSeq?: number, order?: ListRunEventsOrderEnum, after?: number, options?: RawAxiosRequestConfig) {
|
||||
return RunInternalsApiFp(this.configuration).listRunEvents(id, sinceSeq, limit, beforeSeq, order, after, options).then((request) => request(this.axios, this.basePath));
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -186,6 +186,7 @@ export * from './interview-option';
|
|||
export * from './interview-provider-settings';
|
||||
export * from './interview-question-record';
|
||||
export * from './link-run-pull-request-request';
|
||||
export * from './list-run-events200-response';
|
||||
export * from './llm-output-kind';
|
||||
export * from './llm-retry-classification';
|
||||
export * from './llm-retry-classification-after';
|
||||
|
|
@ -241,6 +242,7 @@ export * from './paginated-run-commit-list';
|
|||
export * from './paginated-run-file-list';
|
||||
export * from './paginated-run-list';
|
||||
export * from './paginated-run-stage-list';
|
||||
export * from './paginated-run-stream-list';
|
||||
export * from './paginated-saved-query-list';
|
||||
export * from './paginated-session-list';
|
||||
export * from './paginated-workflow-list-response';
|
||||
|
|
@ -266,7 +268,9 @@ export * from './pending-interview-record';
|
|||
export * from './pending-reason';
|
||||
export * from './permission-level';
|
||||
export * from './petri-access';
|
||||
export * from './petri-admission';
|
||||
export * from './petri-append-request';
|
||||
export * from './petri-graph-ref';
|
||||
export * from './petri-open-request';
|
||||
export * from './petri-open-response';
|
||||
export * from './petri-record';
|
||||
|
|
@ -354,6 +358,9 @@ export * from './run-commit-person';
|
|||
export * from './run-commits-meta';
|
||||
export * from './run-control-action';
|
||||
export * from './run-diff';
|
||||
export * from './run-engine';
|
||||
export * from './run-engine-one-of';
|
||||
export * from './run-engine-one-of1';
|
||||
export * from './run-environment-settings';
|
||||
export * from './run-error';
|
||||
export * from './run-event';
|
||||
|
|
@ -413,6 +420,8 @@ export * from './run-status-running';
|
|||
export * from './run-status-starting';
|
||||
export * from './run-status-submitted';
|
||||
export * from './run-status-succeeded';
|
||||
export * from './run-stream-item';
|
||||
export * from './run-stream-item-kind';
|
||||
export * from './run-superseded-by-props';
|
||||
export * from './run-target';
|
||||
export * from './run-timestamps';
|
||||
|
|
|
|||
32
lib/packages/fabro-api-client/src/models/list-run-events200-response.ts
generated
Normal file
32
lib/packages/fabro-api-client/src/models/list-run-events200-response.ts
generated
Normal file
|
|
@ -0,0 +1,32 @@
|
|||
/* 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 { PaginatedEventList } from './paginated-event-list';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { PaginatedRunStreamList } from './paginated-run-stream-list';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { PaginationMeta } from './pagination-meta';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunStreamItem } from './run-stream-item';
|
||||
|
||||
/**
|
||||
* @type ListRunEvents200Response
|
||||
*/
|
||||
export type ListRunEvents200Response = PaginatedEventList | PaginatedRunStreamList;
|
||||
30
lib/packages/fabro-api-client/src/models/paginated-run-stream-list.ts
generated
Normal file
30
lib/packages/fabro-api-client/src/models/paginated-run-stream-list.ts
generated
Normal file
|
|
@ -0,0 +1,30 @@
|
|||
/* 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 { PaginationMeta } from './pagination-meta';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunStreamItem } from './run-stream-item';
|
||||
|
||||
/**
|
||||
* One page of a Petri run\'s stream, in `stream_seq` order. `event_contract_version` is Petri\'s `EVENT_CONTRACT_VERSION` the server was built against: the version of the event contract every `petri` item follows.
|
||||
*/
|
||||
export interface PaginatedRunStreamList {
|
||||
'data': Array<RunStreamItem>;
|
||||
'meta': PaginationMeta;
|
||||
'event_contract_version': number;
|
||||
}
|
||||
26
lib/packages/fabro-api-client/src/models/petri-admission.ts
generated
Normal file
26
lib/packages/fabro-api-client/src/models/petri-admission.ts
generated
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
/* 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 { PetriGraphRef } from './petri-graph-ref';
|
||||
|
||||
/**
|
||||
* What Petri admitted for a run at create time: the lowered root graph and the pre-lowered child graphs, every one persisted in the blob store before the run exists.
|
||||
*/
|
||||
export interface PetriAdmission {
|
||||
'graph': PetriGraphRef;
|
||||
'children'?: Array<PetriGraphRef>;
|
||||
}
|
||||
26
lib/packages/fabro-api-client/src/models/petri-graph-ref.ts
generated
Normal file
26
lib/packages/fabro-api-client/src/models/petri-graph-ref.ts
generated
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
/* 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.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
/**
|
||||
* An admitted graph in the blob store, verified by digest on load.
|
||||
*/
|
||||
export interface PetriGraphRef {
|
||||
/**
|
||||
* Content-addressed SHA-256 hash of a stored blob. Hex input is case-insensitive; Fabro emits the canonical lowercase form.
|
||||
*/
|
||||
'blob': string;
|
||||
'digest': string;
|
||||
}
|
||||
25
lib/packages/fabro-api-client/src/models/run-engine-one-of.ts
generated
Normal file
25
lib/packages/fabro-api-client/src/models/run-engine-one-of.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.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
export interface RunEngineOneOf {
|
||||
'kind': RunEngineOneOfKindEnum;
|
||||
}
|
||||
|
||||
export const RunEngineOneOfKindEnum = {
|
||||
LEGACY: 'legacy'
|
||||
} as const;
|
||||
|
||||
export type RunEngineOneOfKindEnum = typeof RunEngineOneOfKindEnum[keyof typeof RunEngineOneOfKindEnum];
|
||||
26
lib/packages/fabro-api-client/src/models/run-engine-one-of1.ts
generated
Normal file
26
lib/packages/fabro-api-client/src/models/run-engine-one-of1.ts
generated
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
/* 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 { PetriAdmission } from './petri-admission';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { PetriGraphRef } from './petri-graph-ref';
|
||||
|
||||
/**
|
||||
* @type RunEngineOneOf1
|
||||
*/
|
||||
export type RunEngineOneOf1 = PetriAdmission;
|
||||
30
lib/packages/fabro-api-client/src/models/run-engine.ts
generated
Normal file
30
lib/packages/fabro-api-client/src/models/run-engine.ts
generated
Normal file
|
|
@ -0,0 +1,30 @@
|
|||
/* 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 { PetriGraphRef } from './petri-graph-ref';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunEngineOneOf } from './run-engine-one-of';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunEngineOneOf1 } from './run-engine-one-of1';
|
||||
|
||||
/**
|
||||
* @type RunEngine
|
||||
* The engine a run was created for. `legacy` is the in-process executor; `petri` names the Petri workflow engine and carries what Petri admitted at create time.
|
||||
*/
|
||||
export type RunEngine = RunEngineOneOf | RunEngineOneOf1;
|
||||
|
|
@ -24,6 +24,9 @@ import type { ForkSourceRef } from './fork-source-ref';
|
|||
import type { GitContext } from './git-context';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunEngine } from './run-engine';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
import type { RunProvenance } from './run-provenance';
|
||||
// May contain unused imports in some cases
|
||||
// @ts-ignore
|
||||
|
|
@ -54,4 +57,8 @@ export interface RunSpec {
|
|||
'spec_blob'?: string | null;
|
||||
'git'?: GitContext | null;
|
||||
'fork_source_ref'?: ForkSourceRef | null;
|
||||
/**
|
||||
* The engine the run was created for, with what it admitted. Absent in a spec written before the field existed, which means the legacy executor.
|
||||
*/
|
||||
'engine'?: RunEngine;
|
||||
}
|
||||
|
|
|
|||
26
lib/packages/fabro-api-client/src/models/run-stream-item-kind.ts
generated
Normal file
26
lib/packages/fabro-api-client/src/models/run-stream-item-kind.ts
generated
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
/* 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.
|
||||
*/
|
||||
|
||||
|
||||
|
||||
/**
|
||||
* Which item shape a run stream item carries.
|
||||
*/
|
||||
|
||||
export const RunStreamItemKind = {
|
||||
PETRI: 'petri',
|
||||
PLATFORM: 'platform'
|
||||
} as const;
|
||||
|
||||
export type RunStreamItemKind = typeof RunStreamItemKind[keyof typeof RunStreamItemKind];
|
||||
42
lib/packages/fabro-api-client/src/models/run-stream-item.ts
generated
Normal file
42
lib/packages/fabro-api-client/src/models/run-stream-item.ts
generated
Normal file
|
|
@ -0,0 +1,42 @@
|
|||
/* 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 { RunStreamItemKind } from './run-stream-item-kind';
|
||||
|
||||
/**
|
||||
* One item of a Petri run\'s stream: a Petri `RunEvent` or a Fabro platform record in Fabro\'s envelope. `stream_seq` is the durable per-run delivery sequence the projector assigned when the item\'s record was committed: dense, strictly increasing within the run, and the cursor for `after`. `id` is the item\'s own identity, kept beside the cursor so a client deduplicates by it: for a Petri event the `EventId` as `<log>/<seq>/<index>` (`coordinator/3/0`, `execution 1/23/0`); for a platform record its `seq`. Petri\'s `EventId` is per log and has no platform variant, so it is never the cursor. A `petri` item is a Petri `RunEvent` passed through unchanged: `{id: {log, execution?, seq, index}, origin, recorded_at, observed_at?, context: {invocation, execution, parent?}, subject?, record?, derived?}`. Its vocabulary is Petri\'s public event contract (`crates/core/execution/EVENTS.md` in the Petri repository), not Fabro\'s: the recorded event\'s name is `record.body.event` (`<subject>.<verb>`, e.g. `visit.started`, `step.finished`, `run.finished`), a derived view event\'s is `derived.event`, and the stage a subject names is `(context.execution, subject.firing)` with `subject.node.name` and `subject.visit` as its display label. The server reports the contract version it serves in `PaginatedRunStreamList.event_contract_version`. A `platform` item is a stored platform record: `{seq, recorded_at, record: {kind, ...}, position?: {execution, firing}}`. `record.kind` is one of `run.created`, `run.lifecycle`, `run.title`, `run.parent`, `run.archived`, `run.unarchived`, `run.superseded`, `run.notice`, `interview.answered`, `run.branch`, `git.identity`, `checkpoint`, `pull_request.created`, `notification.sent`, `run.paired`.
|
||||
*/
|
||||
export interface RunStreamItem {
|
||||
'run_id': string;
|
||||
/**
|
||||
* The delivery sequence; the cursor.
|
||||
*/
|
||||
'stream_seq': number;
|
||||
'kind': RunStreamItemKind;
|
||||
/**
|
||||
* The item\'s own identity, for deduplication.
|
||||
*/
|
||||
'id': string;
|
||||
/**
|
||||
* Milliseconds since the Unix epoch when the item\'s record was appended.
|
||||
*/
|
||||
'recorded_at': number;
|
||||
/**
|
||||
* The Petri `RunEvent` or the stored platform record, unchanged.
|
||||
*/
|
||||
'item': { [key: string]: any; };
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue