Delete the agent mirrors and move failover to prompt stages

Pebble's stream is the agent event contract. The run's own agent.mcp.ready,
agent.mcp.failed, and agent.mcp.disconnected events, which mirrored pebble's
McpServer* events, are gone with their props, the sink arms that emitted
them, and their conversion and naming entries; pebble's stored
agent.mcp.server.* events are the only record and feed the stage's fold.

The sink no longer mirrors RouteFailover onto agent.failover either: an
agent stage's moves are pebble's agent.route.failover. The event is now
prompt.failover, emitted only by a one-shot prompt stage that walks its
fallback plan itself, and its props are trimmed to the two routes, the
attempt, and the error; nothing read the rest.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-13 08:21:41 -06:00
parent 232d347d39
commit 7a6eea0399
No known key found for this signature in database
13 changed files with 162 additions and 743 deletions

View file

@ -1414,98 +1414,6 @@ Emitted when a sub-agent is spawned.
| `agent_id` | string | Sub-agent identifier |
| `depth` | number | Nesting depth |
### `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": "...",
"event": "agent.mcp.ready",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"server_name": "github",
"tool_count": 2,
"tools": [
{
"name": "mcp__github__create_issue",
"original_name": "create_issue"
},
{
"name": "mcp__github__list_issues",
"original_name": "list_issues"
}
],
"startup_ms": 842,
"visit": 1
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `server_name` | string | MCP server name |
| `tool_count` | number | Number of tools available |
| `tools` | array | Names-only tool summaries for the ready server, sorted by qualified `name`. Each entry has `name` (Fabro-qualified `mcp__{server}__{tool}` identifier) and `original_name` (server-provided tool name). Descriptions and input schemas are intentionally omitted. The field is omitted from serialized JSON for legacy parity when empty. |
| `startup_ms` | number | Whole milliseconds from the server's launch to its tools being listed. Events written before the field existed read as `0`. |
| `visit` | number | Stage visit count when the server became ready |
### `agent.mcp.failed`
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.mcp.failed",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"server_name": "filesystem",
"error": "Connection refused",
"startup_ms": 4,
"visit": 1
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `server_name` | string | MCP server name |
| `error` | string | Error message |
| `startup_ms` | number | Whole milliseconds from the server's launch to the failure. Events written before the field existed read as `0`. |
| `visit` | number | Stage visit count when the server failed |
### `agent.mcp.disconnected`
An MCP server that was ready lost its connection during the stage. Pebble
publishes the disconnect once per server, from whichever session's tool call
first observed the closed connection, so the event can originate in a
sub-agent. Every later call to that server's tools fails until the session
ends. The stage projection moves the server's status from `ready` to
`disconnected`; its `tool_count` and `invoked` flag are kept.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.mcp.disconnected",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"server_name": "github",
"error": "transport closed",
"visit": 1
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `server_name` | string | MCP server name |
| `error` | string | What closed the connection, as the client observed it |
| `visit` | number | Stage visit count when the disconnect was observed |
### `agent.memory.loaded`
Emitted once per session right after memory discovery, before skills and MCP
@ -1621,48 +1529,53 @@ Emitted whenever a skill is activated in the running session. Sources:
> entirely; slash-skill expansion is reported through `agent.skill.activated`
> with `source == "slash"` instead.
### `agent.failover`
### `prompt.failover`
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.
Emitted by a one-shot prompt stage when it moves to a fallback route. The
prompt stage walks its fallback plan itself, so this is fabro's own event.
An agent stage never emits it: pebble walks the routes and reports each
move as `agent.route.failover`, stored verbatim (below).
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.failover",
"node_id": "code",
"node_label": "code",
"event": "prompt.failover",
"node_id": "summarize",
"node_label": "summarize",
"properties": {
"from_provider": "anthropic",
"from_model": "claude-sonnet-4-20250514",
"to_provider": "openai",
"to_model": "gpt-4o",
"error": "rate limited",
"continuation": "continue_turn"
"attempt": 1,
"error": "rate limited"
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `from_provider` | string | Original provider |
| `from_model` | string | Original model |
| `to_provider` | string | Failover provider |
| `to_model` | string | Failover 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 |
| `from_provider` | string | The provider that failed |
| `from_model` | string | The model that failed |
| `to_provider` | string | The provider the prompt continued on |
| `to_model` | string | The model the prompt continued on |
| `attempt` | number? | How many routes the prompt had moved through, this one included. Absent only on events recorded before it was kept |
| `error` | string | The failure that ended the previous route |
Events recorded before this rename were named `agent.failover` and carried
`original_provider`, `original_model`, `requested_reasoning_effort`,
`effective_reasoning_effort`, and `continuation`; nothing read them.
### `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`. The stage view reads MCP state from
`StageProjection.agent`, which the pebble events feed; the mirrors change
nothing on the stage any more and go next.
`properties` like every other pebble event. They are the only record of an
agent stage's route moves and MCP server outcomes: the stage view reads
both from `StageProjection.agent`, which they feed. Runs recorded before
fabro stored them carry fabro's former mirrors, `agent.failover`,
`agent.mcp.ready`, `agent.mcp.failed`, and `agent.mcp.disconnected`,
which no reader folds any more.
### `agent.route.failover.stopped`

View file

@ -1527,8 +1527,7 @@ mod tests {
use fabro_types::run_event::run::RunFailedProps;
use fabro_types::run_event::{
AgentAcpCancelledProps, AgentAcpCompletedProps, AgentAcpStartedProps,
AgentAcpTimedOutProps, AgentEventProps, AgentMcpDisconnectedProps, AgentMcpFailedProps,
AgentMcpReadyProps, AgentMcpToolSummary, AgentSessionActivatedProps,
AgentAcpTimedOutProps, AgentEventProps, AgentSessionActivatedProps,
AgentSessionDeactivatedProps, CheckpointCompletedProps, InterviewCompletedProps,
InterviewOption, InterviewStartedProps, ParallelBranchCompletedProps,
ParallelBranchStartedProps, RunCompletedProps, RunControlEffectProps, StageCompletedProps,
@ -6859,9 +6858,7 @@ mod tests {
}
/// The stage's embedded fold sees the pebble `McpServer*` events the
/// sink stores, so its MCP view is the whole-session fold's; the
/// `agent.mcp.*` mirrors the sink still emits change nothing on the
/// stage.
/// sink stores, so its MCP view is the whole-session fold's.
#[test]
fn mcp_servers_agree_across_the_two_folds() {
let code = StageId::new("code", 1);
@ -6898,30 +6895,8 @@ mod tests {
let mut run = initialized_projection();
run.apply_event(&stored(1, &code, ready)).unwrap();
run.apply_event(&test_stage_event(
2,
EventBody::AgentMcpReady(AgentMcpReadyProps {
server_name: "github".to_string(),
tool_count: tools.len(),
tools: mirrored_tools(&tools),
startup_ms: 842,
visit: 1,
}),
code.clone(),
))
.unwrap();
run.apply_event(&stored(3, &code, call)).unwrap();
run.apply_event(&stored(4, &code, disconnected)).unwrap();
run.apply_event(&test_stage_event(
5,
EventBody::AgentMcpDisconnected(AgentMcpDisconnectedProps {
server_name: "github".to_string(),
error: "transport closed".to_string(),
visit: 1,
}),
code.clone(),
))
.unwrap();
run.apply_event(&stored(2, &code, call)).unwrap();
run.apply_event(&stored(3, &code, disconnected)).unwrap();
let agent = run.stage(&code).unwrap().agent.as_ref().unwrap();
assert_eq!(agent.mcp_servers, projection.mcp_servers);
@ -6970,22 +6945,11 @@ mod tests {
}
}
fn mirrored_tools(tools: &[McpToolSummary]) -> Vec<AgentMcpToolSummary> {
tools
.iter()
.map(|tool| AgentMcpToolSummary {
name: tool.name.clone(),
original_name: tool.original_name.clone(),
})
.collect()
}
/// The stage keeps `usage` and `model` as its own, derived from
/// `stage.agent` under the rule each assertion states; everything
/// else the stage view shows is read from `agent` directly. The
/// stream is what the sink stores for one agent stage: fabro's own
/// `agent.session.activated` and the `agent.mcp.*` mirrors next to
/// pebble's events.
/// `agent.session.activated` next to pebble's events.
#[test]
fn the_stage_view_reads_the_embedded_fold() {
let code = StageId::new("code", 1);
@ -7030,17 +6994,6 @@ mod tests {
startup_ms: 842,
}),
),
test_stage_event(
4,
EventBody::AgentMcpReady(AgentMcpReadyProps {
server_name: "github".to_string(),
tool_count: tools.len(),
tools: mirrored_tools(&tools),
startup_ms: 842,
visit: 1,
}),
code.clone(),
),
stored(
5,
&code,
@ -7050,16 +7003,6 @@ mod tests {
startup_ms: 3,
}),
),
test_stage_event(
6,
EventBody::AgentMcpFailed(AgentMcpFailedProps {
server_name: "broken".to_string(),
error: "could not launch".to_string(),
startup_ms: 3,
visit: 1,
}),
code.clone(),
),
stored(
7,
&code,
@ -7164,15 +7107,6 @@ mod tests {
error: "transport closed".to_string(),
}),
),
test_stage_event(
20,
EventBody::AgentMcpDisconnected(AgentMcpDisconnectedProps {
server_name: "github".to_string(),
error: "transport closed".to_string(),
visit: 1,
}),
code.clone(),
),
stored(21, &code, root(assistant_message(50, 5))),
stored(22, &code, root(CodingEvent::ProcessingEnd)),
];

