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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 01:25:40 -04:00
parent 1bfe62577d
commit bfc2abebab
No known key found for this signature in database
2 changed files with 50 additions and 7 deletions

View file

@ -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<serde_json::Value> {
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 <first 12 of the run id> <phase>`, 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.

View file

@ -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.