From 1bfe62577df44f13cf6a9cd9f5ddd7e3f1cdbe7c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 01:21:12 -0400 Subject: [PATCH] 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 `.` 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 --- lib/apps/fabro-cli/src/commands/run/attach.rs | 212 +++- lib/apps/fabro-cli/src/commands/run/events.rs | 19 + lib/apps/fabro-cli/src/commands/run/mod.rs | 3 +- .../src/commands/run/petri_stream.rs | 1116 +++++++++++++++++ .../src/commands/run/run_progress/event.rs | 34 +- .../src/commands/run/run_progress/mod.rs | 14 +- .../src/commands/run/run_progress/petri.rs | 264 ++++ lib/apps/fabro-cli/src/server_client.rs | 2 +- lib/apps/fabro-cli/tests/it/scenario/petri.rs | 433 ++++++- 9 files changed, 2035 insertions(+), 62 deletions(-) create mode 100644 lib/apps/fabro-cli/src/commands/run/petri_stream.rs create mode 100644 lib/apps/fabro-cli/src/commands/run/run_progress/petri.rs diff --git a/lib/apps/fabro-cli/src/commands/run/attach.rs b/lib/apps/fabro-cli/src/commands/run/attach.rs index e14f3949d..85216054c 100644 --- a/lib/apps/fabro-cli/src/commands/run/attach.rs +++ b/lib/apps/fabro-cli/src/commands/run/attach.rs @@ -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 { 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 { + 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> { + 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, diff --git a/lib/apps/fabro-cli/src/commands/run/events.rs b/lib/apps/fabro-cli/src/commands/run/events.rs index 062c7fc10..ab35b3c05 100644 --- a/lib/apps/fabro-cli/src/commands/run/events.rs +++ b/lib/apps/fabro-cli/src/commands/run/events.rs @@ -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` diff --git a/lib/apps/fabro-cli/src/commands/run/mod.rs b/lib/apps/fabro-cli/src/commands/run/mod.rs index f59559070..202e75e59 100644 --- a/lib/apps/fabro-cli/src/commands/run/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/mod.rs @@ -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) => { diff --git a/lib/apps/fabro-cli/src/commands/run/petri_stream.rs b/lib/apps/fabro-cli/src/commands/run/petri_stream.rs new file mode 100644 index 000000000..5539ad7cc --- /dev/null +++ b/lib/apps/fabro-cli/src/commands/run/petri_stream.rs @@ -0,0 +1,1116 @@ +//! A Petri run's stream as the CLI reads it: `run events` prints the +//! envelope raw or pretty, and `run attach` follows it live from a cursor. +//! +//! The stream is `GET /runs/{id}/events` for a run whose spec names Petri: +//! one ordered delivery of Petri's own events and Fabro's platform records +//! in the `RunStreamItem` envelope, addressed by `stream_seq`. The CLI +//! never names a Petri type; it reads the item as JSON through +//! [`PetriItem`], which knows where Petri's contract keeps the event name, +//! the subject and the parsed progress payloads. + +#![expect( + clippy::disallowed_types, + reason = "sync CLI `run events` output: blocking std::io::Write is the intended mechanism" +)] +#![expect( + clippy::disallowed_methods, + reason = "sync CLI `run events` output: streams lines to std::io::stdout directly" +)] + +use std::collections::HashMap; +use std::fmt::Write as _; +use std::io::{self, Write}; +use std::time::Duration; + +use anyhow::{Context as _, Result}; +use chrono::{DateTime, TimeZone as _, Utc}; +use fabro_redact::redact_jsonl_line; +use fabro_types::{RunId, RunStreamItem, RunStreamItemKind}; +use fabro_util::terminal::Styles; +use serde_json::Value; +use tokio::time; +use tracing::debug; + +use crate::args::EventsArgs; +use crate::server_client; +use crate::shared::{format_duration_ms, format_usd_micros}; + +/// How long a follower waits before it reconnects after the server ended +/// the stream while the run was still active. +const FOLLOW_RECONNECT_DELAY: Duration = Duration::from_millis(200); + +/// One stream item read as Petri's contract lays it out. +#[derive(Clone, Copy)] +pub(crate) struct PetriItem<'a> { + pub(crate) item: &'a RunStreamItem, +} + +impl<'a> PetriItem<'a> { + pub(crate) fn new(item: &'a RunStreamItem) -> Self { + Self { item } + } + + pub(crate) fn is_petri(self) -> bool { + self.item.kind == RunStreamItemKind::Petri + } + + /// The `.` name of a Petri event, or the kind of a + /// platform record. + pub(crate) fn name(self) -> Option<&'a str> { + self.item.name() + } + + pub(crate) fn recorded_at(self) -> DateTime { + millis(self.item.recorded_at) + } + + fn value(self) -> &'a Value { + &self.item.item + } + + /// The recorded event's fields, beside its `event` tag. + pub(crate) fn body(self) -> Option<&'a Value> { + self.value().pointer("/record/body") + } + + pub(crate) fn derived(self) -> Option<&'a Value> { + self.value().get("derived") + } + + /// Petri's reading of a `step.progress.recorded` payload. + pub(crate) fn parsed(self) -> Option<&'a Value> { + self.derived()?.get("parsed") + } + + /// The `custom` payload of a `step.progress.recorded`, when it is one. + pub(crate) fn custom(self) -> Option<&'a Value> { + self.body()?.pointer("/ev/custom") + } + + /// A `step.progress.recorded` log line: `(stream, line)`. + pub(crate) fn log_line(self) -> Option<(&'a str, &'a str)> { + let log = self.body()?.pointer("/ev/log")?; + Some((log.get("stream")?.as_str()?, log.get("line")?.as_str()?)) + } + + pub(crate) fn subject(self) -> Option<&'a Value> { + self.value().get("subject") + } + + pub(crate) fn node(self) -> Option<&'a Value> { + self.subject()?.get("node") + } + + pub(crate) fn node_name(self) -> Option<&'a str> { + self.node()?.get("name")?.as_str() + } + + /// The node's display label: its `meta.label`, else its name. + pub(crate) fn node_label(self) -> Option<&'a str> { + let node = self.node()?; + node.pointer("/meta/label") + .and_then(Value::as_str) + .filter(|label| !label.is_empty()) + .or_else(|| node.get("name")?.as_str()) + } + + pub(crate) fn node_kind(self) -> Option<&'a str> { + self.node()?.pointer("/meta/kind")?.as_str() + } + + /// Whether the subject's node is a logical stage: not a lowering node + /// (`synthetic`) and not a fork's `parallel.branch` delegate. + pub(crate) fn is_shown_stage(self) -> bool { + let Some(node) = self.node() else { + return false; + }; + let synthetic = node + .pointer("/meta/synthetic") + .and_then(Value::as_bool) + .unwrap_or(false); + !synthetic && self.node_kind() != Some("parallel.branch") + } + + pub(crate) fn visit(self) -> u64 { + self.subject() + .and_then(|subject| subject.get("visit")) + .and_then(Value::as_u64) + .unwrap_or(1) + } + + pub(crate) fn firing(self) -> Option { + self.subject()?.get("firing")?.as_u64() + } + + pub(crate) fn execution(self) -> Option { + self.value().pointer("/context/execution")?.as_u64() + } + + /// The stage key: `(execution, firing)`. + pub(crate) fn stage_key(self) -> Option { + Some(format!("{}:{}", self.execution()?, self.firing()?)) + } + + /// The stored platform record, tagged by `kind`. + pub(crate) fn platform_record(self) -> Option<&'a Value> { + (!self.is_petri()) + .then(|| self.value().get("record")) + .flatten() + } + + pub(crate) fn str_at(self, pointer: &str) -> Option<&'a str> { + self.value().pointer(pointer)?.as_str() + } +} + +/// Whether the item is the platform record of the run's terminal lifecycle +/// transition, which ends the attached stream. +pub(crate) fn is_terminal_lifecycle(item: &RunStreamItem) -> bool { + let view = PetriItem::new(item); + view.platform_record().is_some_and(|record| { + record["kind"].as_str() == Some("run.lifecycle") + && matches!( + record["transition"].as_str(), + Some("succeeded" | "failed" | "dead") + ) + }) +} + +/// The exit code the item decides, if it is one that ends the run: the +/// engine's `run.finished`, or the platform record of the terminal +/// lifecycle transition. +pub(crate) fn exit_code_of(item: &RunStreamItem) -> Option { + let view = PetriItem::new(item); + if view.is_petri() { + if view.name() != Some("run.finished") { + return None; + } + return Some(match view.body()?.get("status")?.as_str()? { + "success" => 0, + _ => 1, + }); + } + let record = view.platform_record()?; + if record["kind"].as_str() != Some("run.lifecycle") { + return None; + } + match record["transition"].as_str()? { + "succeeded" => Some(0), + "failed" | "dead" => Some(1), + _ => None, + } +} + +/// The question a `step.progress.recorded` carries, when Petri parsed one: +/// its id and text. +pub(crate) fn question_of(item: &RunStreamItem) -> Option<(&str, &str)> { + let view = PetriItem::new(item); + let parsed = view.parsed()?; + if parsed.get("kind")?.as_str()? != "question" { + return None; + } + let question = parsed.get("question")?; + Some(( + question.get("id")?.as_str()?, + question.get("text")?.as_str()?, + )) +} + +/// Whether the item closes the question with this id: a delivered answer, +/// its expiry, or Fabro's record of who answered. +pub(crate) fn resolves_question(item: &RunStreamItem, question_id: &str) -> bool { + let view = PetriItem::new(item); + if let Some(record) = view.platform_record() { + return record["kind"].as_str() == Some("interview.answered") + && record["question"].as_str() == Some(question_id); + } + match view.name() { + Some("control.requested") => view.str_at("/derived/answer/question") == Some(question_id), + Some("step.progress.recorded") => view.parsed().is_some_and(|parsed| { + parsed["kind"].as_str() == Some("question_expired") + && parsed["question"].as_str() == Some(question_id) + }), + _ => false, + } +} + +/// The raw line `run events` prints for an item: the envelope as JSON, +/// redacted. +pub(crate) fn raw_line(item: &RunStreamItem) -> Result { + let line = serde_json::to_string(item)?; + Ok(redact_jsonl_line(&line)) +} + +/// What the pretty printer remembers between items: when each firing +/// started, and when the run did. +#[derive(Default)] +pub(crate) struct PrettyState { + clock: StageClock, +} + +/// When each firing's visit started, by stage key, so its completion can +/// show a duration. +#[derive(Default)] +pub(crate) struct StageClock { + starts: HashMap, + run_start: Option, +} + +impl StageClock { + /// Note the item and answer how long its stage or the run has been + /// running, when the item ends one. + pub(crate) fn observe(&mut self, item: &RunStreamItem) -> Option { + let view = PetriItem::new(item); + match view.name()? { + "run.started" if view.is_petri() => { + self.run_start = Some(item.recorded_at); + None + } + "visit.started" => { + if let Some(key) = view.stage_key() { + self.starts.insert(key, item.recorded_at); + } + None + } + "visit.completed" | "branch.completed" => { + let key = view.stage_key()?; + let start = self.starts.remove(&key)?; + Some(item.recorded_at.saturating_sub(start)) + } + "run.finished" if view.is_petri() => { + Some(item.recorded_at.saturating_sub(self.run_start?)) + } + _ => None, + } + } +} + +/// The pretty line for an item, or nothing for one the terminal does not +/// show. Petri events render by `.` with the stage's label; +/// platform records by their kind. +pub(crate) fn format_pretty( + item: &RunStreamItem, + styles: &Styles, + state: &mut PrettyState, +) -> Option { + let elapsed = state.clock.observe(item); + let view = PetriItem::new(item); + let ts = view.recorded_at().format("%H:%M:%S").to_string(); + let ts = styles.dim.apply_to(&ts).to_string(); + if let Some(record) = view.platform_record() { + return format_platform_record(&ts, record, styles); + } + // A revisited node shows its visit beside its label, as a stage id does. + let label = match (view.node_label(), view.visit()) { + (Some(label), 1) => label.to_string(), + (Some(label), visit) => format!("{label}@{visit}"), + (None, _) => "?".to_string(), + }; + let label = label.as_str(); + match view.name()? { + "run.started" => Some(format!( + "{ts} {}", + styles.dim.apply_to("Engine: petri run started") + )), + "visit.started" if view.is_shown_stage() => Some(format!( + "{ts} {} {}", + styles.bold_cyan.apply_to("\u{25b6}"), + styles.bold.apply_to(label), + )), + "visit.completed" if view.is_shown_stage() => { + let derived = view.derived()?; + let outcome = derived.get("outcome")?; + let executed = derived + .get("executed") + .and_then(Value::as_bool) + .unwrap_or(true); + let status = outcome.get("status").and_then(Value::as_str).unwrap_or("?"); + let duration = format_duration_ms(elapsed.unwrap_or(0)); + if !executed || status == "skipped" { + return Some(format!( + "{ts} {} {} {}", + styles.dim.apply_to("\u{2298}"), + styles.bold.apply_to(label), + styles.dim.apply_to("skipped"), + )); + } + match status { + "success" | "partial_success" => { + let mut line = format!( + "{ts} {} {}", + styles.green.apply_to("\u{2713}"), + styles.bold.apply_to(label), + ); + if status == "partial_success" { + let _ = write!(line, " {}", styles.yellow.apply_to("partial")); + } + let _ = write!(line, " {duration}"); + if let Some(usage) = usage_summary(outcome, styles) { + let _ = write!(line, " {usage}"); + } + Some(line) + } + other => { + let error = outcome + .pointer("/failure/message") + .and_then(Value::as_str) + .unwrap_or(other); + Some(format!( + "{ts} {} {} {}", + styles.red.apply_to("\u{2717}"), + styles.bold.apply_to(label), + styles.red.apply_to(error), + )) + } + } + } + "retry.scheduled" => { + let derived = view.derived()?; + let attempt = derived.get("next_attempt").and_then(Value::as_u64)?; + let delay = 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); + Some(format!( + "{ts} {} {}: retrying (attempt {attempt}, delay {})", + styles.yellow.apply_to("\u{21bb}"), + label, + format_duration_ms(delay), + )) + } + "route.applied" => { + let derived = view.derived()?; + let target = derived.pointer("/target/name").and_then(Value::as_str)?; + let transition = derived + .get("transition") + .and_then(Value::as_str) + .unwrap_or("") + .to_lowercase(); + let back = derived + .get("back") + .and_then(Value::as_bool) + .unwrap_or(false); + let detail = if back { " (loop)" } else { "" }; + // Petri records the route after the next visit started, so the + // line names both ends of the edge. + Some(format!( + "{ts} {} {} {} {}{}", + styles.dim.apply_to(view.node_name().unwrap_or("?")), + styles.dim.apply_to("\u{2192}"), + target, + styles.dim.apply_to(&transition), + styles.dim.apply_to(detail), + )) + } + "fork.started" => { + let branches = view + .derived()? + .get("branches") + .and_then(Value::as_array) + .map_or(0, Vec::len); + Some(format!( + "{ts} {} {} {}", + styles.bold_cyan.apply_to("\u{2442}"), + styles.bold.apply_to(label), + styles.dim.apply_to(format!("{branches} branches")), + )) + } + "branch.completed" => { + let result = view.derived()?.get("result")?; + let name = result + .pointer("/node/name") + .and_then(Value::as_str) + .unwrap_or(label); + let status = result.get("status").and_then(Value::as_str).unwrap_or("?"); + let (glyph, style) = if status == "success" { + ("\u{2713}", &styles.green) + } else { + ("\u{2717}", &styles.red) + }; + Some(format!( + "{ts} {} branch {} {} {}", + style.apply_to(glyph), + name, + styles.dim.apply_to(status), + styles + .dim + .apply_to(format_duration_ms(elapsed.unwrap_or(0))), + )) + } + "fork.completed" => { + let derived = view.derived()?; + let results = derived + .get("results") + .and_then(Value::as_array) + .map_or(0, Vec::len); + let disposition = derived + .get("disposition") + .and_then(Value::as_str) + .unwrap_or("joined"); + Some(format!( + "{ts} {} {} {}", + styles.dim.apply_to("\u{2442}"), + styles + .dim + .apply_to(format!("{results} branches {disposition}")), + styles.bold.apply_to(label), + )) + } + "step.progress.recorded" => format_progress(&ts, view, label, styles), + "control.requested" => { + let answer = view.derived()?.get("answer")?; + let value = answer + .get("choice") + .or_else(|| answer.get("text")) + .and_then(Value::as_str) + .map_or_else(|| answer.to_string(), str::to_string); + let late = !view + .derived()? + .get("deliverable") + .and_then(Value::as_bool) + .unwrap_or(true); + let suffix = if late { " (late)" } else { "" }; + Some(format!( + "{ts} {} answered: {}{}", + styles.dim.apply_to("\u{21b3}"), + value, + styles.dim.apply_to(suffix), + )) + } + "invocation.cancel.requested" => { + let reason = view + .body()? + .get("reason") + .and_then(Value::as_str) + .unwrap_or("?"); + Some(format!( + "{ts} {} {}", + styles.bold_red.apply_to("\u{2717} Cancel requested"), + styles.dim.apply_to(reason), + )) + } + "run.finished" => { + let status = view.body()?.get("status").and_then(Value::as_str)?; + let duration = format_duration_ms(elapsed.unwrap_or(0)); + Some(match status { + "success" => format!( + "{ts} {} {}", + styles.bold_green.apply_to("\u{2713} SUCCEEDED"), + styles.bold.apply_to(&duration), + ), + "cancelled" => format!( + "{ts} {} {}", + styles.bold_red.apply_to("\u{2717} CANCELLED"), + styles.bold.apply_to(&duration), + ), + _ => format!( + "{ts} {} {}", + styles.bold_red.apply_to("\u{2717} FAILED"), + styles.bold.apply_to(&duration), + ), + }) + } + _ => None, + } +} + +/// The token and cost summary of a finished visit, from the Pebble or +/// prompt usage in its metrics. +fn usage_summary(outcome: &Value, styles: &Styles) -> Option { + let custom = outcome.pointer("/metrics/custom")?; + let usage = custom + .get("pebble.usage") + .or_else(|| custom.get("prompt.usage"))?; + let tokens = usage.get("tokens")?; + let input = tokens.get("input").and_then(Value::as_u64).unwrap_or(0); + let output = tokens.get("output").and_then(Value::as_u64).unwrap_or(0); + let total = input.saturating_add(output); + let mut parts = Vec::new(); + if let Some(cost) = usage.pointer("/cost/usd_micros").and_then(Value::as_u64) { + parts.push(format_usd_micros(cost)); + } + if total > 0 { + parts.push(format!("{total} toks")); + } + (!parts.is_empty()).then(|| styles.dim.apply_to(parts.join(" ")).to_string()) +} + +/// The line for a `step.progress.recorded`: a question, a log line, an +/// agent envelope, a prompt's completion. +fn format_progress(ts: &str, view: PetriItem<'_>, label: &str, styles: &Styles) -> Option { + if let Some(parsed) = view.parsed() { + match parsed.get("kind").and_then(Value::as_str) { + Some("question") => { + let question = parsed.get("question")?; + let text = question.get("text").and_then(Value::as_str).unwrap_or(""); + let mut line = format!( + "{ts} {} {}: {}", + styles.yellow.apply_to("?"), + styles.bold.apply_to(label), + text, + ); + if let Some(options) = question.get("options").and_then(Value::as_array) { + let labels: Vec<&str> = options + .iter() + .filter_map(|option| option.get("label").and_then(Value::as_str)) + .collect(); + if !labels.is_empty() { + let _ = write!(line, " {}", styles.dim.apply_to(labels.join(" "))); + } + } + return Some(line); + } + Some("question_expired") => { + let default = parsed + .get("default") + .and_then(Value::as_str) + .map_or_else(String::new, |default| format!(" (default {default})")); + return Some(format!( + "{ts} {} question expired{}", + styles.dim.apply_to("\u{21b3}"), + styles.dim.apply_to(&default), + )); + } + _ => {} + } + } + if let Some((_, line)) = view.log_line() { + return Some(format!( + "{ts} {} {}", + styles.dim.apply_to("\u{2502}"), + styles.dim.apply_to(line), + )); + } + let custom = view.custom()?; + match custom.get("kind").and_then(Value::as_str)? { + "pebble" => format_envelope(ts, custom.get("event")?, label, styles), + "attractor.prompt.completed" => { + let response = custom.get("response").and_then(Value::as_str).unwrap_or(""); + let header = format!("{ts} {} {}", "\u{1f4ac}", styles.bold.apply_to(label)); + let body = indented(styles, response, " "); + Some(format!("{header}\n{body}\n")) + } + "attractor.checkout" => { + let repository = custom + .get("repository") + .and_then(Value::as_str) + .unwrap_or("?"); + let commit = custom.get("commit").and_then(Value::as_str).unwrap_or(""); + Some(format!( + "{ts} Checkout: {} {}", + repository, + styles.dim.apply_to(short_sha(commit)), + )) + } + _ => None, + } +} + +/// The line for a Pebble coding-agent envelope: the assistant's text, a +/// tool call's start and end. +fn format_envelope(ts: &str, envelope: &Value, label: &str, styles: &Styles) -> Option { + let event = envelope.get("event")?.as_object()?; + let (variant, fields) = event.iter().next()?; + match variant.as_str() { + "AssistantMessage" => { + let model = fields.get("model").and_then(Value::as_str).unwrap_or("?"); + let text = fields.get("text").and_then(Value::as_str).unwrap_or(""); + let header = format!( + "{ts} {} {} {}", + "\u{1f4ac}", + styles.bold.apply_to(label), + styles.dim.apply_to(format!("[{model}]")), + ); + let body = indented(styles, text, " "); + Some(format!("{header}\n{body}\n")) + } + "ToolCallStarted" => { + let tool = fields + .get("tool_name") + .and_then(Value::as_str) + .unwrap_or("?"); + Some(format!( + "{ts} {} {}", + styles.dim.apply_to("\u{2699}"), + styles.dim.apply_to(tool), + )) + } + "ToolCallCompleted" => { + let tool = fields + .get("tool_name") + .and_then(Value::as_str) + .unwrap_or("?"); + let is_error = fields + .get("is_error") + .and_then(Value::as_bool) + .unwrap_or(false); + let (glyph, style) = if is_error { + ("\u{2717}", &styles.red) + } else { + ("\u{2713}", &styles.green) + }; + Some(format!("{ts} {} {}", style.apply_to(glyph), tool)) + } + _ => None, + } +} + +/// The line for a platform record, by its kind. +fn format_platform_record(ts: &str, record: &Value, styles: &Styles) -> Option { + let kind = record.get("kind")?.as_str()?; + match kind { + "run.created" => { + let spec = record.get("spec"); + let name = record + .get("title") + .and_then(Value::as_str) + .or_else(|| spec?.pointer("/settings/workflow/name")?.as_str()) + .or_else(|| spec?.get("workflow_slug")?.as_str()) + .unwrap_or("Run"); + let run_id = spec + .and_then(|spec| spec.get("run_id")) + .and_then(Value::as_str) + .unwrap_or("?"); + let header = format!( + "{ts} {} {} {}", + styles.bold_cyan.apply_to("\u{25b6}"), + styles.bold.apply_to(name), + styles.dim.apply_to(run_id), + ); + match spec + .and_then(|spec| spec.pointer("/settings/run/goal")) + .and_then(Value::as_str) + { + Some(goal) if !goal.is_empty() => { + let body = indented(styles, goal, " "); + Some(format!("{header}\n{body}\n")) + } + _ => Some(header), + } + } + "run.lifecycle" => { + let transition = record.get("transition").and_then(Value::as_str)?; + let reason = record + .get("reason") + .and_then(Value::as_str) + .map_or_else(String::new, |reason| format!(": {reason}")); + Some(format!( + "{ts} {}", + styles + .dim + .apply_to(format!("\u{00b7} {transition}{reason}")) + )) + } + "run.notice" => { + let level = record + .get("level") + .and_then(Value::as_str) + .unwrap_or("info"); + let code = record.get("code").and_then(Value::as_str).unwrap_or(""); + let message = record.get("message").and_then(Value::as_str).unwrap_or(""); + let label = match level { + "warn" => styles.yellow.apply_to("Warning:").to_string(), + "error" => styles.bold_red.apply_to("Error:").to_string(), + _ => styles.bold.apply_to("Info:").to_string(), + }; + let code_suffix = if code.is_empty() { + String::new() + } else { + format!(" {}", styles.dim.apply_to(format!("[{code}]"))) + }; + Some(format!("{ts} {label} {message}{code_suffix}")) + } + "checkpoint" => { + let sha = record + .get("git_commit_sha") + .and_then(Value::as_str) + .map_or_else(|| "(no commit)".to_string(), short_sha); + Some(format!( + "{ts} {} {}", + styles.dim.apply_to("\u{2398} Checkpoint"), + styles.dim.apply_to(sha), + )) + } + "pull_request.created" => { + let url = record + .get("html_url") + .and_then(Value::as_str) + .unwrap_or("?"); + let draft = record + .get("draft") + .and_then(Value::as_bool) + .unwrap_or(false); + let suffix = if draft { " (draft)" } else { "" }; + Some(format!( + "{ts} Pull request: {}{}", + url, + styles.dim.apply_to(suffix) + )) + } + "interview.answered" => { + let principal = record.get("principal"); + let who = principal + .and_then(|principal| principal.get("login")) + .or_else(|| principal?.get("kind")) + .and_then(Value::as_str) + .unwrap_or("?"); + Some(format!( + "{ts} {} answered by {}", + styles.dim.apply_to("\u{21b3}"), + styles.dim.apply_to(who), + )) + } + "run.title" => { + let title = record.get("title").and_then(Value::as_str).unwrap_or(""); + Some(format!("{ts} Title: {title}")) + } + "run.branch" => { + let branch = record + .get("run_branch") + .and_then(Value::as_str) + .unwrap_or("?"); + let base = record + .get("base_sha") + .and_then(Value::as_str) + .map_or_else(String::new, |sha| format!(" from {}", short_sha(sha))); + Some(format!( + "{ts} Branch: {}{}", + branch, + styles.dim.apply_to(&base) + )) + } + "git.identity" => { + let identity = record.get("identity")?; + let name = identity.get("name").and_then(Value::as_str).unwrap_or("?"); + let email = identity.get("email").and_then(Value::as_str).unwrap_or("?"); + let source = identity + .get("source") + .and_then(Value::as_str) + .unwrap_or("?"); + Some(format!( + "{ts} Git identity: {name} <{email}> {}", + styles.dim.apply_to(source) + )) + } + other => Some(format!( + "{ts} {}", + styles.dim.apply_to(format!("\u{00b7} {other}")) + )), + } +} + +fn short_sha(sha: &str) -> String { + sha.chars().take(7).collect() +} + +fn indented(styles: &Styles, text: &str, indent: &str) -> String { + let wrap_width = Styles::terminal_width().saturating_sub(indent.len()); + styles + .render_markdown_width(text, wrap_width) + .lines() + .map(|line| format!("{indent}{line}")) + .collect::>() + .join("\n") +} + +fn millis(recorded_at: u64) -> DateTime { + Utc.timestamp_millis_opt(i64::try_from(recorded_at).unwrap_or(i64::MAX)) + .single() + .unwrap_or_default() +} + +/// `run events` for a Petri run: the stream, filtered by `--since` and +/// `--tail`, raw or pretty, then followed live when asked. +pub(crate) async fn print_events( + client: &server_client::Client, + run_id: &RunId, + args: &EventsArgs, + since: Option>, + pretty: bool, + styles: &Styles, +) -> Result<()> { + let items = client + .list_run_stream(run_id, 0) + .await + .context("Failed to list the run's stream")?; + let last_seq = items.last().map_or(0, |item| item.stream_seq); + let mut selected: Vec<&RunStreamItem> = items + .iter() + .filter(|item| since.is_none_or(|cutoff| PetriItem::new(item).recorded_at() >= cutoff)) + .collect(); + if let Some(tail) = args.tail { + let start = selected.len().saturating_sub(tail); + selected.drain(..start); + } + + let stdout = io::stdout(); + let mut out = stdout.lock(); + let mut state = PrettyState::default(); + for item in selected { + write_item(&mut out, item, pretty, styles, &mut state)?; + } + out.flush()?; + + if args.follow { + Box::pin(follow(client, run_id, last_seq, pretty, styles, state)).await?; + } + Ok(()) +} + +fn write_item( + out: &mut dyn Write, + item: &RunStreamItem, + pretty: bool, + styles: &Styles, + state: &mut PrettyState, +) -> Result<()> { + if pretty { + if let Some(line) = format_pretty(item, styles, state) { + writeln!(out, "{line}")?; + } + } else { + writeln!(out, "{}", raw_line(item)?)?; + } + Ok(()) +} + +/// Follow the stream live from `after`: attach, print each item, and on a +/// stream the server ended before the run's terminal record, reconnect +/// from the last `stream_seq` printed unless the run has concluded. +async fn follow( + client: &server_client::Client, + run_id: &RunId, + after: u64, + pretty: bool, + styles: &Styles, + mut state: PrettyState, +) -> Result<()> { + let stdout = io::stdout(); + let mut out = stdout.lock(); + let mut cursor = after; + loop { + let mut stream = client.attach_run_stream(run_id, Some(cursor)).await?; + while let Some(item) = stream.next_item().await? { + cursor = item.stream_seq; + write_item(&mut out, &item, pretty, styles, &mut state)?; + out.flush()?; + if is_terminal_lifecycle(&item) { + return Ok(()); + } + } + let run_state = client + .get_run_state(run_id) + .await + .context("Failed to read run state from server while following the stream")?; + if run_state.status.is_terminal() { + // The server ended the stream after its grace: print what + // landed since, if anything, and stop. + for item in client.list_run_stream(run_id, cursor).await? { + write_item(&mut out, &item, pretty, styles, &mut state)?; + } + out.flush()?; + debug!("Run reached terminal status and the stream ended, stopping follow"); + return Ok(()); + } + debug!( + cursor, + "the attached stream ended before the run did; reconnecting" + ); + time::sleep(FOLLOW_RECONNECT_DELAY).await; + } +} + +#[cfg(test)] +mod tests { + use fabro_types::fixtures; + use serde_json::json; + + use super::*; + + fn item(kind: RunStreamItemKind, stream_seq: u64, value: Value) -> RunStreamItem { + RunStreamItem { + run_id: fixtures::RUN_1, + stream_seq, + kind, + id: stream_seq.to_string(), + recorded_at: 1_789_706_579_000 + stream_seq * 1000, + item: value, + } + } + + fn petri(stream_seq: u64, value: Value) -> RunStreamItem { + item(RunStreamItemKind::Petri, stream_seq, value) + } + + fn platform(stream_seq: u64, record: &Value) -> RunStreamItem { + item( + RunStreamItemKind::Platform, + stream_seq, + json!({"seq": stream_seq, "recorded_at": 0, "record": record}), + ) + } + + fn subject(name: &str, kind: &str) -> Value { + json!({ + "node": {"id": 2, "name": name, "kind": "attractor/x", "meta": {"label": name, "kind": kind}}, + "firing": 2, "visit": 1, "attempt": 1, "generation": 0, "branch": {"role": "none"} + }) + } + + fn view_event(name: &str, node: &str, kind: &str, derived: Value) -> Value { + let mut derived = derived; + derived["event"] = json!(name); + json!({ + "id": {"log": "execution", "execution": 0, "seq": 5, "index": 1}, + "origin": "derived", + "context": {"invocation": 0, "execution": 0}, + "subject": subject(node, kind), + "derived": derived + }) + } + + #[test] + fn a_stage_renders_its_start_and_its_end_with_the_elapsed_time() { + let styles = Styles::new(false); + let mut state = PrettyState::default(); + let started = petri(1, view_event("visit.started", "say", "command", json!({}))); + let completed = petri( + 4, + view_event( + "visit.completed", + "say", + "command", + json!({ + "outcome": {"status": "success", "metrics": {"duration_ms": 42}}, + "executed": true, "attempts": 1 + }), + ), + ); + let start_line = format_pretty(&started, &styles, &mut state).expect("a start line"); + assert!(start_line.contains("\u{25b6} say"), "{start_line}"); + let end_line = format_pretty(&completed, &styles, &mut state).expect("an end line"); + assert!(end_line.contains("\u{2713} say"), "{end_line}"); + assert!( + end_line.contains(" 3s"), + "the visit took three seconds: {end_line}" + ); + } + + #[test] + fn a_fork_delegate_is_not_a_stage_but_its_branch_completion_is_shown() { + let styles = Styles::new(false); + let mut state = PrettyState::default(); + let delegate = petri( + 1, + view_event("visit.started", "a", "parallel.branch", json!({})), + ); + assert!(format_pretty(&delegate, &styles, &mut state).is_none()); + let completed = petri( + 3, + view_event( + "branch.completed", + "a", + "parallel.branch", + json!({ + "result": {"node": {"name": "a"}, "status": "success"} + }), + ), + ); + let line = format_pretty(&completed, &styles, &mut state).expect("a branch line"); + assert!(line.contains("branch a"), "{line}"); + assert!(line.contains(" 2s"), "{line}"); + } + + #[test] + fn a_platform_notice_and_the_terminal_lifecycle_record_render_by_kind() { + let styles = Styles::new(false); + let mut state = PrettyState::default(); + let notice = platform( + 2, + &json!({"kind": "run.notice", "level": "warn", "code": "x.y", "message": "careful"}), + ); + let line = format_pretty(¬ice, &styles, &mut state).expect("a notice line"); + assert!(line.contains("Warning: careful [x.y]"), "{line}"); + let finished = platform( + 3, + &json!({"kind": "run.lifecycle", "transition": "succeeded", "status": {"kind": "succeeded"}}), + ); + assert!(is_terminal_lifecycle(&finished)); + assert_eq!(exit_code_of(&finished), Some(0)); + let line = format_pretty(&finished, &styles, &mut state).expect("a lifecycle line"); + assert!(line.contains("\u{00b7} succeeded"), "{line}"); + } + + #[test] + fn a_question_is_found_and_closed_by_its_answer_or_by_who_answered() { + let asked = petri( + 5, + json!({ + "origin": "external", + "context": {"invocation": 0, "execution": 0}, + "subject": subject("gate", "human"), + "record": {"seq": 12, "body": {"event": "step.progress.recorded", "firing": 2, + "ev": {"custom": {"$question": {"id": "gate#2", "text": "Go?"}}}}}, + "derived": {"parsed": {"kind": "question", "question": {"id": "gate#2", "text": "Go?", + "options": [{"key": "Y", "label": "[Y] Yes"}], "kind": "yes_no"}}} + }), + ); + assert_eq!(question_of(&asked), Some(("gate#2", "Go?"))); + let styles = Styles::new(false); + let mut state = PrettyState::default(); + let line = format_pretty(&asked, &styles, &mut state).expect("a question line"); + assert!(line.contains("? gate: Go? [Y] Yes"), "{line}"); + + let answered = petri( + 6, + json!({ + "origin": "external", + "context": {"invocation": 0, "execution": 0}, + "subject": subject("gate", "human"), + "record": {"seq": 14, "body": {"event": "control.requested", "firing": 2, + "ctl": {"deliver": {"$answer": {"question": "gate#2", "choice": "N"}}}}}, + "derived": {"deliverable": true, "answer": {"question": "gate#2", "choice": "N"}} + }), + ); + assert!(resolves_question(&answered, "gate#2")); + assert!(!resolves_question(&answered, "gate#3")); + let who = platform( + 7, + &json!({"kind": "interview.answered", "question": "gate#2", "principal": {"kind": "user", "login": "dev"}}), + ); + assert!(resolves_question(&who, "gate#2")); + let line = format_pretty(&who, &styles, &mut state).expect("an answered-by line"); + assert!(line.contains("answered by dev"), "{line}"); + } + + #[test] + fn the_engine_finish_decides_the_exit_code() { + let finished = petri( + 9, + json!({ + "origin": "external", "context": {}, + "record": {"seq": 6, "body": {"event": "run.finished", "status": "failed"}} + }), + ); + assert_eq!(exit_code_of(&finished), Some(1)); + let styles = Styles::new(false); + let mut state = PrettyState::default(); + let line = format_pretty(&finished, &styles, &mut state).expect("a finish line"); + assert!(line.contains("\u{2717} FAILED"), "{line}"); + } + + #[test] + fn the_raw_line_is_the_envelope_as_json() { + let notice = platform( + 2, + &json!({"kind": "run.notice", "level": "info", "message": "m"}), + ); + let line = raw_line(¬ice).expect("a line"); + let value: Value = serde_json::from_str(&line).expect("json"); + assert_eq!(value["stream_seq"], 2); + assert_eq!(value["kind"], "platform"); + assert_eq!(value["item"]["record"]["kind"], "run.notice"); + } +} diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs index e98535eac..93a7fc7cb 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs @@ -388,7 +388,23 @@ fn agent_progress_event( stored: &RunEvent, event: &CodingEvent, ) -> Option { - 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>, + event: &CodingEvent, +) -> Option { 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 diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index 04b838b95..020ea4332 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -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 { diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/petri.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/petri.rs new file mode 100644 index 000000000..502b5c062 --- /dev/null +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/petri.rs @@ -0,0 +1,264 @@ +//! The progress lines a Petri run's stream items mean: the mapping from +//! Petri's `.` 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 { + 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| 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 { + 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 { + 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 { + 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, + } +} diff --git a/lib/apps/fabro-cli/src/server_client.rs b/lib/apps/fabro-cli/src/server_client.rs index 0d20ac36d..d8ad0a9cc 100644 --- a/lib/apps/fabro-cli/src/server_client.rs +++ b/lib/apps/fabro-cli/src/server_client.rs @@ -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; diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index d6f14a141..777ee2be8 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -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 { - 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 { + 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 +/// `.` name (`question` and `question_expired` for the parsed +/// progress payloads), a platform lifecycle record as +/// `lifecycle:`, another platform record by its kind. +fn stream_names(items: &[serde_json::Value]) -> Vec { + 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 `, 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 { + 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 `.` 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 `.` 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 = lines + .iter() + .map(|line| line["stream_seq"].as_u64().expect("a stream_seq")) + .collect(); + let expected: Vec = (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 = 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(); +}