Read a Petri run through the CLI from its stream

`run events` on a Petri run prints the run stream: raw, the envelope as
one JSON line per item; `--pretty`, Petri's events by `<subject>.<verb>`
with the stage's label (a visit's start and end with its elapsed time,
the route both ends of the edge, a fork's branches, a question with its
options and its answer, log lines, the agent's messages and tool calls,
the engine's finish) and the platform records by kind (the run's
creation, its lifecycle, a checkpoint's commit, a pull request, a
notice, who answered). `--follow` attaches from the last `stream_seq`
printed and reconnects from its cursor when the server ends the stream
before the run's terminal record.

`run attach` on a Petri run replays the stream through the progress
renderer (a new mapping from stream items onto the progress events the
renderer draws, sharing the coding-agent mapping with the legacy
envelope), follows it live from its cursor with the same reconnect, asks
a question the stream carries at the terminal, and exits with the status
the engine's finish or the terminal lifecycle record decides. `wait` and
`inspect` read the projection unchanged.

The CLI never names a Petri type: `PetriItem` reads the item as JSON
where Petri's contract keeps the event name, the subject and the parsed
progress payloads.

The CLI's Petri scenarios read the stream instead of the legacy events
(the lifecycle records, the question and who answered it, the expiry),
and three new ones cover a finished run through `events` (raw, tail,
and a `--pretty` snapshot), `attach`, `wait` and `inspect`; `attach`
answering a gate from the terminal; and `events --follow` to the run's
end.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 01:21:12 -04:00
parent 3941a24a28
commit 1bfe62577d
No known key found for this signature in database
9 changed files with 2035 additions and 62 deletions

View file