View file

@ -848,42 +848,6 @@ fn event_body_from_event(event: &Event) -> EventBody {
visit: *visit,
})
}
Event::AgentMcpReady {
visit,
server_name,
tool_count,
tools,
startup_ms,
..
} => EventBody::AgentMcpReady(fabro_types::AgentMcpReadyProps {
server_name: server_name.clone(),
tool_count: *tool_count,
tools: tools.clone(),
startup_ms: *startup_ms,
visit: *visit,
}),
Event::AgentMcpFailed {
visit,
server_name,
error,
startup_ms,
..
} => EventBody::AgentMcpFailed(fabro_types::AgentMcpFailedProps {
server_name: server_name.clone(),
error: error.clone(),
startup_ms: *startup_ms,
visit: *visit,
}),
Event::AgentMcpDisconnected {
visit,
server_name,
error,
..
} => EventBody::AgentMcpDisconnected(fabro_types::AgentMcpDisconnectedProps {
server_name: server_name.clone(),
error: error.clone(),
visit: *visit,
}),
Event::AgentInterruptInjected { visit, .. } => {
EventBody::AgentInterruptInjected(fabro_types::AgentInterruptInjectedProps {
visit: *visit,

View file

@ -653,35 +653,6 @@ pub enum Event {
visit: u32,
session_id: String,
},
/// An MCP server configured for a stage connected and listed its tools.
AgentMcpReady {
node_id: String,
visit: u32,
server_name: String,
tool_count: usize,
tools: Vec<fabro_types::AgentMcpToolSummary>,
/// Whole milliseconds from launch to the tools being listed.
#[serde(default)]
startup_ms: u64,
},
/// An MCP server configured for a stage failed to start or connect.
AgentMcpFailed {
node_id: String,
visit: u32,
server_name: String,
error: String,
/// Whole milliseconds from launch to the failure.
#[serde(default)]
startup_ms: u64,
},
/// An MCP server that was ready lost its connection during the stage;
/// its tools fail until the session ends.
AgentMcpDisconnected {
node_id: String,
visit: u32,
server_name: String,
error: String,
},
/// A run-level interrupt was delivered to a concrete steerable agent
/// session/stage.
AgentInterruptInjected {
@ -1493,18 +1464,13 @@ impl Event {
Self::Failover { stage, props } => {
warn!(
stage,
original_provider = ?props.original_provider,
original_model = ?props.original_model,
attempt = ?props.attempt,
from_provider = %props.from_provider,
from_model = %props.from_model,
to_provider = %props.to_provider,
to_model = %props.to_model,
requested_reasoning_effort = ?props.requested_reasoning_effort,
effective_reasoning_effort = ?props.effective_reasoning_effort,
continuation = ?props.continuation,
error = %props.error,
"LLM provider failover"
"Prompt stage moved to a fallback route"
);
}
Self::CommandStarted {
@ -1561,42 +1527,6 @@ impl Event {
} => {
debug!(node_id, visit, session_id, "Agent session deactivated");
}
Self::AgentMcpReady {
node_id,
visit,
server_name,
tool_count,
startup_ms,
..
} => {
debug!(
node_id,
visit, server_name, tool_count, startup_ms, "MCP server ready"
);
}
Self::AgentMcpFailed {
node_id,
visit,
server_name,
error,
startup_ms,
} => {
warn!(
node_id,
visit, server_name, error, startup_ms, "MCP server failed"
);
}
Self::AgentMcpDisconnected {
node_id,
visit,
server_name,
error,
} => {
warn!(
node_id,
visit, server_name, error, "MCP server disconnected"
);
}
Self::AgentInterruptInjected {
node_id,
visit,

View file

@ -83,15 +83,12 @@ pub fn event_name(event: &Event) -> Cow<'static, str> {
Event::StallWatchdogTimeout { .. } => "watchdog.timeout",
Event::ArtifactCaptured { .. } => "artifact.captured",
Event::SshAccessReady { .. } => "ssh.ready",
Event::Failover { .. } => "agent.failover",
Event::Failover { .. } => "prompt.failover",
Event::CommandStarted { .. } => "command.started",
Event::CommandCompleted { .. } => "command.completed",
Event::AgentSessionActivated { .. } => "agent.session.activated",
Event::AgentToolsAvailable { .. } => "agent.tools.available",
Event::AgentSessionDeactivated { .. } => "agent.session.deactivated",
Event::AgentMcpReady { .. } => "agent.mcp.ready",
Event::AgentMcpFailed { .. } => "agent.mcp.failed",
Event::AgentMcpDisconnected { .. } => "agent.mcp.disconnected",
Event::AgentInterruptInjected { .. } => "agent.interrupt.injected",
Event::AgentPairUserMessage { .. } => "agent.pair.user_message",
Event::AgentPairSystemMessage { .. } => "agent.pair.system_message",
@ -133,13 +130,18 @@ mod tests {
"parallel.branch.started"
);
assert_eq!(
event_name(&Event::AgentMcpDisconnected {
node_id: "code".to_string(),
visit: 1,
server_name: "github".to_string(),
error: "transport closed".to_string(),
event_name(&Event::Failover {
stage: "code".to_string(),
props: fabro_types::FailoverProps {
from_provider: "anthropic".to_string(),
from_model: "claude-fable-5".to_string(),
to_provider: "openai".to_string(),
to_model: "gpt-5.6-sol".to_string(),
attempt: Some(1),
error: "overloaded".to_string(),
},
}),
"agent.mcp.disconnected"
"prompt.failover"
);
assert_eq!(
event_name(&Event::Agent {

View file

@ -132,10 +132,7 @@ fn stored_event_fields_for_variant(event: &Event) -> StoredEventFields {
| Event::AgentAcpCompleted { node_id, .. }
| Event::AgentAcpCancelled { node_id, .. }
| Event::AgentAcpTimedOut { node_id, .. } => node_stored_fields(Some(node_id.clone())),
Event::AgentAcpStarted { node_id, visit, .. }
| Event::AgentMcpReady { node_id, visit, .. }
| Event::AgentMcpFailed { node_id, visit, .. }
| Event::AgentMcpDisconnected { node_id, visit, .. } => {
Event::AgentAcpStarted { node_id, visit, .. } => {
let node_id_str = node_id.clone();
let node_label = default_node_label(Some(&node_id_str), None);
StoredEventFields {

View file

@ -14,7 +14,6 @@ use fabro_types::FailoverProps;
use lithos_llm::catalog::ProviderId;
use lithos_llm::types::ReasoningEffort;
use pebble_coding_agent::FallbackRoute;
use pebble_coding_agent::events::FailoverContinuation;
use super::controls::EffectiveRequestControls;
use crate::event::{Emitter, Event, StageScope};
@ -111,42 +110,24 @@ impl FallbackPlan {
})
.collect()
}
}
/// The `agent.failover` payload for a move from `from` to `to`, both
/// `provider/model` selectors, on this plan.
///
/// `from` may be a route that failed during activation without serving
/// traffic; `error` says why it was abandoned. Consecutive payloads
/// chain: one's `to` is the next one's `from`. `continuation` is how
/// pebble said the new route carried the prompt on; a one-shot stage,
/// which re-sends its request itself, has none to report.
pub(crate) fn failover_props(
&self,
from: &str,
to: &str,
attempt: u32,
error: &str,
continuation: Option<FailoverContinuation>,
) -> FailoverProps {
let (from_provider, from_model) = split_selector(from);
let (to_provider, to_model) = split_selector(to);
let effective_reasoning_effort = std::iter::once(&self.original)
.chain(self.remaining.iter())
.find(|route| route.selector() == to)
.and_then(|route| route.controls.reasoning_effort);
FailoverProps {
original_provider: Some(self.original.target.provider.to_string()),
original_model: Some(self.original.target.model.to_string()),
attempt: Some(attempt),
from_provider,
from_model,
to_provider,
to_model,
requested_reasoning_effort: self.original.controls.reasoning_effort,
effective_reasoning_effort,
error: error.to_string(),
continuation: continuation.map(|continuation| continuation.as_str().to_string()),
}
/// The `prompt.failover` payload for a one-shot stage's move from `from` to
/// `to`, both `provider/model` selectors.
///
/// `from` may be a route that failed during activation without serving
/// traffic; `error` says why it was abandoned. Consecutive payloads chain:
/// one's `to` is the next one's `from`.
pub(crate) fn failover_props(from: &str, to: &str, attempt: u32, error: &str) -> FailoverProps {
let (from_provider, from_model) = split_selector(from);
let (to_provider, to_model) = split_selector(to);
FailoverProps {
from_provider,
from_model,
to_provider,
to_model,
attempt: Some(attempt),
error: error.to_string(),
}
}
@ -271,8 +252,10 @@ pub(crate) fn fallback_plan(
)
}
/// Emit `agent.failover` for the plan's most recent
/// Emit `prompt.failover` for the plan's most recent
/// [`FallbackPlan::advance`], on a one-shot stage that walks the plan itself.
/// An agent stage never emits it: pebble walks the routes and reports each
/// move as `agent.route.failover`.
pub(crate) fn emit_failover(
node: &Node,
emitter: &Emitter,
@ -283,12 +266,11 @@ pub(crate) fn emit_failover(
emitter.emit_scoped(
&Event::Failover {
stage: node.id.clone(),
props: plan.failover_props(
props: failover_props(
&plan.previous().selector(),
&plan.current().selector(),
plan.attempt(),
error,
None,
),
},
stage_scope,
@ -384,61 +366,26 @@ mod tests {
}
#[test]
fn failover_props_carry_the_continuation_pebble_reported() {
let policy =
ModelFallbackPolicy::new(BTreeMap::from([("claude-fable-5".to_string(), vec![
FallbackTarget::new("openai", "gpt-5.6-sol"),
])]));
let (plan, notices) = fallback_plan(
&enabled_fallback_catalog(),
&policy,
"claude-fable-5",
&builtin::anthropic(),
EffectiveRequestControls {
reasoning_effort: Some(ReasoningEffort::Medium),
speed: None,
},
);
assert!(notices.is_empty());
let continued = plan.failover_props(
fn failover_props_name_both_routes_and_the_attempt() {
let props = failover_props(
"anthropic/claude-fable-5",
"openai/gpt-5.6-sol",
1,
"overloaded",
Some(FailoverContinuation::ContinueTurn),
);
assert_eq!(continued.continuation.as_deref(), Some("continue_turn"));
assert_eq!(continued.original_provider.as_deref(), Some("anthropic"));
assert_eq!(continued.original_model.as_deref(), Some("claude-fable-5"));
assert_eq!(continued.attempt, Some(1));
assert_eq!(continued.from_provider, "anthropic");
assert_eq!(continued.from_model, "claude-fable-5");
assert_eq!(continued.to_provider, "openai");
assert_eq!(continued.to_model, "gpt-5.6-sol");
assert_eq!(
continued.requested_reasoning_effort,
Some(ReasoningEffort::Medium)
);
assert_eq!(continued.error, "overloaded");
assert_eq!(props, FailoverProps {
from_provider: "anthropic".to_string(),
from_model: "claude-fable-5".to_string(),
to_provider: "openai".to_string(),
to_model: "gpt-5.6-sol".to_string(),
attempt: Some(1),
error: "overloaded".to_string(),
});
let replayed = plan.failover_props(
"anthropic/claude-fable-5",
"openai/gpt-5.6-sol",
1,
"overloaded",
Some(FailoverContinuation::ReplayPrompt),
);
assert_eq!(replayed.continuation.as_deref(), Some("replay_prompt"));
// A one-shot stage walks the plan itself and reports no continuation.
let one_shot = plan.failover_props(
"anthropic/claude-fable-5",
"openai/gpt-5.6-sol",
1,
"overloaded",
None,
);
assert_eq!(one_shot.continuation, None);
// A selector with no slash is all model.
let bare = failover_props("local-model", "openai/gpt-5.6-sol", 2, "down");
assert_eq!(bare.from_provider, "");
assert_eq!(bare.from_model, "local-model");
assert_eq!(bare.attempt, Some(2));
}
}

View file

@ -6,8 +6,8 @@
//! ends and resumed by the next, which binds its own event scope, hooks, and
//! interviewer. Model failover is pebble's: the stage hands it the resolved
//! fallback routes, pebble keeps the conversation as it stands and asks the
//! next route to continue it, and this module mirrors each move as the run's
//! `agent.failover` event.
//! next route to continue it, and reports each move as its own
//! `agent.route.failover` event, stored like every other.
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex, PoisonError};
@ -24,15 +24,15 @@ use fabro_mcp::pebble::pebble_servers;
use fabro_sandbox::{RunSandbox, SecretRedactor};
use fabro_types::settings::run::RunModelControls;
use fabro_types::{
AgentMcpToolSummary, AgentProfileKind, BilledModelUsage, ModelRef, PermissionLevel,
SessionCapability, StageId, StageTiming, UsdMicros, billing,
AgentProfileKind, BilledModelUsage, ModelRef, PermissionLevel, SessionCapability, StageId,
StageTiming, UsdMicros, billing,
};
use fabro_util::home::Home;
use lithos_llm::catalog::{ModelId, ProviderId};
use lithos_llm::types::{Message as LlmMessage, Role, TokenCounts};
use pebble_agent::ToolMiddleware;
use pebble_coding_agent::environment::Environment;
use pebble_coding_agent::events::{CodingAgentEvent, CodingEvent, EventSink, EventSinkError};
use pebble_coding_agent::events::{CodingAgentEvent, EventSink, EventSinkError};
use pebble_coding_agent::extensions::HumanInputProvider;
use pebble_coding_agent::projection::{DescendantAccount, SessionProjection};
use pebble_coding_agent::state::Message;
@ -165,18 +165,13 @@ fn classify_agent_error(error: pebble_coding_agent::Error) -> AgentErrorDisposit
/// 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, 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.
/// `SessionProjection` rebuilt from the log sees what the live one saw.
/// Pebble's stream is the agent event contract; fabro emits an agent event
/// of its own only for a fact pebble cannot know.
struct WorkflowEventSink {
emitter: Arc<Emitter>,
node_id: String,
scope: StageScope,
/// The stage's resolved plan, for the controls and origin the mirrored
/// failover event names.
plan: FallbackPlan,
/// Pebble's fold of every event this sink recorded: the stage's one
/// account of what its agent and subagents spent, wrote, and ran. The
/// store folds the same events the same way, so the stage's billing at
@ -200,88 +195,6 @@ impl EventSink for WorkflowEventSink {
// Every event, including streaming deltas, resets the run's activity
// watchdog.
self.emitter.touch();
match &event.event {
// The failed route's accounting (`usage`, `cost_usd_micros`,
// `inference_ms`, `tool_ms`) is not mirrored: the stage's totals
// already include it through the prompt report, and no run event
// of fabro's own carries per-route usage yet.
CodingEvent::RouteFailover {
from,
to,
attempt,
error,
usage: _,
cost_usd_micros: _,
inference_ms: _,
tool_ms: _,
continuation,
} => {
self.emitter.emit_scoped(
&Event::Failover {
stage: self.node_id.clone(),
props: self.plan.failover_props(
from,
to,
*attempt,
&error.message,
Some(*continuation),
),
},
&self.scope,
);
}
CodingEvent::McpServerReady {
server,
tools,
startup_ms,
} => {
self.emitter.emit_scoped(
&Event::AgentMcpReady {
node_id: self.node_id.clone(),
visit: self.scope.visit,
server_name: server.clone(),
tool_count: tools.len(),
tools: tools
.iter()
.map(|tool| AgentMcpToolSummary {
name: tool.name.clone(),
original_name: tool.original_name.clone(),
})
.collect(),
startup_ms: *startup_ms,
},
&self.scope,
);
}
CodingEvent::McpServerFailed {
server,
error,
startup_ms,
} => {
self.emitter.emit_scoped(
&Event::AgentMcpFailed {
node_id: self.node_id.clone(),
visit: self.scope.visit,
server_name: server.clone(),
error: error.clone(),
startup_ms: *startup_ms,
},
&self.scope,
);
}
CodingEvent::McpServerDisconnected { server, error } => {
self.emitter.emit_scoped(
&Event::AgentMcpDisconnected {
node_id: self.node_id.clone(),
visit: self.scope.visit,
server_name: server.clone(),
error: error.clone(),
},
&self.scope,
);
}
_ => {}
}
// 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.
@ -721,7 +634,6 @@ impl PebbleBackend {
emitter: Arc::clone(bindings.emitter),
node_id: bindings.node_id.to_string(),
scope: bindings.stage_scope.clone(),
plan: plan.clone(),
projection: Mutex::new(SessionProjection::new()),
});
let event_sink = Arc::clone(&sink) as Arc<dyn EventSink>;

View file

@ -44,7 +44,7 @@ use fabro_workflow::test_support::WorkflowRunner;
use httpmock::Method::POST;
use httpmock::MockServer;
use lithos_llm::catalog::ProviderId;
use pebble_coding_agent::events::{CodingEvent, FailoverStop};
use pebble_coding_agent::events::{CodingEvent, FailoverContinuation, FailoverStop};
use tokio_util::sync::CancellationToken;
const MODEL: &str = "mock-model";
@ -1139,18 +1139,17 @@ async fn an_mcp_tool_is_available_to_the_stage() {
assert_eq!(echoed.calls_async().await, 1, "{:?}", names(&stage.events));
assert_eq!(work_stage(&state).response.as_deref(), Some("Echoed"));
let ready = stage
.events
.lock()
.unwrap()
.iter()
.find_map(|event| match &event.body {
EventBody::AgentMcpReady(props) => Some(props.clone()),
// The server's outcome is pebble's own event, stored like every other.
let ready = coding_events(&stage.events)
.into_iter()
.find_map(|(_, event)| match event {
CodingEvent::McpServerReady { server, tools, .. } => Some((server, tools)),
_ => None,
})
.expect("the MCP server reports ready");
assert_eq!(ready.server_name, "echo");
assert_eq!(ready.tool_count, 1);
assert_eq!(ready.0, "echo");
assert_eq!(ready.1.len(), 1);
assert_eq!(count(&stage.events, "agent.mcp.server.ready"), 1);
let completed = coding_events(&stage.events)
.into_iter()
.find_map(|(_, event)| match event {
@ -1264,29 +1263,35 @@ async fn failover_continues_the_conversation_without_rerunning_tools() {
work_stage(&state).response.as_deref(),
Some("Recovered on backup")
);
let failover = stage
.events
.lock()
.unwrap()
.iter()
.find_map(|event| match &event.body {
EventBody::Failover(props) => Some(props.clone()),
// The move is pebble's own event, stored verbatim; fabro emits no
// failover event of its own for an agent stage.
let failover = coding_events(&stage.events)
.into_iter()
.find_map(|(_, event)| match event {
CodingEvent::RouteFailover {
from,
to,
error,
continuation,
..
} => Some((from, to, error, continuation)),
_ => None,
})
.expect("the failover is emitted");
assert_eq!(failover.from_provider, "primary");
assert_eq!(failover.to_provider, "backup");
assert_eq!(failover.to_model, "backup-model");
.expect("the failover is stored");
assert!(failover.0.starts_with("primary/"), "got {}", failover.0);
assert_eq!(failover.1, "backup/backup-model");
assert!(
failover.error.contains("primary key revoked"),
failover.2.message.contains("primary key revoked"),
"got {}",
failover.error
failover.2.message
);
assert_eq!(
failover.continuation.as_deref(),
Some("continue_turn"),
failover.3,
FailoverContinuation::ContinueTurn,
"the primary committed a tool result, so the backup continued the turn"
);
assert_eq!(count(&stage.events, "agent.route.failover"), 1);
assert_eq!(count(&stage.events, "prompt.failover"), 0);
let tool_completions = coding_events(&stage.events)
.into_iter()
.filter(|(_, event)| matches!(event, CodingEvent::ToolCallCompleted { .. }))
@ -1390,9 +1395,10 @@ async fn an_exhausted_fallback_chain_stores_the_stopped_failover() {
}
);
// The move to the backup is fabro's own event; the stop on the backup
// is pebble's, stored under its derived name after the error it reports.
assert_eq!(count(&stage.events, "agent.failover"), 1);
// The move to the backup and the stop on the backup are both pebble's,
// stored under their derived names; the stop follows the error it
// reports.
assert_eq!(count(&stage.events, "agent.route.failover"), 1);
assert_eq!(count(&stage.events, "agent.route.failover.stopped"), 1);
let stopped_at = position(&stage.events, "agent.route.failover.stopped").unwrap();
assert!(work_stage_event(&stage.events, stopped_at));

View file

@ -131,10 +131,10 @@ pub use run::{
RunServerProvenance, RunSpec,
};
pub use run_event::{
AgentEventProps, AgentMcpToolSummary, AgentToolsAvailableProps, CODING_EVENT_NAMES, EventBody,
FailoverProps, InterviewOption, MetadataSnapshotFailureKind, MetadataSnapshotPhase, RunEvent,
RunNoticeCode, RunNoticeLevel, RunPairEndedReason, RunPairFailedReason, RunRunnableSource,
SessionCapability, coding_event_name, is_coding_event_name, sandbox_driver_event_name,
AgentEventProps, AgentToolsAvailableProps, CODING_EVENT_NAMES, EventBody, FailoverProps,
InterviewOption, MetadataSnapshotFailureKind, MetadataSnapshotPhase, RunEvent, RunNoticeCode,
RunNoticeLevel, RunPairEndedReason, RunPairFailedReason, RunRunnableSource, SessionCapability,
coding_event_name, is_coding_event_name, sandbox_driver_event_name,
};
pub use run_failure::RunFailure;
pub use run_id::{RunId, fixtures};

View file

@ -236,48 +236,6 @@ pub struct AgentSteerDroppedProps {
pub count: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentMcpReadyProps {
pub server_name: String,
pub tool_count: usize,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub tools: Vec<AgentMcpToolSummary>,
/// Whole milliseconds from the server's launch to its tools being
/// listed. Events written before the field existed read as `0`.
#[serde(default)]
pub startup_ms: u64,
pub visit: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentMcpToolSummary {
pub name: String,
pub original_name: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentMcpFailedProps {
pub server_name: String,
pub error: String,
/// Whole milliseconds from the server's launch to the failure. Events
/// written before the field existed read as `0`.
#[serde(default)]
pub startup_ms: u64,
pub visit: u32,
}
/// An MCP server that was ready lost its connection during the stage; every
/// later call to its tools fails until the session ends. Pebble reports the
/// disconnect once per server, from whichever session's tool call first
/// observed the closed connection.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentMcpDisconnectedProps {
pub server_name: String,
/// What closed the connection, as the client observed it.
pub error: String,
pub visit: u32,
}
#[cfg(test)]
mod tests {
use std::time::{Duration, UNIX_EPOCH};

View file

@ -1,4 +1,3 @@
use lithos_llm::types::ReasoningEffort;
use serde::{Deserialize, Serialize};
use super::ExecOutputTail;
@ -248,35 +247,20 @@ pub struct SshAccessReadyProps {
pub ssh_command: String,
}
/// A one-shot prompt stage moved to a fallback route. The stage walks its
/// plan itself, so this is fabro's own event; an agent stage's moves are
/// pebble's `agent.route.failover`, stored verbatim.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct FailoverProps {
/// `original_*` and `attempt` are `Option` only because failover events
/// recorded before model-keyed fallbacks lack them. New events always set
/// them; stored events are immutable, so absence stays a supported input.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub original_provider: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub original_model: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attempt: Option<u32>,
pub from_provider: String,
pub from_model: String,
pub to_provider: String,
pub to_model: String,
pub from_model: String,
pub to_provider: String,
pub to_model: String,
/// How many routes the prompt had moved through, this one included.
/// `None` only on events recorded before it was kept.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub requested_reasoning_effort: Option<ReasoningEffort>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub effective_reasoning_effort: Option<ReasoningEffort>,
pub error: String,
/// How the new route carried the prompt on, as pebble reported it:
/// `replay_prompt` when nothing the prompt committed was in the
/// conversation and the new route was asked the prompt again, or
/// `continue_turn` when the conversation held assistant output or tool
/// results and the new route continued from there. Absent on events
/// written before pebble reported it, and on one-shot prompt stages,
/// which walk the plan themselves and always re-send the request.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub continuation: Option<String>,
pub attempt: Option<u32>,
pub error: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]

View file

@ -223,12 +223,6 @@ pub enum EventBody {
AgentSteerBuffered(AgentSteerBufferedProps),
#[serde(rename = "agent.steer.dropped")]
AgentSteerDropped(AgentSteerDroppedProps),
#[serde(rename = "agent.mcp.ready")]
AgentMcpReady(AgentMcpReadyProps),
#[serde(rename = "agent.mcp.failed")]
AgentMcpFailed(AgentMcpFailedProps),
#[serde(rename = "agent.mcp.disconnected")]
AgentMcpDisconnected(AgentMcpDisconnectedProps),
#[serde(rename = "subgraph.started")]
SubgraphStarted(SubgraphStartedProps),
#[serde(rename = "subgraph.completed")]
@ -270,7 +264,7 @@ pub enum EventBody {
ArtifactCaptured(ArtifactCapturedProps),
#[serde(rename = "ssh.ready")]
SshAccessReady(SshAccessReadyProps),
#[serde(rename = "agent.failover")]
#[serde(rename = "prompt.failover")]
Failover(FailoverProps),
#[serde(rename = "cli.ensure.started")]
CliEnsureStarted(CliEnsureStartedProps),
@ -516,9 +510,6 @@ impl EventBody {
Self::AgentInterruptInjected(_) => "agent.interrupt.injected",
Self::AgentSteerBuffered(_) => "agent.steer.buffered",
Self::AgentSteerDropped(_) => "agent.steer.dropped",
Self::AgentMcpReady(_) => "agent.mcp.ready",
Self::AgentMcpFailed(_) => "agent.mcp.failed",
Self::AgentMcpDisconnected(_) => "agent.mcp.disconnected",
Self::SubgraphStarted(_) => "subgraph.started",
Self::SubgraphCompleted(_) => "subgraph.completed",
Self::SandboxInitializing(_) => "sandbox.initializing",
@ -534,7 +525,7 @@ impl EventBody {
Self::StallWatchdogTimeout(_) => "watchdog.timeout",
Self::ArtifactCaptured(_) => "artifact.captured",
Self::SshAccessReady(_) => "ssh.ready",
Self::Failover(_) => "agent.failover",
Self::Failover(_) => "prompt.failover",
Self::CliEnsureStarted(_) => "cli.ensure.started",
Self::CliEnsureCompleted(_) => "cli.ensure.completed",
Self::CliEnsureFailed(_) => "cli.ensure.failed",
@ -650,9 +641,6 @@ fn is_known_event_name(event: &str) -> bool {
| "agent.interrupt.injected"
| "agent.steer.buffered"
| "agent.steer.dropped"
| "agent.mcp.ready"
| "agent.mcp.failed"
| "agent.mcp.disconnected"
| "subgraph.started"
| "subgraph.completed"
| "sandbox.initializing"
@ -674,7 +662,7 @@ fn is_known_event_name(event: &str) -> bool {
| "watchdog.timeout"
| "artifact.captured"
| "ssh.ready"
| "agent.failover"
| "prompt.failover"
| "cli.ensure.started"
| "cli.ensure.completed"
| "cli.ensure.failed"
@ -1251,12 +1239,12 @@ mod tests {
}
#[test]
fn historical_failover_event_defaults_new_route_context() {
fn historical_prompt_failover_event_defaults_its_attempt() {
let line = json!({
"id": "evt_failover",
"ts": "2026-04-04T12:00:00.000Z",
"run_id": fixtures::RUN_1,
"event": "agent.failover",
"event": "prompt.failover",
"properties": {
"from_provider": "anthropic",
"from_model": "claude-fable-5",
@ -1268,51 +1256,37 @@ mod tests {
let parsed = RunEvent::from_value(line).unwrap();
let EventBody::Failover(props) = parsed.body else {
panic!("expected agent.failover");
panic!("expected prompt.failover");
};
assert_eq!(props.original_provider, None);
assert_eq!(props.original_model, None);
assert_eq!(props.attempt, None);
assert_eq!(props.requested_reasoning_effort, None);
assert_eq!(props.effective_reasoning_effort, None);
assert_eq!(props.continuation, None);
assert_eq!(props.to_model, "gpt-5.6-sol");
}
#[test]
fn failover_event_round_trips_its_continuation() {
fn prompt_failover_event_round_trips() {
let body = EventBody::Failover(FailoverProps {
original_provider: Some("anthropic".to_string()),
original_model: Some("claude-fable-5".to_string()),
attempt: Some(1),
from_provider: "anthropic".to_string(),
from_model: "claude-fable-5".to_string(),
to_provider: "openai".to_string(),
to_model: "gpt-5.6-sol".to_string(),
requested_reasoning_effort: None,
effective_reasoning_effort: None,
error: "overloaded".to_string(),
continuation: Some("continue_turn".to_string()),
from_model: "claude-fable-5".to_string(),
to_provider: "openai".to_string(),
to_model: "gpt-5.6-sol".to_string(),
attempt: Some(1),
error: "overloaded".to_string(),
});
let value = serde_json::to_value(&body).unwrap();
assert_eq!(value["event"], "agent.failover");
assert_eq!(value["properties"]["continuation"], "continue_turn");
assert_eq!(value["event"], "prompt.failover");
assert_eq!(
value["properties"],
json!({
"from_provider": "anthropic",
"from_model": "claude-fable-5",
"to_provider": "openai",
"to_model": "gpt-5.6-sol",
"attempt": 1,
"error": "overloaded"
})
);
let parsed: EventBody = serde_json::from_value(value).unwrap();
assert_eq!(parsed, body);
// A one-shot stage, or an event written before pebble reported the
// continuation, omits the field rather than writing `null`.
let EventBody::Failover(mut props) = body else {
unreachable!()
};
props.continuation = None;
let value = serde_json::to_value(EventBody::Failover(props)).unwrap();
assert!(
value["properties"]
.as_object()
.unwrap()
.get("continuation")
.is_none()
);
}
#[test]
@ -2303,108 +2277,6 @@ mod tests {
}
}
#[test]
fn agent_mcp_ready_serializes_with_tool_summaries() {
let body = EventBody::AgentMcpReady(AgentMcpReadyProps {
server_name: "github".to_string(),
tool_count: 1,
tools: vec![AgentMcpToolSummary {
name: "mcp__github__create_issue".to_string(),
original_name: "create_issue".to_string(),
}],
startup_ms: 0,
visit: 1,
});
let value = serde_json::to_value(&body).unwrap();
assert_eq!(value["event"], "agent.mcp.ready");
assert_eq!(
value["properties"]["tools"][0]["name"],
"mcp__github__create_issue"
);
assert_eq!(
value["properties"]["tools"][0]["original_name"],
"create_issue"
);
}
#[test]
fn agent_mcp_ready_and_failed_carry_startup_ms_and_default_it_when_absent() {
let ready = EventBody::AgentMcpReady(AgentMcpReadyProps {
server_name: "github".to_string(),
tool_count: 0,
tools: Vec::new(),
startup_ms: 842,
visit: 1,
});
let value = serde_json::to_value(&ready).unwrap();
assert_eq!(value["properties"]["startup_ms"], 842);
let failed = EventBody::AgentMcpFailed(AgentMcpFailedProps {
server_name: "filesystem".to_string(),
error: "could not launch `npx`".to_string(),
startup_ms: 4,
visit: 1,
});
let value = serde_json::to_value(&failed).unwrap();
assert_eq!(value["properties"]["startup_ms"], 4);
// Events written before pebble reported startup time.
let legacy: EventBody = serde_json::from_value(json!({
"event": "agent.mcp.failed",
"properties": {
"server_name": "filesystem",
"error": "Connection refused",
"visit": 1
}
}))
.unwrap();
match legacy {
EventBody::AgentMcpFailed(props) => assert_eq!(props.startup_ms, 0),
other => panic!("unexpected body: {other:?}"),
}
}
#[test]
fn agent_mcp_disconnected_round_trips() {
let body = EventBody::AgentMcpDisconnected(AgentMcpDisconnectedProps {
server_name: "github".to_string(),
error: "transport closed".to_string(),
visit: 1,
});
let value = serde_json::to_value(&body).unwrap();
assert_eq!(value["event"], "agent.mcp.disconnected");
assert_eq!(
value["properties"],
json!({
"server_name": "github",
"error": "transport closed",
"visit": 1
})
);
let parsed: EventBody = serde_json::from_value(value).unwrap();
assert_eq!(parsed, body);
}
#[test]
fn agent_mcp_ready_omits_tools_when_empty() {
let body = EventBody::AgentMcpReady(AgentMcpReadyProps {
server_name: "github".to_string(),
tool_count: 0,
tools: Vec::new(),
startup_ms: 0,
visit: 1,
});
let value = serde_json::to_value(&body).unwrap();
assert!(
value["properties"]
.as_object()
.unwrap()
.get("tools")
.is_none(),
"empty tools should be omitted for legacy parity"
);
}
#[test]
fn agent_tools_available_round_trips_without_parameter_schemas() {
let body = EventBody::AgentToolsAvailable(AgentToolsAvailableProps {