Describe the run stream where the docs described the legacy event log

`docs/internal/events-strategy.md` is now the run stream strategy: the
two logs (Petri's records and Fabro's platform records), the projector
that folds them and assigns `stream_seq`, how to record a fact Petri
cannot know, how to read the stream, and the Ask Fabro session log.
`docs/internal/events.md` (the 106-event catalog) and the event schema
v2 shape document described the deleted `EventBody` model and are
deleted; AGENTS.md routes to the strategy for platform records and
stream consumers. The testing strategy's `progress.jsonl` rules name
records and stream items instead, and the public API nav drops the
removed per-stage events endpoint.

The interview adapter's module docs and the fabro-petri README no
longer claim the adapter posts `interview.*` events: readers see a
question in Petri's own progress record, and the server records who
answered as the `interview.answered` platform record.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 14:45:13 -04:00
parent 3e904c7121
commit 2f4888c199
No known key found for this signature in database
6 changed files with 121 additions and 251 deletions

View file

@ -273,7 +273,7 @@ When working in an area covered by a strategy doc, read the relevant document
**before** making changes:
- **`docs/internal/logging-strategy.md`** — read when adding `tracing` calls (`info!`, `debug!`, `warn!`, `error!`), working on error handling paths, or adding new operations that should be observable
- **`docs/internal/events-strategy.md`** — read when adding or modifying `Event` variants, touching `Emitter`/`emit()`, changing `progress.jsonl` output, or adding new workflow stage types
- **`docs/internal/events-strategy.md`** — read when adding a platform record kind, changing the projection fold or the run stream, or writing a consumer that matches on stream items
- **`docs/internal/testing-strategy.md`** — read when adding or reorganizing tests, choosing between unit vs `tests/it`, deciding whether a test belongs in `cmd` vs `workflow` vs `scenario`, or deciding how to structure snapshots and fixtures
- **`docs/internal/server-secrets-strategy.md`** — read when adding or changing server-level secrets, startup validation, install-time secret persistence, or subprocess env inheritance/scrubbing
- **`docs/internal/migrations-strategy.md`** — read when adding or changing temporary compatibility migrations, startup/file rewrites, migration runners, backups, or removal deadlines

View file