@ -31,7 +31,7 @@ use fabro_workflow::run_status::RunStatus;
use tokio::signal::ctrl_c;
use tokio::time::{Duration as TokioDuration, sleep};
use super::run_progress;
use super::{petri_stream, run_progress};
use crate::server_client;
const INTERVIEW_UNANSWERED_MESSAGE: &str =
@ -39,6 +39,9 @@ const INTERVIEW_UNANSWERED_MESSAGE: &str =
const JSON_INTERVIEW_MESSAGE: &str = "This run is waiting for human input, but --json is non-interactive. Reattach without --json to answer it.";
const ATTACH_PREMATURE_EOF_MESSAGE: &str = "Attach stream ended before terminal run event.";
const PROMPT_READ_POLL_INTERVAL: TokioDuration = TokioDuration::from_millis(50);
/// How long a Petri attach waits before it reconnects to the stream the
/// server ended while the run was still active.
const STREAM_RECONNECT_DELAY: TokioDuration = TokioDuration::from_millis(200);
enum PromptRead {
Line(String),
@ -182,6 +185,22 @@ pub(crate) async fn attach_run_with_client(
) -> Result<ExitCode> {
let state = client.get_run_state(run_id).await?;
let auto_approve = state.spec.settings.run.execution.approval == ApprovalMode::Auto;
if state.spec.engine.is_petri() {
return Box::pin(attach_petri_run_with_client(
client,
run_id,
&state,
styles,
AttachOptions {
auto_approve,
verbose: live_verbose,
kill_on_detach,
json_output,
},
printer,
))
.await;
}
let events = client.list_run_events(run_id, None, None).await?;
let replay_events = events.clone();
let next_seq = events.last().map_or(1, |event| event.seq.saturating_add(1));
@ -321,6 +340,197 @@ async fn attach_live_run_with_client(
}
}
/// Attach to a Petri run: replay its stream through the progress renderer,
/// then follow it live from the last `stream_seq` seen. A question on the
/// stream is asked at the terminal and answered through the questions API.
/// When the server ends the stream before the run's terminal record, the
/// attach reconnects from its cursor, so no item is missed or repeated.
async fn attach_petri_run_with_client(
client: &server_client::Client,
run_id: &RunId,
state: &server_client::RunProjection,
styles: &'static Styles,
opts: AttachOptions,
printer: Printer,
) -> Result<ExitCode> {
let is_tty = std::io::stderr().is_terminal();
let mut progress_ui = run_progress::ProgressUI::new(is_tty, opts.verbose);
let ctrl_c_signal = ctrl_c();
tokio::pin!(ctrl_c_signal);
let items = client.list_run_stream(run_id, 0).await?;
let mut cursor = items.last().map_or(0, |item| item.stream_seq);
let mut replayed_exit_code = None;
for item in &items {
emit_stream_item(&mut progress_ui, item, opts.json_output)?;
if let Some(code) = petri_stream::exit_code_of(item) {
replayed_exit_code = Some(ExitCode::from(code));
}
}
if let Some(exit_code) = replayed_exit_code.or_else(|| {
state_is_terminal(state).then(|| state_exit_code(state).unwrap_or(ExitCode::from(1)))
}) {
finish_progress(&mut progress_ui, opts.json_output);
return Ok(exit_code);
}
loop {
let mut stream = client.attach_run_stream(run_id, Some(cursor)).await?;
if let Some(exit_code) = Box::pin(handle_pending_petri_interview(
client,
run_id,
&mut stream,
&mut cursor,
&opts,
&mut progress_ui,
styles,
printer,
))
.await?
{
return Ok(exit_code);
}
loop {
let next_item = tokio::select! {
_ = &mut ctrl_c_signal => {
handle_detach_signal(client, run_id, opts.kill_on_detach, printer).await;
finish_progress(&mut progress_ui, opts.json_output);
return Ok(ExitCode::from(1));
}
result = stream.next_item() => result?,
};
let Some(item) = next_item else {
break;
};
cursor = item.stream_seq;
emit_stream_item(&mut progress_ui, &item, opts.json_output)?;
if let Some(code) = petri_stream::exit_code_of(&item) {
finish_progress(&mut progress_ui, opts.json_output);
return Ok(ExitCode::from(code));
}
if petri_stream::question_of(&item).is_some() {
if let Some(exit_code) = Box::pin(handle_pending_petri_interview(
client,
run_id,
&mut stream,
&mut cursor,
&opts,
&mut progress_ui,
styles,
printer,
))
.await?
{
return Ok(exit_code);
}
}
}
// The server ended the stream. A run that concluded has nothing
// more to send past what the grace let through; otherwise this is
// a lost connection, and the attach resumes from its cursor.
let state = client.get_run_state(run_id).await?;
if state_is_terminal(&state) {
for item in client.list_run_stream(run_id, cursor).await? {
emit_stream_item(&mut progress_ui, &item, opts.json_output)?;
}
finish_progress(&mut progress_ui, opts.json_output);
return Ok(state_exit_code(&state).unwrap_or(ExitCode::from(1)));
}
sleep(STREAM_RECONNECT_DELAY).await;
}
}
/// Ask the run's pending question, if one is listed, while the stream keeps
/// flowing: an answer given elsewhere, or the run ending, ends the prompt.
async fn handle_pending_petri_interview(
client: &server_client::Client,
run_id: &RunId,
stream: &mut server_client::RunStreamItemStream,
cursor: &mut u64,
opts: &AttachOptions,
progress_ui: &mut run_progress::ProgressUI,
styles: &'static Styles,
printer: Printer,
) -> Result<Option<ExitCode>> {
let Some(question) = client.list_run_questions(run_id).await?.into_iter().next() else {
return Ok(None);
};
if json_pending_interview_requires_manual_input(opts.json_output, opts.auto_approve) {
fabro_util::printerr!(printer, "{JSON_INTERVIEW_MESSAGE}");
return Ok(Some(ExitCode::from(1)));
}
if opts.json_output {
return Ok(None);
}
hide_progress(progress_ui, opts.json_output);
let ask = ask_attach_question(api_question_to_question(&question), styles);
tokio::pin!(ask);
let ctrl_c_signal = ctrl_c();
tokio::pin!(ctrl_c_signal);
let answer = loop {
let next_item = tokio::select! {
answer = &mut ask => {
break answer;
}
_ = &mut ctrl_c_signal => {
handle_detach_signal(client, run_id, opts.kill_on_detach, printer).await;
show_progress(progress_ui, opts.json_output);
return Ok(Some(ExitCode::from(1)));
}
result = stream.next_item() => result?,
};
// The stream ended under the prompt: the caller reconnects and asks
// again if the question is still pending.
let Some(item) = next_item else {
show_progress(progress_ui, opts.json_output);
return Ok(None);
};
*cursor = item.stream_seq;
emit_stream_item(progress_ui, &item, opts.json_output)?;
if let Some(code) = petri_stream::exit_code_of(&item) {
show_progress(progress_ui, opts.json_output);
return Ok(Some(ExitCode::from(code)));
}
if petri_stream::resolves_question(&item, &question.id) {
show_progress(progress_ui, opts.json_output);
return Ok(None);
}
};
show_progress(progress_ui, opts.json_output);
if answer_requires_reattach(&answer) {
fabro_util::printerr!(printer, "{INTERVIEW_UNANSWERED_MESSAGE}");
return Ok(Some(ExitCode::from(1)));
}
submit_server_interview_answer(client, run_id, &question.id, &answer).await?;
Ok(None)
}
fn emit_stream_item(
progress_ui: &mut run_progress::ProgressUI,
item: &fabro_types::RunStreamItem,
json_output: bool,
) -> Result<()> {
if json_output {
let stdout = std::io::stdout();
let mut handle = stdout.lock();
writeln!(handle, "{}", petri_stream::raw_line(item)?)?;
} else {
progress_ui.handle_stream_item(item);
}
Ok(())
}
async fn handle_pending_server_interview(
client: &server_client::Client,
run_id: &RunId,

View file

@ -20,6 +20,7 @@ use fabro_util::terminal::Styles;
use tokio::time;
use tracing::{debug, info};
use super::petri_stream;
use crate::args::EventsArgs;
use crate::command_context::CommandContext;
use crate::server_client;
@ -42,6 +43,24 @@ pub(crate) async fn run(
None => None,
};
// A Petri run's events are its stream, in the stream envelope.
let state = client
.get_run_state(&run_id)
.await
.context("Failed to read run state from server")?;
if state.spec.engine.is_petri() {
let pretty = args.pretty && !ctx.json_output();
return Box::pin(petri_stream::print_events(
client.as_ref(),
&run_id,
args,
since_cutoff,
pretty,
styles,
))
.await;
}
let events = match (args.tail, since_cutoff.is_none()) {
(Some(tail), true) => {
// With --tail 0 --follow, fetch one event anyway so `last_seq`

View file

@ -20,6 +20,7 @@ pub(crate) mod fork;
pub(crate) mod logs;
pub(crate) mod output;
pub(crate) mod overrides;
pub(crate) mod petri_stream;
mod petri_worker;
pub(crate) mod preview;
mod remote_workflow;
@ -126,7 +127,7 @@ pub(crate) async fn dispatch(
RunCommands::Diff(args) => diff::run(args, base_ctx).await,
RunCommands::Events(args) => {
let styles = Styles::detect_stdout();
events::run(&args, &styles, base_ctx).await
Box::pin(events::run(&args, &styles, base_ctx)).await
}
RunCommands::Logs(args) => logs::run(&args, base_ctx).await,
RunCommands::Resume(args) => {

File diff suppressed because it is too large Load diff

View file

@ -388,7 +388,23 @@ fn agent_progress_event(
stored: &RunEvent,
event: &CodingEvent,
) -> Option<ProgressEvent> {
let root_session = stored.parent_session_id.is_none();
coding_progress_event(
node_id,
stored.parent_session_id.is_none(),
Some(stored.ts),
event,
)
}
/// The progress line for one coding agent event, given whether it came
/// from the root session and when it was recorded: the mapping the legacy
/// envelope and a Petri stream envelope share.
pub(super) fn coding_progress_event(
node_id: String,
root_session: bool,
timestamp: Option<DateTime<Utc>>,
event: &CodingEvent,
) -> Option<ProgressEvent> {
match event {
CodingEvent::AssistantMessage { model, .. } => Some(ProgressEvent::AssistantMessage {
stage_node_id: node_id,
@ -401,10 +417,10 @@ fn agent_progress_event(
arguments,
} => Some(ProgressEvent::ToolCallStarted {
stage_node_id: node_id,
tool_name: tool_name.clone(),
tool_call_id: tool_call_id.clone(),
arguments: arguments.clone(),
timestamp: Some(stored.ts),
tool_name: tool_name.clone(),
tool_call_id: tool_call_id.clone(),
arguments: arguments.clone(),
timestamp,
}),
CodingEvent::ToolCallCompleted {
tool_call_id,
@ -412,10 +428,10 @@ fn agent_progress_event(
..
} => Some(ProgressEvent::ToolCallCompleted {
stage_node_id: node_id,
tool_call_id: tool_call_id.clone(),
is_error: *is_error,
duration_ms: None,
timestamp: Some(stored.ts),
tool_call_id: tool_call_id.clone(),
is_error: *is_error,
duration_ms: None,
timestamp,
}),
CodingEvent::Warning { kind, details, .. } if kind == "context_window" => {
let usage_percent = details

View file

@ -3,10 +3,11 @@
reason = "sync CLI run-progress renderer: writes to std::io::stderr directly"
)]
use fabro_types::{RunEvent, RunNoticeCode};
use fabro_types::{RunEvent, RunNoticeCode, RunStreamItem};
mod event;
mod info_display;
mod petri;
mod renderer;
mod setup_display;
mod stage_display;
@ -14,6 +15,7 @@ mod styles;
use event::{ProgressEvent, from_json_line, from_run_event};
use info_display::InfoDisplay;
use petri::PetriProgressState;
use renderer::ProgressRenderer;
use setup_display::SetupDisplay;
use stage_display::StageDisplay;
@ -24,6 +26,7 @@ pub(crate) struct ProgressUI {
setup: SetupDisplay,
info: InfoDisplay,
saw_metadata_snapshot_failure: bool,
petri: PetriProgressState,
}
impl ProgressUI {
@ -46,6 +49,7 @@ impl ProgressUI {
setup: SetupDisplay::new(verbose),
info: InfoDisplay::new(verbose),
saw_metadata_snapshot_failure: false,
petri: PetriProgressState::default(),
}
}
@ -95,6 +99,14 @@ impl ProgressUI {
}
}
/// One item of a Petri run's stream: the progress lines it means, if
/// any, rendered as a legacy event's would be.
pub(crate) fn handle_stream_item(&mut self, item: &RunStreamItem) {
for progress_event in petri::progress_events(item, &mut self.petri) {
self.dispatch(progress_event);
}
}
fn dispatch(&mut self, event: ProgressEvent) {
let renderer = &self.renderer;
match event {

View file

@ -0,0 +1,264 @@
//! The progress lines a Petri run's stream items mean: the mapping from
//! Petri's `<subject>.<verb>` events and Fabro's platform records onto the
//! [`ProgressEvent`]s the renderer already draws for a legacy run.
//!
//! A stage is a firing whose node is a logical stage (`VIEWS.md`); its
//! display key is the node's name, as a legacy stage's `node_id` is. A
//! fork's `parallel.branch` delegates are the branches of the parallel
//! group, never stages of their own.
use fabro_types::run_event::RunNoticeLevel;
use fabro_types::{CodingAgentEvent, RunStreamItem, StageOutcome, StageTiming};
use serde_json::Value;
use super::event::{ProgressEvent, ProgressUsage, coding_progress_event};
use crate::commands::run::petri_stream::{PetriItem, StageClock};
/// What the mapping remembers between items: when each firing started.
#[derive(Default)]
pub(super) struct PetriProgressState {
clock: StageClock,
}
/// The progress events one stream item means, in order.
pub(super) fn progress_events(
item: &RunStreamItem,
state: &mut PetriProgressState,
) -> Vec<ProgressEvent> {
let elapsed = state.clock.observe(item);
let view = PetriItem::new(item);
if let Some(record) = view.platform_record() {
return platform_progress_event(record).into_iter().collect();
}
let Some(name) = view.name() else {
return Vec::new();
};
let node_id = view.node_name().unwrap_or("?").to_string();
let label = view.node_label().unwrap_or("?").to_string();
match name {
"visit.started" => {
if view.is_shown_stage() {
vec![ProgressEvent::StageStarted {
node_id,
name: label,
script: None,
}]
} else if view.node_kind() == Some("parallel.branch") {
vec![ProgressEvent::ParallelBranchStarted { branch: node_id }]
} else {
Vec::new()
}
}
"visit.completed" if view.is_shown_stage() => {
let Some(derived) = view.derived() else {
return Vec::new();
};
let outcome = derived.get("outcome");
let executed = derived
.get("executed")
.and_then(Value::as_bool)
.unwrap_or(true);
let status = outcome
.and_then(|outcome| outcome.get("status"))
.and_then(Value::as_str)
.unwrap_or("?");
let timing = StageTiming {
wall_time_ms: elapsed.unwrap_or(0),
..StageTiming::default()
};
let completed =
|status: &str, usage: Option<ProgressUsage>| ProgressEvent::StageCompleted {
node_id: node_id.clone(),
name: label.clone(),
timing,
status: status.to_string(),
usage,
};
if !executed || status == "skipped" {
return vec![completed("skipped", None)];
}
match status {
"success" => vec![completed("succeeded", outcome.and_then(usage_of))],
"partial_success" => {
vec![completed("partially_succeeded", outcome.and_then(usage_of))]
}
"cancelled" => vec![completed("cancelled", None)],
other => {
let error = outcome
.and_then(|outcome| outcome.pointer("/failure/message"))
.and_then(Value::as_str)
.unwrap_or(other)
.to_string();
vec![ProgressEvent::StageFailed {
node_id,
name: label,
error,
}]
}
}
}
"retry.scheduled" => {
let Some(derived) = view.derived() else {
return Vec::new();
};
let attempt = derived
.get("next_attempt")
.and_then(Value::as_u64)
.unwrap_or(0);
let delay_ms = derived
.pointer("/base_delay/secs")
.and_then(Value::as_u64)
.map(|secs| secs.saturating_mul(1000))
.or_else(|| derived.get("base_delay").and_then(Value::as_u64))
.unwrap_or(0);
vec![ProgressEvent::StageRetrying {
name: label,
attempt,
max_attempts: attempt,
delay_ms,
}]
}
"fork.started" => vec![ProgressEvent::ParallelStarted],
"branch.completed" => {
let Some(result) = view.derived().and_then(|derived| derived.get("result")) else {
return Vec::new();
};
let branch = result
.pointer("/node/name")
.and_then(Value::as_str)
.unwrap_or(&node_id)
.to_string();
let status = match result.get("status").and_then(Value::as_str) {
Some("success") => StageOutcome::Succeeded,
Some("partial_success") => StageOutcome::PartiallySucceeded,
_ => StageOutcome::Failed {
retry_requested: false,
},
};
vec![ProgressEvent::ParallelBranchCompleted {
branch,
duration_ms: elapsed.unwrap_or(0),
status,
}]
}
"fork.completed" => vec![ProgressEvent::ParallelCompleted],
"route.applied" => {
let Some(derived) = view.derived() else {
return Vec::new();
};
let Some(to_node) = derived.pointer("/target/name").and_then(Value::as_str) else {
return Vec::new();
};
let back = derived
.get("back")
.and_then(Value::as_bool)
.unwrap_or(false);
if back {
vec![ProgressEvent::LoopRestart {
from_node: node_id,
to_node: to_node.to_string(),
}]
} else {
vec![ProgressEvent::EdgeSelected {
from_node: node_id,
to_node: to_node.to_string(),
label: None,
condition: derived
.get("transition")
.and_then(Value::as_str)
.filter(|transition| *transition != "Continue")
.map(str::to_lowercase),
}]
}
}
"step.progress.recorded" => envelope_progress_event(view, node_id).into_iter().collect(),
_ => Vec::new(),
}
}
/// The progress line of a Pebble coding-agent envelope on a stage's
/// stream, if the terminal shows it.
fn envelope_progress_event(view: PetriItem<'_>, node_id: String) -> Option<ProgressEvent> {
let custom = view.custom()?;
if custom.get("kind").and_then(Value::as_str) != Some("pebble") {
return None;
}
let envelope: CodingAgentEvent = serde_json::from_value(custom.get("event")?.clone()).ok()?;
let root_session = envelope.parent_session_id.is_none();
coding_progress_event(
node_id,
root_session,
Some(view.recorded_at()),
&envelope.event,
)
}
/// The stage's model usage from a finished visit's metrics, when the step
/// reported one.
fn usage_of(outcome: &Value) -> Option<ProgressUsage> {
let custom = outcome.pointer("/metrics/custom")?;
let usage = custom
.get("pebble.usage")
.or_else(|| custom.get("prompt.usage"))?;
let tokens = usage.get("tokens")?;
Some(ProgressUsage {
input_tokens: tokens.get("input").and_then(Value::as_u64).unwrap_or(0),
output_tokens: tokens.get("output").and_then(Value::as_u64).unwrap_or(0),
cost: usage
.pointer("/cost/usd_micros")
.and_then(Value::as_u64)
.map(|micros| micros as f64 / 1_000_000.0),
})
}
/// The progress line a platform record means, if the terminal shows it.
fn platform_progress_event(record: &Value) -> Option<ProgressEvent> {
match record.get("kind")?.as_str()? {
"run.created" => Some(ProgressEvent::RunCreated {
web_url: record
.get("web_url")
.and_then(Value::as_str)
.map(str::to_string),
}),
"run.branch" => Some(ProgressEvent::WorkflowStarted {
worktree_dir: None,
base_branch: record
.get("run_branch")
.and_then(Value::as_str)
.map(str::to_string),
base_sha: record
.get("base_sha")
.and_then(Value::as_str)
.map(str::to_string),
}),
"run.notice" => Some(ProgressEvent::RunNotice {
level: match record.get("level").and_then(Value::as_str) {
Some("warn") => RunNoticeLevel::Warn,
Some("error") => RunNoticeLevel::Error,
_ => RunNoticeLevel::Info,
},
code: record
.get("code")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
message: record
.get("message")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
}),
"pull_request.created" => Some(ProgressEvent::PullRequestCreated {
pr_url: record
.get("html_url")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
draft: record
.get("draft")
.and_then(Value::as_bool)
.unwrap_or(false),
}),
_ => None,
}
}

View file

@ -7,7 +7,7 @@ use fabro_client::{
AuthEntry, AuthStore, Credential, OAuthSession, ServerTarget, TransportConnector,
apply_bearer_token_auth,
};
pub(crate) use fabro_client::{Client, RunEventStream};
pub(crate) use fabro_client::{Client, RunEventStream, RunStreamItemStream};
use fabro_config::Storage;
use fabro_config::bind::Bind;
pub(crate) use fabro_types::RunProjection;

View file

@ -22,8 +22,9 @@
#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")]
use std::env;
use std::io::{Read as _, Write as _};
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::process::{Child, Command, Output, Stdio};
use std::time::{Duration, Instant};
use fabro_client::ServerTarget;
@ -32,14 +33,12 @@ use fabro_petri::SqliteRunStore;
use fabro_petri::engine::{self, RunStatus};
use fabro_petri::petri::RunKey;
use fabro_static::EnvVars;
use fabro_store::EventEnvelope;
use fabro_test::{apply_test_isolation, expect_reqwest_json, isolated_storage_dir, test_context};
use fabro_types::EventBody;
use fabro_test::{
apply_test_isolation, expect_reqwest_json, fabro_snapshot, isolated_storage_dir, test_context,
};
use crate::cmd::support::created_run_id;
use crate::support::{
TEST_DEV_TOKEN, TEST_SESSION_SECRET, parse_event_envelopes, seed_dev_token_auth,
};
use crate::support::{TEST_DEV_TOKEN, TEST_SESSION_SECRET, seed_dev_token_auth};
const HOST_PLUGIN: &str = "sandbox-driver-host";
const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS";
@ -353,17 +352,70 @@ async fn wait_for_status(server: &RunningServer, run_id: &str, expected: &[&str]
}
}
async fn run_events(server: &RunningServer, run_id: &str) -> Vec<EventEnvelope> {
parse_event_envelopes(&run_json(server, &format!("runs/{run_id}/events")).await)
/// The run's stream, as `GET /runs/{id}/events` serves a Petri run: every
/// item in `stream_seq` order, in the stream envelope.
async fn run_stream(server: &RunningServer, run_id: &str) -> Vec<serde_json::Value> {
let mut items = Vec::new();
let mut after = 0;
loop {
let page = run_json(
server,
&format!("runs/{run_id}/events?after={after}&limit=1000"),
)
.await;
let data = page["data"]
.as_array()
.cloned()
.expect("the stream page has a data array");
let Some(last) = data.last() else {
break;
};
after = last["stream_seq"].as_u64().expect("a stream_seq");
let has_more = page["meta"]["has_more"].as_bool().unwrap_or(false);
items.extend(data);
if !has_more {
break;
}
}
items
}
fn event_names(events: &[EventEnvelope]) -> Vec<&str> {
events
/// What each stream item is, for an assertion: a Petri event by its
/// `<subject>.<verb>` name (`question` and `question_expired` for the parsed
/// progress payloads), a platform lifecycle record as
/// `lifecycle:<transition>`, another platform record by its kind.
fn stream_names(items: &[serde_json::Value]) -> Vec<String> {
items
.iter()
.map(|envelope| envelope.event.event_name())
.map(|line| {
let item = &line["item"];
if line["kind"] == "platform" {
let record = &item["record"];
return match record["kind"].as_str().unwrap_or("?") {
"run.lifecycle" => {
format!("lifecycle:{}", record["transition"].as_str().unwrap_or("?"))
}
kind => kind.to_string(),
};
}
if let Some(kind) = item["derived"]["parsed"]["kind"].as_str() {
if matches!(kind, "question" | "question_expired") {
return kind.to_string();
}
}
item["record"]["body"]["event"]
.as_str()
.or_else(|| item["derived"]["event"].as_str())
.unwrap_or("?")
.to_string()
})
.collect()
}
fn count_of(names: &[String], expected: &str) -> usize {
names.iter().filter(|name| *name == expected).count()
}
/// The pid of the worker subprocess the server launched for the run: the
/// worker retitles itself `fabro <first 12 of the run id> <phase>`, so that
/// is what the process table shows.
@ -432,18 +484,12 @@ async fn a_petri_run_executes_in_the_server_launched_worker() {
let state = run_json(&server, &format!("runs/{run_id}/state")).await;
assert_eq!(state["spec"]["engine"]["kind"], "petri", "state: {state}");
let events = run_events(&server, &run_id).await;
let names = event_names(&events);
assert_eq!(
names
.iter()
.filter(|name| **name == "run.completed")
.count(),
1,
"{names:?}"
);
let names = stream_names(&run_stream(&server, &run_id).await);
assert_eq!(count_of(&names, "lifecycle:succeeded"), 1, "{names:?}");
assert_eq!(count_of(&names, "run.finished"), 1, "{names:?}");
assert!(
names.contains(&"run.starting") && names.contains(&"run.running"),
names.iter().any(|name| name == "lifecycle:starting")
&& names.iter().any(|name| name == "lifecycle:running"),
"{names:?}"
);
@ -519,40 +565,35 @@ async fn a_petri_run_resumes_in_a_new_worker_after_the_server_restarts() {
std::fs::write(&gate, "go").expect("the gate opens");
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
let events = run_events(&server, &run_id).await;
let names = event_names(&events);
let items = run_stream(&server, &run_id).await;
let names = stream_names(&items);
assert_eq!(
status,
"succeeded",
"events: {names:?}\nserver stderr:\n{}",
"stream: {names:?}\nserver stderr:\n{}",
server.stderr_text()
);
assert_eq!(
names
.iter()
.filter(|name| **name == "run.completed")
.count(),
1,
"{names:?}"
);
assert_eq!(count_of(&names, "lifecycle:succeeded"), 1, "{names:?}");
assert_eq!(count_of(&names, "run.finished"), 1, "{names:?}");
// `fabro run` asked for the first start; the restart asked for a
// resume, after the run had been running.
let first_running = names
.iter()
.position(|name| *name == "run.running")
.position(|name| name == "lifecycle:running")
.expect("the run ran before the crash");
let resume_request = events
let resume_request = items
.iter()
.position(|envelope| {
matches!(
&envelope.event.body,
EventBody::RunStartRequested(props) if props.resume
)
.position(|line| {
let record = &line["item"]["record"];
line["kind"] == "platform"
&& record["kind"] == "run.lifecycle"
&& record["transition"] == "start_requested"
&& record["source"] == "resume"
})
.expect("the restart asked for a resume");
assert!(resume_request > first_running, "{names:?}");
assert_eq!(
names.iter().filter(|name| **name == "run.running").count(),
count_of(&names, "lifecycle:running"),
2,
"the run ran once before and once after the restart: {names:?}"
);
@ -706,11 +747,11 @@ async fn a_human_gate_in_the_worker_is_answered_through_the_api() {
markers.join("no").exists() && !markers.join("yes").exists(),
"the no branch ran"
);
let names = run_events(&server, &run_id).await;
let names = event_names(&names);
let names = stream_names(&run_stream(&server, &run_id).await);
assert!(
names.contains(&"interview.started") && names.contains(&"interview.completed"),
"{names:?}"
names.iter().any(|name| name == "question")
&& names.iter().any(|name| name == "interview.answered"),
"the question and who answered it are on the stream: {names:?}"
);
assert!(questions(&server, &run_id).await.is_empty());
let store = server.petri_store().await;
@ -805,9 +846,303 @@ async fn an_unanswered_gate_in_the_worker_expires_with_its_default() {
markers.join("no").exists() && !markers.join("yes").exists(),
"the default ran"
);
let events = run_events(&server, &run_id).await;
let names = event_names(&events);
assert!(names.contains(&"interview.timeout"), "{names:?}");
let names = stream_names(&run_stream(&server, &run_id).await);
assert!(
names.iter().any(|name| name == "question_expired"),
"{names:?}"
);
assert!(questions(&server, &run_id).await.is_empty());
server.shutdown();
}
/// A CLI command against the server, as `run_detached` seeds its auth.
fn cli(context: &fabro_test::TestContext, server: &RunningServer, args: &[&str]) -> Output {
let target = server.target();
let output = context
.command()
.args(args)
.args(["--server", &target])
.output()
.expect("the CLI command executes");
assert!(
output.status.success(),
"`fabro {}` failed\nstdout:\n{}\nstderr:\n{}",
args.join(" "),
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
output
}
fn ndjson(output: &Output) -> Vec<serde_json::Value> {
String::from_utf8_lossy(&output.stdout)
.lines()
.filter(|line| !line.trim().is_empty())
.map(|line| serde_json::from_str(line).expect("a JSON line"))
.collect()
}
/// The `<subject>.<verb>` name of a Petri item, or the kind of a platform
/// record, from a raw stream line.
fn stream_line_name(line: &serde_json::Value) -> String {
let item = &line["item"];
if line["kind"] == "platform" {
return format!(
"platform:{}",
item["record"]["kind"].as_str().unwrap_or("?")
);
}
item["record"]["body"]["event"]
.as_str()
.or_else(|| item["derived"]["event"].as_str())
.unwrap_or("?")
.to_string()
}
/// The snapshot filters for `events --pretty` over a Petri run: clocks,
/// durations and the run id vary per run.
fn pretty_filters(context: &fabro_test::TestContext) -> Vec<(String, String)> {
let mut filters = context.filters();
filters.push((r"\b\d{2}:\d{2}:\d{2}\b".to_string(), "[CLOCK]".to_string()));
filters.push((
r"\b\d+(\.\d+)?(ms|s)\b".to_string(),
"[DURATION]".to_string(),
));
filters
}
/// A finished Petri run reads back through the CLI: `events` prints the
/// stream envelope raw, dense in `stream_seq`; `events --pretty` renders
/// the stages by `<subject>.<verb>` with their labels and the platform
/// records by kind; `attach` replays it and exits with the run's status;
/// `wait` and `runs inspect` read the projection.
#[tokio::test(flavor = "multi_thread")]
async fn a_finished_petri_run_reads_back_through_the_cli() {
if host_plugin().is_none() {
return;
}
let context = test_context!();
let server = RunningServer::start().await;
let workspace = write_petri_workspace(&context, "echo hello from petri");
let run_id = run_detached(&context, &server, &workspace);
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
assert_eq!(
status,
"succeeded",
"server stderr:\n{}",
server.stderr_text()
);
let target = server.target();
// Raw: the envelope, one item per line, dense and in order.
let raw = cli(&context, &server, &["events", &run_id]);
let lines = ndjson(&raw);
let seqs: Vec<u64> = lines
.iter()
.map(|line| line["stream_seq"].as_u64().expect("a stream_seq"))
.collect();
let expected: Vec<u64> = (1..=seqs.len() as u64).collect();
assert_eq!(seqs, expected, "stream_seq is dense");
for line in &lines {
assert_eq!(line["run_id"], run_id, "{line}");
assert!(
line["id"].is_string() && line["recorded_at"].is_u64(),
"{line}"
);
}
let names: Vec<String> = lines.iter().map(stream_line_name).collect();
for expected in [
"platform:run.created",
"platform:run.lifecycle",
"run.started",
"visit.started",
"visit.completed",
"run.finished",
] {
assert!(
names.iter().any(|name| name == expected),
"{expected} is on the stream: {names:?}"
);
}
assert_eq!(
names.last().map(String::as_str),
Some("platform:run.lifecycle"),
"the terminal lifecycle record ends the stream: {names:?}"
);
// Tail: the last two items only.
let tail = cli(&context, &server, &["events", "--tail", "2", &run_id]);
assert_eq!(ndjson(&tail).len(), 2);
// Pretty: stages and platform records.
let mut cmd = context.command();
cmd.args(["events", "--pretty", "--server", &target, &run_id]);
fabro_snapshot!(pretty_filters(&context), cmd, @r"
success: true
exit_code: 0
----- stdout -----
[CLOCK] ▶ Run one command [ULID]
[CLOCK] · submitted
[CLOCK] · start_requested
[CLOCK] · runnable
[CLOCK] · starting
[CLOCK] · running
[CLOCK] Engine: petri run started
[CLOCK] ▶ start
[CLOCK] │ checkout: [TEMP_DIR]/petri-workspace is not a Git repository; the workspace starts empty
[CLOCK] ✓ start [DURATION]
[CLOCK] ▶ say
[CLOCK] start → say continue
[CLOCK] │ hello from petri
[CLOCK] ✓ say [DURATION]
[CLOCK] ▶ exit
[CLOCK] say → exit continue
[CLOCK] ✓ exit [DURATION]
[CLOCK] ✓ SUCCEEDED [DURATION]
[CLOCK] · succeeded
----- stderr -----
");
// Attach replays the finished run and exits with its status.
let attach = cli(&context, &server, &["attach", &run_id]);
let stderr = String::from_utf8_lossy(&attach.stderr);
assert!(stderr.contains("say"), "the stage is drawn: {stderr}");
// Wait reads the projection's status and conclusion.
let wait = cli(&context, &server, &["wait", &run_id]);
let stderr = String::from_utf8_lossy(&wait.stderr);
assert!(stderr.contains("Succeeded"), "{stderr}");
// Inspect reads the projection, whose spec names the engine.
let inspect = cli(&context, &server, &["inspect", &run_id]);
let inspected: serde_json::Value =
serde_json::from_slice(&inspect.stdout).expect("inspect prints JSON");
let entry = &inspected[0];
assert_eq!(entry["run_id"], run_id, "{entry}");
assert_eq!(entry["run_spec"]["engine"]["kind"], "petri", "{entry}");
assert_eq!(entry["conclusion"]["status"], "succeeded", "{entry}");
server.shutdown();
}
/// `attach` on a Petri run with a human gate asks the question at the
/// terminal and answers it through the questions API; the answer routes
/// the gate and the attach exits with the run's status.
#[tokio::test(flavor = "multi_thread")]
async fn attach_asks_a_petri_gate_at_the_terminal_and_answers_it() {
if host_plugin().is_none() {
return;
}
let context = test_context!();
let server = RunningServer::start().await;
let markers = context.temp_dir.join("markers");
std::fs::create_dir_all(&markers).expect("the marker dir creates");
let workspace = write_petri_workflow(&context, &gate_dot(&markers, ""));
let run_id = run_detached_with(&context, &server, &workspace, &[]);
wait_for_questions(&server, &run_id, 1).await;
let target = server.target();
let mut attach_cmd = Command::new(env!("CARGO_BIN_EXE_fabro"));
apply_test_isolation(&mut attach_cmd, &context.home_dir);
attach_cmd
.current_dir(&context.temp_dir)
.args(["attach", "--server", &target, &run_id])
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
let mut child = attach_cmd.spawn().expect("attach spawns");
{
let mut stdin = child.stdin.take().expect("attach stdin is piped");
stdin.write_all(b"N\n").expect("the answer writes");
}
let output = child
.wait_with_output()
.expect("attach exits once the run ends");
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
output.status.success(),
"attach failed\nstderr:\n{stderr}\nserver stderr:\n{}",
server.stderr_text()
);
assert!(stderr.contains("Go?"), "the question was asked: {stderr}");
assert!(
markers.join("no").exists() && !markers.join("yes").exists(),
"the no branch ran"
);
let status = wait_for_status(&server, &run_id, &["succeeded", "failed"]).await;
assert_eq!(status, "succeeded");
server.shutdown();
}
/// `events --follow` on a Petri run follows the stream live from its
/// cursor: the items already stored print first, the ones committed while
/// the run goes on follow, and the terminal lifecycle record ends it.
#[tokio::test(flavor = "multi_thread")]
async fn events_follow_streams_a_petri_run_live_to_its_end() {
if host_plugin().is_none() {
return;
}
let context = test_context!();
let server = RunningServer::start().await;
let gate = context.temp_dir.join("go");
let workspace = write_petri_workspace(
&context,
&format!(
"while [ ! -f {} ]; do sleep 0.05; done; echo released",
gate.display()
),
);
let run_id = run_detached(&context, &server, &workspace);
wait_for_status(&server, &run_id, &["running"]).await;
let target = server.target();
let mut follow_cmd = Command::new(env!("CARGO_BIN_EXE_fabro"));
apply_test_isolation(&mut follow_cmd, &context.home_dir);
follow_cmd
.current_dir(&context.temp_dir)
.args([
"events", "--follow", "--pretty", "--server", &target, &run_id,
])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
let mut child = follow_cmd.spawn().expect("events --follow spawns");
// Let the follower attach before the run is released.
tokio::time::sleep(Duration::from_millis(500)).await;
std::fs::write(&gate, b"").expect("the release marker writes");
let deadline = Instant::now() + RUN_TIMEOUT;
let status = loop {
if let Some(status) = child.try_wait().expect("the follower polls") {
break status;
}
assert!(
Instant::now() < deadline,
"events --follow did not end with the run; server stderr:\n{}",
server.stderr_text()
);
tokio::time::sleep(POLL).await;
};
let mut stdout = String::new();
child
.stdout
.take()
.expect("stdout is piped")
.read_to_string(&mut stdout)
.expect("stdout reads");
let mut stderr = String::new();
child
.stderr
.take()
.expect("stderr is piped")
.read_to_string(&mut stderr)
.expect("stderr reads");
assert!(status.success(), "events --follow failed: {stderr}");
assert!(stdout.contains("▶ say"), "the stage started: {stdout}");
assert!(stdout.contains("│ released"), "the live log line: {stdout}");
assert!(stdout.contains("✓ SUCCEEDED"), "the finish: {stdout}");
assert!(
stdout.trim_end().ends_with("· succeeded"),
"the terminal lifecycle record ends the follow: {stdout}"
);
server.shutdown();
}