From bfc2abebabb94e9b99130f4e20bcc60b7b4dc15b Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 01:25:40 -0400 Subject: [PATCH] Wait for the terminal lifecycle record before reading a stream's end Fabro's terminal `run.lifecycle` record lands a moment after Petri's `run.finished`: the worker exits, the server records the status, the projector folds it. A CLI scenario that asserts on the end of the stream now waits for that record instead of reading the stream as soon as the runs row turns `succeeded`, which the projector writes from the engine's finish alone. The fabro-petri README names the projector's stream reader and commit signal, the server's reconnect test with its fixture capture, and the CLI scenarios that read a run back through the stream. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-cli/tests/it/scenario/petri.rs | 28 ++++++++++++++++-- lib/components/fabro-petri/README.md | 29 +++++++++++++++---- 2 files changed, 50 insertions(+), 7 deletions(-) diff --git a/lib/apps/fabro-cli/tests/it/scenario/petri.rs b/lib/apps/fabro-cli/tests/it/scenario/petri.rs index 777ee2be8..68c594619 100644 --- a/lib/apps/fabro-cli/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-cli/tests/it/scenario/petri.rs @@ -416,6 +416,29 @@ fn count_of(names: &[String], expected: &str) -> usize { names.iter().filter(|name| *name == expected).count() } +/// The run's whole stream once it is settled: Fabro's terminal lifecycle +/// record lands a moment after the engine's finish (the worker exits, the +/// server records the status, the projector folds it), so a reader that +/// wants the end of the stream waits for that record. +async fn settled_stream(server: &RunningServer, run_id: &str) -> Vec { + let deadline = Instant::now() + RUN_TIMEOUT; + loop { + let items = run_stream(server, run_id).await; + let names = stream_names(&items); + if names + .iter() + .any(|name| matches!(name.as_str(), "lifecycle:succeeded" | "lifecycle:failed")) + { + return items; + } + assert!( + Instant::now() < deadline, + "run {run_id} never recorded its terminal lifecycle transition: {names:?}" + ); + tokio::time::sleep(POLL).await; + } +} + /// 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. @@ -484,7 +507,7 @@ 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 names = stream_names(&run_stream(&server, &run_id).await); + let names = stream_names(&settled_stream(&server, &run_id).await); assert_eq!(count_of(&names, "lifecycle:succeeded"), 1, "{names:?}"); assert_eq!(count_of(&names, "run.finished"), 1, "{names:?}"); assert!( @@ -565,7 +588,7 @@ 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 items = run_stream(&server, &run_id).await; + let items = settled_stream(&server, &run_id).await; let names = stream_names(&items); assert_eq!( status, @@ -932,6 +955,7 @@ async fn a_finished_petri_run_reads_back_through_the_cli() { "server stderr:\n{}", server.stderr_text() ); + settled_stream(&server, &run_id).await; let target = server.target(); // Raw: the envelope, one item per line, dense and in order. diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 451686873..d613eccab 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -87,7 +87,14 @@ Every adapter the integration plan describes lands here. every Petri run, at startup. A run that executes in the server process goes through `Projector::observe_store`, which signals after each append. A torn tail (a record Petri cannot read) holds the view where it stands - and reports the run incomplete with the reason. + and reports the run incomplete with the reason. The projector also serves + the stream back (`Projector::stream_after`, one `RunStreamItem` per row: + `run_id`, `stream_seq`, `kind`, the item's own `id`, `recorded_at`, the + item) and signals its readers after each committed pass + (`Projector::subscribe`), which is how `GET /runs/{id}/events` pages a + Petri run by `after` and `GET /runs/{id}/attach` follows it live. The + version of Petri's event contract the stream carries is + `petri::EVENT_CONTRACT_VERSION`. - The platform adapters the plan adds after it: hooks and the run tools. ### What the projection leaves default @@ -181,13 +188,25 @@ serving the projection over Petri's records; a human gate is answered through the questions API; and Petri's diagnostics refuse a run at create. The server's `petri_runs` unit tests cover the lease ending at worker exit and the restart reconcile that relaunches a worker in resume mode. +`lib/apps/fabro-server/tests/it/scenario/petri_stream.rs` covers the stream: +a client attached to a two-branch parallel run disconnects once both +branches started, a platform notice is recorded while both branch scripts +run, the client reconnects from its last `stream_seq`, and the union of +what it saw is the whole stream, every item once, in order, with the notice +between the branch events and the same as the paged listing. With +`FABRO_CAPTURE_PETRI_FIXTURES` set, the scenarios write their settled +projection and stream under `apps/fabro-web/app/test-fixtures/petri/`, +which the web app's rendering tests read. The worker path is covered with the real binary in `lib/apps/fabro-cli/tests/it/scenario/petri.rs`: a command-only Petri run executes in the worker a foreground server launched, its records reach `petri_records` over the HTTP store and its lease ends with the worker; and a run whose server and worker are both killed mid-stage resumes in a new -worker after the server restarts, with one `run.completed`; a human gate in -the worker is answered through the questions API over the control channel; -two parallel gates each bind their own answer; and an unanswered gate -expires with its default. +worker after the server restarts, with one terminal lifecycle record; a +human gate in the worker is answered through the questions API over the +control channel; two parallel gates each bind their own answer; and an +unanswered gate expires with its default. The same file reads a finished +run back through the CLI (`events` raw, tail and `--pretty`, `attach`, +`wait`, `inspect`), answers a gate from an attached terminal, and follows +a run live with `events --follow` to its end.