From 217a5860b6b5c5eb453ff159d98c2f009d86189a Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 13 Sep 2026 07:25:38 -0600 Subject: [PATCH] Store ProcessingEnd and the mirrored pebble events in the run log Pebble's SessionProjection reads ProcessingEnd to complete a prompt and mark the session idle, so a projection rebuilt from the run's log needs it: one small event per prompt. The four pebble events the sink mirrored onto fabro's own agent.failover and agent.mcp.* are now stored verbatim as well, so the fold sees the route moves and the MCP outcomes; the mirrors stay until every reader is on the projection. Co-Authored-By: Claude Fable 5.1 --- docs/internal/events.md | 28 +++++++++++++++++-- .../fabro-workflow/src/handler/llm/pebble.rs | 23 +++++++-------- 2 files changed, 35 insertions(+), 16 deletions(-) diff --git a/docs/internal/events.md b/docs/internal/events.md index 1f0747002..bcab1a858 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -951,7 +951,12 @@ Object-lifecycle event. `session_id` and `parent_session_id` are envelope fields } ``` -No properties. +No properties. One per prompt, when the agent has nothing more to do for +it. Pebble's `SessionProjection`, embedded in `StageProjection.agent`, reads +it to mark the prompt complete and the session idle, so the projection +rebuilt from the run's log needs it. Runs recorded before fabro stored it +never have it; their `agent.activity` stays `running`, and +`StageProjection.state` is the authority on whether the stage is done. ### `agent.input` @@ -1402,6 +1407,10 @@ Emitted when a sub-agent is spawned. ### `agent.mcp.ready` +Fabro's mirror of pebble's `McpServerReady`. The pebble event itself is +also stored, as `agent.mcp.server.ready`; see the section on stored pebble +events below. + ```json { "id": "...", "ts": "...", "run_id": "...", @@ -1605,7 +1614,10 @@ Emitted whenever a skill is activated in the running session. Sources: ### `agent.failover` -Emitted when the agent fails over to a different LLM provider/model. +Emitted when the agent fails over to a different LLM provider/model. On an +agent stage this is fabro's mirror of pebble's `RouteFailover`, which is +also stored as `agent.route.failover`; a one-shot prompt stage, which walks +the fallback plan without pebble, emits only this event. ```json { @@ -1633,10 +1645,20 @@ Emitted when the agent fails over to a different LLM provider/model. | `error` | string | Error that triggered failover | | `continuation` | string? | How the new route carried the prompt on, as pebble reported it: `replay_prompt` (nothing the prompt committed was in the conversation, so the new route was asked the prompt again) or `continue_turn` (the conversation held assistant output or tool results, so the new route continued from there). Absent on events written before pebble reported it and on one-shot prompt stages, which re-send their request themselves | +### `agent.route.failover`, `agent.mcp.server.ready`, `agent.mcp.server.failed`, `agent.mcp.server.disconnected` + +Pebble's `RouteFailover`, `McpServerReady`, `McpServerFailed`, and +`McpServerDisconnected` events, stored verbatim with pebble's envelope in +`properties` like every other pebble event. Fabro also mirrors each onto +its own `agent.failover`, `agent.mcp.ready`, `agent.mcp.failed`, and +`agent.mcp.disconnected`, which the store folds into `StageProjection`'s +`mcp_servers`; the pebble events feed `StageProjection.agent`. The mirrors +go once every reader is on `agent`. + ### `agent.route.failover.stopped` Pebble's `RouteFailoverStopped` event, stored verbatim like every other -pebble event fabro does not mirror. An agent stage with fallback routes +pebble event. An agent stage with fallback routes publishes it when a model failure ends the prompt on its current route anyway: the failure does not qualify for failover (`reason: "ineligible"`) or every route has been taken (`reason: "exhausted"`). It follows the diff --git a/lib/components/fabro-workflow/src/handler/llm/pebble.rs b/lib/components/fabro-workflow/src/handler/llm/pebble.rs index b21e3e60b..07f4231f0 100644 --- a/lib/components/fabro-workflow/src/handler/llm/pebble.rs +++ b/lib/components/fabro-workflow/src/handler/llm/pebble.rs @@ -163,13 +163,12 @@ fn classify_agent_error(error: pebble_coding_agent::Error) -> AgentErrorDisposit // --- Event sink ----------------------------------------------------------- /// Pebble's durable event sink for one stage: every agent event becomes a -/// run event in the run's log before the agent goes on. A route failover and -/// an MCP server's outcome or disconnect are facts the run already has -/// events for, so those are mirrored onto the run's own `agent.failover`, -/// `agent.mcp.ready`, `agent.mcp.failed`, and `agent.mcp.disconnected` -/// events instead of being stored twice. A failover that stops short, with -/// the chain exhausted or the error ineligible, has no event of fabro's own -/// and is stored as pebble's `agent.route.failover.stopped`. +/// run event in the run's log before the agent goes on, so the stage's +/// `SessionProjection` rebuilt from the log sees what the live one saw. A +/// route failover and an MCP server's outcome or disconnect are also +/// mirrored onto the run's own `agent.failover`, `agent.mcp.ready`, +/// `agent.mcp.failed`, and `agent.mcp.disconnected` events, which the store +/// still folds; those mirrors go once every reader is on the projection. struct WorkflowEventSink { emitter: Arc, node_id: String, @@ -214,7 +213,6 @@ impl EventSink for WorkflowEventSink { }, &self.scope, ); - return Ok(()); } CodingEvent::McpServerReady { server, @@ -238,7 +236,6 @@ impl EventSink for WorkflowEventSink { }, &self.scope, ); - return Ok(()); } CodingEvent::McpServerFailed { server, @@ -255,7 +252,6 @@ impl EventSink for WorkflowEventSink { }, &self.scope, ); - return Ok(()); } CodingEvent::McpServerDisconnected { server, error } => { self.emitter.emit_scoped( @@ -267,12 +263,13 @@ impl EventSink for WorkflowEventSink { }, &self.scope, ); - return Ok(()); } _ => {} } - // Deltas and the prompt's own durability barrier are not run history. - if event.event.is_streaming_noise() || matches!(event.event, CodingEvent::ProcessingEnd) { + // Streaming deltas are not run history. `ProcessingEnd` is: pebble's + // `SessionProjection` reads it to complete the prompt and mark the + // session idle, so a projection rebuilt from the run's log needs it. + if event.event.is_streaming_noise() { return Ok(()); } self.emitter