@ -1,224 +1,105 @@
# Fabro Events Strategy
# Fabro Run Stream Strategy
Fabro emits structured **workflow run events** during execution for observability. Events are the durable audit trail for a run: they drive the run store, SSE streaming, CLI progress rendering, and optional JSONL sinks.
A run's history is two logs, and its public event API is one ordered stream
over both:
Events are distinct from tracing logs. Tracing is developer diagnostics; events are product-facing state transitions and activity records that other systems consume.
- **Petri's records.** The engine writes every fact about execution: the run
starting and finishing, each step's firing, progress, output and outcome, a
question asked and answered, a scope acquired. Fabro stores them unchanged
in `petri_records` through `fabro-petri`'s `SqliteRunStore`, and reads them
through Petri's event contract (`RunEvent`, with its `derived` view).
- **Platform records.** Facts Fabro knows and Petri does not: the lifecycle
before and after the engine (`run.created`, `run.lifecycle`, `run.title`,
`run.parent`, `run.archived`, `run.superseded`, `run.notice`), who answered
a question (`interview.answered`), the branch and git identity a run works
under, a checkpoint commit, the pull request requests and outcomes, a
notification sent, a pairing. They are `PlatformRecord` values in
`fabro-store::platform_records`, stored in `platform_records` with a
per-run `seq`.
Detached runs rely on this distinction. If something needs to be visible after reattach, emit a `Event` rather than only logging to stderr or `detach.log`.
The **projector** (`fabro-petri::projection`) folds both logs into the run's
`RunProjection`, the view `GET /runs/{id}/state` serves, and assigns each
record it consumes a `stream_seq` in `petri_stream`. That stream is what
`GET /runs/{id}/events` and the attach stream serve, item by item, as
`RunStreamItem`: `{run_id, stream_seq, kind: petri|platform, id, recorded_at,
item}`. `stream_seq` is the cursor a client resumes from; `id` is the item's
own identity (`<log>/<seq>/<index>` for a Petri event, the record's `seq` for
a platform record) for deduplication.
## Architecture
Tracing logs are separate. Tracing is developer diagnostics; the stream is
the product-facing record other systems consume. If something must be
visible after a reattach, it has to be a record, not a log line.
```text
Engine/Handler -> Event -> Emitter::emit()
|- trace(raw event)
|- canonicalize -> RunEvent
`- on_event(&RunEvent)
|- run store
|- SSE
|- optional JSONL/debug sinks
`- CLI / tests / metrics listeners
```
## Recording a fact
The canonical `RunEvent` is built exactly once in the `fabro-workflow::event` module.
Petri's own facts need nothing from Fabro: the engine records them and the
projector's fold reads them. Add Fabro code only for a fact Petri cannot
know.
- `Event` (in `fabro-workflow`) is the internal typed event emitted by engine and handlers.
- `Emitter` owns an immutable `run_id` and converts `Event` into `RunEvent` via `to_run_event_at()`.
- `RunEvent` (in `fabro-types`) holds envelope metadata plus a typed `body: EventBody`. It has no cached JSON fields; the wire format is produced only during serialization.
- Every listener receives `&RunEvent`, not `&Event`.
- Bypass paths that cannot go through the emitter must call `to_run_event()` once and reuse the same `RunEvent` for every sink.
To record such a fact:
## Canonical Envelope
1. Add a variant to `PlatformRecord` and its kind to `PlatformRecordKind` in
`lib/components/fabro-store/src/platform_records.rs`. The kind is the
`kind` tag on the wire, lowercase dot notation (`pull_request.created`).
A record that belongs to a stage names its `execution` and `firing`.
2. Append it through the server's `run_records` module (or the worker's
client), which commits the record and wakes the projector. Never write
`platform_records` from anywhere else.
3. Fold it in `fabro-petri::projection` when the projection should show it.
A record nobody reads from the projection still reaches the stream.
4. Update the readers that match on record kinds: the CLI's pretty stream
rendering (`petri_stream.rs`), the web app's stream handling, the Slack
service, and the tests or fixtures that name kinds.
Each serialized `RunEvent` uses this canonical envelope:
Do not add a platform record that restates a Petri record. The projection
already carries what the engine knows; read it there.
```json
{
"id": "01960d0c-5d16-7d6e-8f61-9fd6f4a532b5",
"ts": "2026-03-30T12:00:01.000Z",
"run_id": "01JQ...",
"event": "agent.tool.started",
"session_id": "ses_child",
"parent_session_id": "ses_parent",
"node_id": "code",
"node_label": "Code",
"actor": {
"kind": "agent",
"session_id": "ses_child",
"parent_session_id": "ses_parent",
"model": "gpt-5.2"
},
"properties": {
"tool_name": "read_file",
"tool_call_id": "call_1",
"arguments": {"path": "src/main.rs"}
}
}
```
## Reading the stream
Always-present fields:
Servers and workers hold the stream through the projector: `stream_after`
for a page, `subscribe` for live items, `stream_head` for the cursor to
start from. The server's `stream_follower` reads every run's stream once and
fans it out to the in-memory run map and the global broadcast that `/attach`
and the Slack service take their items from.
| Field | Type | Notes |
|---|---|---|
| `id` | string | UUIDv7 event id |
| `ts` | string | UTC timestamp with millisecond precision |
| `run_id` | string | Workflow run id |
| `event` | string | Lowercase dot-notation event name |
Clients read `GET /runs/{id}/events?after=<stream_seq>` for a page and the
attach stream for live items; `fabro-client` exposes `list_run_stream`,
`list_run_stream_until` and `attach_run_stream`.
Optional top-level fields:
When matching items:
| Field | When present |
|---|---|
| `session_id` | Agent/session events |
| `parent_session_id` | Forwarded child-session events |
| `node_id` | Events tied to a graph node or branch |
| `node_label` | Display label for `node_id`; omitted when not applicable |
| `actor` | The principal responsible for the event |
- A Petri event's name is `item.record.body.event` (`run.started`,
`step.started`, `step.progress.recorded`, `step.finished`,
`run.finished`); its parsed meaning is under `item.derived` (a pending
question is `derived.parsed.kind == "question"`).
- A platform record's kind is `item.record.kind`.
- The run has ended when a platform `run.lifecycle` record's `transition`
is `succeeded`, `failed` or `dead`. Petri's `run.finished` precedes it and
carries the engine's own status.
Everything else lives inside `properties`.
Never rebuild an item downstream: pass the `RunStreamItem` through as read.
Important rules:
## Agent events
- Optional envelope fields are omitted, not serialized as `null`.
- Event-specific fields do not get flattened into the top level.
- Actor identity normally lives only in top-level `actor: Principal`; do not duplicate it in
event-specific properties. The exception is `run.created`, whose
`properties.provenance.subject` is the durable run creator stored in `RunSpec`; its envelope
`actor` is derived from the same principal.
- User actors must carry canonical IdP identity through `Principal::User { identity, login, auth_method }`, not a login-only string.
- `EventPayload` validation requires `id`, `ts`, `run_id`, and `event`.
Pebble's `CodingAgentEvent` stream is the agent event contract. Petri stores
each event a coding agent publishes for a step as that step's progress, and
the projector folds them into `StageProjection.agent` with pebble's
`SessionProjection`. Read `StageProjection.agent`, or the stored progress
record itself, instead of folding the stream again. Fabro adds nothing of
its own to this stream.
## Naming
## Ask Fabro sessions
The external event name is lowercase dot notation, for example:
Ask Fabro sessions are not runs. Their events (`run.session.*`) live in
their own log, `run_session_events`, through `RunSessionEventStore`, numbered
per session and served by the sessions API. They never enter a run's stream.
- `run.started`
- `stage.completed`
- `agent.tool.started`
- `sandbox.ready`
- `parallel.branch.completed`
## Persistence guarantees
`event_name()` in the `fabro-workflow::event` module is exhaustive. Do not use wildcard fallthroughs when adding new variants.
A record is committed before it is visible: the projector reads only what
the store has committed, and the stream's `stream_seq` is assigned in the
same transaction as the projection that consumed the record. A client that
resumes from its last `stream_seq` sees every item exactly once.
## Node And Session Metadata
`node_id` is the stable graph identifier. `node_label` is the human-facing display name. Stage events should surface both through the envelope when applicable.
Agent events now use explicit session links:
- `session_id` identifies the session that originally emitted the event.
- `parent_session_id` identifies the immediate parent session for forwarded child events.
- Nested sub-agents preserve immediate parentage across boundaries.
`AgentEvent::SubAgentEvent` no longer exists. Child activity is forwarded as normal agent events with session linkage in the envelope.
## Direct-Write Paths
Most events flow through `Emitter::emit()`. The remaining direct-write paths must use:
1. `to_run_event(run_id, event)`
2. Serialize and redact once
3. Reuse that exact `RunEvent` for every sink
Never build the same `RunEvent` twice if multiple sinks receive it.
## Adding A New Event
### 1. Add the typed event
Add a variant to `Event`, `AgentEvent`, or `SandboxLifecycle` as appropriate. Sandbox
facts come from two places: the pipeline emits `Initializing`, `Ready`, and
`InitializeFailed` around bringing the sandbox up, and the sandbox driver's own events
(operations and their outcome, progress inside a create such as an image pull, snapshot
builds, state observations, notices) are stored whole as `Event::SandboxDriver` by the
`DriverEventRecorder` in the `fabro-workflow::event` module. Their names derive from the
event (`fabro_types::sandbox_driver_event_name`): `<subject>.<action>.<phase>` such as
`sandbox.stop.completed` or `snapshot.create.started`, `<subject>.state`, and
`<subject>.notice`; their `properties` are the driver's event as the driver serializes
it, so the driver's `Event` is part of fabro's stored format. Fabro-sandbox emits no
events of its own.
### 2. Add tracing
Extend `Event::trace()` so the raw event is observable in tracing output.
### 3. Add an external name
Extend `event_name()` with the new lowercase dot-notation string.
### 4. Add the `EventBody` variant
Add a variant to `EventBody` in `fabro-types/src/run_event/mod.rs` with a corresponding props struct. Use `#[serde(rename = "dotted.name")]` matching the external name from step 3.
### 5. Map envelope fields and construct `EventBody`
Update `stored_event_fields()` and `event_body_from_event()` in the `fabro-workflow::event` module:
- Move `node_id`, `node_label`, `session_id`, and `parent_session_id` into the envelope when appropriate.
- Construct the `EventBody` variant directly from the `Event` fields.
- For `Event::Agent` sub-variants, merge `visit` into the inner props and lift `stage` to `node_id`.
- For `Event::Sandbox` sub-variants, unwrap and flatten into the corresponding `EventBody` variant.
### 6. Emit it
Prefer `Emitter::emit(&Event::...)`.
Use `to_run_event()` only for true bypass paths.
For cache-backed lifecycle work, emit slow-path start events only when the operation actually misses cache or waits on remote state. Completion events should represent a real ensure step (inspect, build, pull, or poll), not a configured no-op.
### 7. Update consumers
Check:
- CLI progress parsing
- `fabro events`
- store validation
- tests or fixtures that inspect event names or fields
## Agent Events
Pebble's `CodingAgentEvent` stream is the agent event contract. The worker's
event sink stores every event the coding agent publishes for a stage, except
streaming deltas, verbatim as `EventBody::Agent` under a name derived from
its variant (`fabro_types::coding_event_name`), and the store folds those
events into `StageProjection.agent` with pebble's `SessionProjection`. Do not
add a fabro event that restates a pebble event, and do not add a second fold
of the stream: read `StageProjection.agent`, or the stored pebble event
itself, instead.
Fabro emits an agent event of its own only for a fact pebble cannot know.
Today those are `agent.session.activated`, `agent.session.deactivated`,
`agent.tools.available`, `agent.pair.user_message`,
`agent.pair.system_message`, `agent.interrupt.injected`,
`agent.steer.buffered`, `agent.steer.dropped`, the `agent.acp.*` family, and
`prompt.failover` for a one-shot prompt stage that walks its fallback plan
without pebble. A new fabro agent event needs the same justification: name
the fact pebble does not have.
## Consumer Guidance
When writing Rust consumers (listeners, store projections, CLI progress):
- Match on `event.body` using `EventBody::*` variants. This gives you typed access to event-specific fields. For a pebble event, match `EventBody::Agent(props)` and then `props.coding_event()`.
- For a stage's agent facts (usage, route, MCP servers, skills, todos, subagents, files, failovers, compactions), read `StageProjection.agent` rather than folding the events again.
- Use `event.node_id`, `event.node_label`, `event.session_id`, and `event.parent_session_id` for envelope metadata.
- Only use `event.event_name()` or `event.properties()` for generic/display purposes (logging, forwarding). These involve serialization and should not be used on hot paths.
When writing external JSON consumers (SSE clients, JSONL parsers):
- Match on the `"event"` field for the dot-notation event name.
- Read event-specific data from `"properties"`.
- Read stage/branch identity from `"node_id"` and `"node_label"`.
- Read agent hierarchy from `"session_id"` and `"parent_session_id"`.
Do not rebuild or mutate the `RunEvent` in downstream listeners.
## Bypass And Persistence Guarantees
Any JSONL sink, the run store, and SSE should reflect the same canonical envelope bytes after redaction.
An active workflow treats any run-event sink write failure as fatal. It cancels execution and
attempts to persist `run.failed` through the direct sink path. Persistence-error logs must include
the full source chain so an HTTP status or transport failure remains visible.
`status.json` remains the authoritative completion signal for detached runs. Terminal run status should only be written after all post-run work is finished.
A worker cannot continue past a record it failed to append: the store's
error reaches the engine and fails the run.

