diff --git a/AGENTS.md b/AGENTS.md index 0f3f512ab..a7f26cb0b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 diff --git a/docs/internal/events-strategy.md b/docs/internal/events-strategy.md index 5d036ed05..16610a2e2 100644 --- a/docs/internal/events-strategy.md +++ b/docs/internal/events-strategy.md @@ -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 (`//` 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=` 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`): `..` such as -`sandbox.stop.completed` or `snapshot.create.started`, `.state`, and -`.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. diff --git a/docs/internal/testing-strategy.md b/docs/internal/testing-strategy.md index b5fa9488a..dbbc65d92 100644 --- a/docs/internal/testing-strategy.md +++ b/docs/internal/testing-strategy.md @@ -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 diff --git a/docs/public/docs.json b/docs/public/docs.json index 601655663..9c92f2fad 100644 --- a/docs/public/docs.json +++ b/docs/public/docs.json @@ -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" ] } diff --git a/lib/components/fabro-petri/README.md b/lib/components/fabro-petri/README.md index 14b23f3bb..ccfd36793 100644 --- a/lib/components/fabro-petri/README.md +++ b/lib/components/fabro-petri/README.md @@ -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 diff --git a/lib/components/fabro-petri/src/interview.rs b/lib/components/fabro-petri/src/interview.rs index 19727f29c..09c2bbfa7 100644 --- a/lib/components/fabro-petri/src/interview.rs +++ b/lib/components/fabro-petri/src/interview.rs @@ -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, observed: Arc, @@ -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(", ");