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