View file

@ -122,7 +122,7 @@ Disallowed setup in `fabro-cli/tests/it`:
- writing `run.json` directly
- writing `status.json` directly
- writing `progress.jsonl` directly
- writing Petri records or platform records directly
- writing `conclusion.json` directly
- writing runtime interview files directly
- writing cached workflow files into run dirs directly
@ -175,7 +175,7 @@ Good structured snapshot targets:
- `status.json`
- `inspect` output
- `live.json`
- compacted `progress.jsonl` event sequences
- compacted run stream item sequences
- workflow conclusions and checkpoint summaries
### Keep direct assertions for relational invariants
@ -343,7 +343,7 @@ Before merging a test change, check:
Avoid these patterns in CLI integration tests:
- manually creating fake run directories
- writing `progress.jsonl` lines by hand
- writing run stream items or records by hand
- writing runtime interview files by hand
- writing asset manifests by hand
- scattering the same workflow setup across many files instead of using fixtures

View file

@ -237,7 +237,6 @@
"pages": [
"GET /api/v1/runs/{id}/checkpoint",
"GET /api/v1/runs/{id}/stages",
"GET /api/v1/runs/{id}/stages/{stageId}/events",
"GET /api/v1/runs/{id}/settings"
]
}

View file

@ -45,14 +45,13 @@ Every adapter the integration plan describes lands here.
`/questions/{qid}/answer` is validated against that pending record and
reaches the worker's control interviewer over the control bus (or the
in-process one directly) under the same id, mapped onto Petri's answer.
The adapter still posts the legacy `interview.*` events (through the
worker's run event sink, or the run's database in the server process)
with that id and the projection's stage label, for the readers that
follow the event stream rather than the projection: Slack, `run attach`
and the web app's Q&A renderer. The store derives the `interview.answered`
platform record, with the answering principal, from `interview.completed`.
An expired or cancelled question is completed as `interview.timeout` or
`interview.interrupted`; an auto-approved run answers itself.
The readers that follow the run's stream rather than the projection
(Slack, `run attach`, the web app's Q&A renderer) see the question in
Petri's own progress record (`derived.parsed.kind == "question"`) and
the answer record that closes it. The server records who answered as
the `interview.answered` platform record when it accepts the answer. An
expired or cancelled question ends without an answer and Petri's gate
fails closed; an auto-approved run answers itself.
- `secrets`: Petri's `SecretProvider` over the vault's token entries, so a
`{{ secrets.NAME }}` reference resolves at spawn into a command's
environment and is masked in every record; a sensitive answer registers

View file

@ -5,8 +5,8 @@
//! dispatcher hands this adapter one [`InterviewRequest`] per question, on
//! its own task, with the question's identity (invocation path, execution,
//! firing, attempt, node, occurrence, ask). The adapter surfaces the
//! question to Fabro the way a legacy `human` stage does, waits for the
//! answer the way the legacy worker does, and hands Petri the reply.
//! question to Fabro as a `human` stage's question, waits for the answer
//! through the questions API, and hands Petri the reply.
//!
//! # One identity
//!
@ -26,27 +26,17 @@
//!
//! The projection derives the pending question, its answer and its expiry
//! from Petri's records alone; the run's own record is the source of truth
//! and nothing the adapter posts is folded into it. The adapter still
//! posts the legacy `interview.*` events through a [`QuestionSink`], with
//! Petri's id and the projection's stage label, for the readers that
//! follow the run's event stream rather than its projection:
//! and nothing the adapter posts is folded into it. The readers that follow
//! the run's stream rather than its projection (the server's Slack service,
//! `fabro run attach`, the web app's Q&A renderer) see the question in
//! Petri's own event: the progress record whose `derived.parsed.kind` is
//! `question`, and the answer record that closes it.
//!
//! - the server's Slack service posts a question to the channel on
//! `interview.started` and finishes it on `interview.completed`,
//! `interview.timeout` or `interview.interrupted`;
//! - `fabro run attach` polls the questions API when `interview.started`
//! arrives and stops waiting on the question's closing event;
//! - the web app's human Q&A renderer pairs `interview.started` with its
//! closing event by question id in the stage's event list;
//! - the server clears its record of an accepted answer on the closing event,
//! so the transport can be claimed again.
//!
//! The worker's [`EventSinkQuestions`] appends them over the run event sink
//! the worker already carries lifecycle events on, and the server's
//! in-process path appends them through [`DatabaseQuestions`]. For a Petri
//! run the store derives the `interview.answered` platform record from
//! `interview.completed`: the question's Petri id and the principal that
//! answered, which is the actor the adapter stamps on the event.
//! The adapter reports what happens to each question as a [`QuestionNotice`]
//! to a [`QuestionSink`], when something observes it: a test's board. The
//! server records who answered as the `interview.answered` platform record
//! when it accepts the answer, since that is a Fabro fact Petri's answer
//! record does not carry.
//!
//! # How the answer comes back
//!
@ -69,9 +59,9 @@
//! none) and reports the expiry itself; the dispatcher then fires the
//! adapter's cancel token, as it does when the firing ends without an
//! answer or the run is cancelled. The adapter returns promptly with
//! [`InterviewReply::Cancelled`] and posts `interview.timeout` when the
//! gate reported the expiry, else `interview.interrupted`, so the readers
//! above see the question end. The expiry report is seen by the adapter's
//! [`InterviewReply::Cancelled`] and reports the question as expired when
//! the gate reported the expiry, else as interrupted, so an observer sees
//! the question end. The expiry report is seen by the adapter's
//! own observer ([`FabroInterviewer::observer`]), which the run registers
//! ahead of the dispatcher so the report is noted before the token fires.
//! The dispatcher races the reply against the same token and may drop the
@ -85,10 +75,10 @@
//! # Auto-approval
//!
//! A run whose `[run.execution] approval` is `auto` answers every question
//! at once as the legacy runner's auto-approve interviewer does (`yes`,
//! the first option, or `auto-approved` text), attributed to the engine.
//! The question is still posted and completed, so the run's stream shows
//! what was decided, and the projection closes it on the delivered answer.
//! at once as `--auto-approve` always has (`yes`, the first option, or
//! `auto-approved` text), attributed to the engine. The question is still
//! asked and answered through Petri, so the run's stream shows what was
//! decided, and the projection closes it on the delivered answer.
use std::collections::{BTreeSet, HashMap, HashSet};
use std::sync::{Arc, Mutex, PoisonError};
@ -152,7 +142,7 @@ impl QuestionIdentity {
}
}
/// A question as Fabro shows it: the fields of `interview.started`.
/// A question as Fabro shows it.
#[derive(Clone, Debug, PartialEq)]
pub struct AskedQuestion {
/// Petri's id for the question, the one id Fabro knows it by.
@ -421,8 +411,9 @@ impl Interviewer for FabroInterviewer {
return InterviewReply::Cancelled;
};
// `submission.actor` is who answered, a Fabro fact Petri's answer
// record does not carry: it goes out on `interview.completed`, from
// which the store derives the `interview.answered` platform record.
// record does not carry: the server records it as the
// `interview.answered` platform record when it accepts the answer;
// here it only reaches an observer.
let Some(answer) = petri_answer(&submission.answer, &request.question) else {
outstanding.close_unanswered(&reason_of(&submission.answer.value));
return InterviewReply::Cancelled;
@ -444,8 +435,8 @@ impl Interviewer for FabroInterviewer {
/// A question the adapter is waiting on. When the wait ends without an
/// answer, whether the adapter saw the cancel or the dispatcher dropped
/// the reply future first, the end of the question is posted from here
/// on its own task: `interview.timeout` when the gate reported the
/// expiry, else `interview.interrupted`.
/// on its own task: as expired when the gate reported the expiry, else as
/// interrupted.
struct Outstanding {
sink: Arc<dyn QuestionSink>,
observed: Arc<Observed>,
@ -663,8 +654,8 @@ fn strip_accelerator(label: &str) -> &str {
}
}
/// The answer as `interview.completed` records it. A sensitive text
/// answer is never written out: the dispatcher registers it as a secret.
/// The answer as the notice records it. A sensitive text answer is never
/// written out: the dispatcher registers it as a secret.
fn describe(answer: &Answer, question: &Question) -> String {
if !answer.choices.is_empty() {
return answer.choices.join(", ");