Merge pull request #865 from fabro-sh/agent-events-docs

State the agent event contract and show the agent sidebar in demo mode
This commit is contained in:
Bryan Helmkamp 2026-09-13 10:36:58 -04:00 • committed by GitHub
commit c9820cf0ad
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 376 additions and 83 deletions

View file

@ -175,11 +175,32 @@ Check:
- 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.
- 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.

View file

@ -883,7 +883,36 @@ Emitted when execution loops back to an earlier node.
## Agent events
Most agent activity events are stage-scoped and carry `node_id` (the workflow stage), `node_label`, `stage_id`, `session_id`, and `parent_session_id` in the envelope. Session object lifecycle events are the exception: `agent.session.started` and `agent.session.ended` are not stage-scoped and intentionally omit `node_id`, `node_label`, `stage_id`, and `visit`.
Pebble's `CodingAgentEvent` stream is the agent event contract. Every event
the coding agent publishes for a stage, except streaming deltas, is stored
verbatim as an `EventBody::Agent` under a name derived from its variant
(`agent.message`, `agent.tool.started`, `agent.route.failover`,
`agent.mcp.server.ready`, `todo.created`, and so on; the full list is
`CODING_EVENT_NAMES`). Its `properties` are pebble's own envelope, so
pebble's event types are part of fabro's stored format, and the store folds
the same events into `StageProjection.agent` with pebble's
`SessionProjection`, the one fold of that stream.
Fabro emits an agent event of its own only for a fact pebble cannot know:
- `agent.session.activated` and `agent.session.deactivated`: the stage's
route, controls, permission level, and steering capabilities, as fabro
resolved them.
- `agent.tools.available`: the tool catalog fabro handed the agent.
- `agent.pair.user_message` and `agent.pair.system_message`: pair mode.
- `agent.interrupt.injected`, `agent.steer.buffered`, `agent.steer.dropped`:
run-level steering as it reaches, waits for, or misses a session.
- `agent.acp.started`, `agent.acp.completed`, `agent.acp.cancelled`,
`agent.acp.timed_out`: an external ACP agent process, which pebble does
not run.
- `prompt.failover`: a one-shot prompt stage moving to a fallback route,
which it does without pebble.
Every agent activity event is stage-scoped and carries `node_id` (the
workflow stage), `node_label`, `stage_id`, `session_id`, and
`parent_session_id` in the envelope. Pebble's session lifecycle events
(`agent.session.started`, `agent.session.ended`) are stored with the stage
that ran the session like the rest.
### `agent.session.started`

View file

@ -364,30 +364,33 @@ V2 keeps the current durable family surface broadly intact.
### Agent Durable Events
- `agent.session.started`
- `agent.session.ended`
- `agent.processing.end`
- `agent.input`
- `agent.message`
- `agent.tool.started`
- `agent.tool.completed`
- `agent.error`
- `agent.warning`
- `agent.loop.detected`
- `agent.steering.injected`
- `agent.compaction.started`
- `agent.compaction.completed`
- `agent.llm.started`
- `agent.llm.first_output`
- `agent.llm.retry`
- `agent.sub.spawned`
- `agent.sub.completed`
- `agent.sub.failed`
- `agent.sub.closed`
- `agent.mcp.ready`
- `agent.mcp.failed`
- `agent.mcp.disconnected`
- `agent.failover`
Pebble's events, stored verbatim under the names `fabro_types::CODING_EVENT_NAMES`
lists (every `CodingEvent` variant except the streaming deltas):
- `agent.session.started`, `agent.session.ended`, `agent.processing.end`
- `agent.input`, `agent.message`
- `agent.llm.started`, `agent.llm.first_output`, `agent.llm.retry`
- `agent.tool.started`, `agent.tool.completed`, `agent.tool.process.completed`, `agent.tool.rounds.exhausted`
- `agent.error`, `agent.warning`, `agent.loop.detected`
- `agent.steering.injected`, `agent.round.interrupted`
- `agent.compaction.started`, `agent.compaction.completed`, `agent.compaction.failed`, `agent.compaction.cancelled`
- `agent.route.failover`, `agent.route.failover.stopped`
- `agent.mcp.server.ready`, `agent.mcp.server.failed`, `agent.mcp.server.disconnected`
- `agent.sub.spawned`, `agent.sub.turn.started`, `agent.sub.completed`, `agent.sub.failed`, `agent.sub.closed`
- `agent.memory.loaded`, `agent.skills.discovered`, `agent.skill.activated`
- `todo.created`, `todo.updated`, `todo.deleted`
Fabro's own, for facts pebble cannot know:
- `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`
- `agent.acp.started`, `agent.acp.completed`, `agent.acp.cancelled`, `agent.acp.timed_out`
- `prompt.failover` (a one-shot prompt stage's move to a fallback route)
The former mirrors `agent.mcp.ready`, `agent.mcp.failed`,
`agent.mcp.disconnected`, and `agent.failover` are no longer emitted; runs
recorded with them read them back as generic events.
### Git

View file

@ -449,6 +449,18 @@ pub(crate) async fn get_run_status(
}
}
/// The demo run's projection: every stage with its state, and the agent
/// stage carrying the coding agent's fold of its stored events, so the stage
/// sidebar shows MCP servers, a failover, a subagent, files, skills, and a
/// compaction in demo mode.
pub(crate) async fn get_run_state(
_auth: RequiredUser,
State(_state): State<Arc<AppState>>,
Path(_id): Path<String>,
) -> Response {
(StatusCode::OK, Json(runs::run_state())).into_response()
}
#[derive(Debug, serde::Deserialize)]
pub(crate) struct ResolveRunParams {
selector: String,
@ -1441,10 +1453,18 @@ mod runs {
]
}
/// The agent stage's stored events: what pebble reports for one prompt
/// that finds MCP servers, activates a skill, reads and writes files,
/// delegates to a subagent, moves to a fallback route, and compacts,
/// plus fabro's own `stage.prompt`.
pub(super) fn stage_events() -> Vec<fabro_types::EventEnvelope> {
use fabro_types::run_event::stage::StagePromptProps;
use fabro_types::{AgentEventProps, EventBody, EventEnvelope, RunEvent};
use pebble_coding_agent::events::{CodingAgentEvent, CodingEvent, TokenUsage};
use pebble_coding_agent::events::{
CodingAgentEvent, CodingEvent, CompactionReason, ErrorData, ErrorKind,
FailoverContinuation, InputSource, McpToolSummary, SkillActivationSource, SkillSummary,
TokenUsage,
};
let run_id = demo_run_id(1);
let node_id = "detect-drift";
@ -1476,28 +1496,40 @@ mod runs {
CodingAgentEvent::new("ses_demo_detect_drift", event, ts.into()),
))
};
let message = |text: &str| {
agent(CodingEvent::AssistantMessage {
let subagent = |event: CodingEvent| {
EventBody::Agent(AgentEventProps::new(
node_id,
1,
CodingAgentEvent::new("ses_demo_sub_1", event, ts.into())
.with_parent_session_id("ses_demo_detect_drift"),
))
};
let answer =
|model: &str, text: &str, input: u64, output: u64| CodingEvent::AssistantMessage {
text: text.into(),
model: "claude-opus-4.6".into(),
usage: TokenUsage::default(),
model: model.into(),
usage: TokenUsage {
input,
output,
..TokenUsage::default()
},
cost_usd_micros: None,
cost_source: None,
tool_call_count: 0,
context_window: None,
reasoning: None,
})
};
let tool_started = |tool_call_id: &str, path: &str| {
};
let message = |text: &str| agent(answer("claude-opus-4.6", text, 1_200, 180));
let call_started = |tool: &str, tool_call_id: &str, arguments: serde_json::Value| {
agent(CodingEvent::ToolCallStarted {
tool_name: "read_file".into(),
tool_name: tool.into(),
tool_call_id: tool_call_id.into(),
arguments: serde_json::json!({ "path": path }),
arguments,
})
};
let tool_completed = |tool_call_id: &str, output: &str| {
let call_completed = |tool: &str, tool_call_id: &str, output: &str| {
agent(CodingEvent::ToolCallCompleted {
tool_name: "read_file".into(),
tool_name: tool.into(),
tool_call_id: tool_call_id.into(),
output: serde_json::json!(output),
metadata: pebble_agent::ToolOutputMetadata::default(),
@ -1508,52 +1540,222 @@ mod runs {
output_bytes_omitted: 0,
})
};
let tool_started = |tool_call_id: &str, path: &str| {
call_started(
"read_file",
tool_call_id,
serde_json::json!({ "path": path }),
)
};
let tool_completed =
|tool_call_id: &str, output: &str| call_completed("read_file", tool_call_id, output);
let started = |provider: &str, model: &str| CodingEvent::SessionStarted {
provider: Some(provider.into()),
model: Some(model.into()),
};
vec![
make_envelope(
1,
"evt-detect-drift-1",
EventBody::StagePrompt(StagePromptProps {
visit: 1,
text: "You are a drift detection agent. Compare the production and staging environments and identify any configuration or code drift.".into(),
mode: None,
provider: None,
model: None,
reasoning_effort: None,
speed: None,
}),
let prompt = "You are a drift detection agent. Compare the production and staging environments and identify any configuration or code drift.";
let report = "# Drift report\n\n- redis.max_connections: 200 (production) vs 100 (staging)\n- redis.tls: enabled vs disabled\n- iam.session_duration: 3600s vs 1800s\n";
let events = vec![
EventBody::StagePrompt(StagePromptProps {
visit: 1,
text: prompt.into(),
mode: None,
provider: None,
model: None,
reasoning_effort: None,
speed: None,
}),
agent(started("anthropic", "claude-opus-4.6")),
agent(CodingEvent::McpServerReady {
server: "github".into(),
tools: vec![
McpToolSummary {
name: "mcp__github__list_issues".into(),
original_name: "list_issues".into(),
},
McpToolSummary {
name: "mcp__github__create_issue".into(),
original_name: "create_issue".into(),
},
],
startup_ms: 842,
}),
agent(CodingEvent::McpServerFailed {
server: "atlassian".into(),
error: "auth failed: the API token has expired".into(),
startup_ms: 3,
}),
agent(CodingEvent::SkillsDiscovered {
profile: "anthropic".into(),
source_dirs: vec![".fabro/skills".into()],
skills: vec![
SkillSummary {
name: "drift-triage".into(),
description: "Rank configuration drift by blast radius".into(),
},
SkillSummary {
name: "terraform".into(),
description: "Read and plan Terraform modules".into(),
},
],
skipped: Vec::new(),
}),
agent(CodingEvent::UserInput {
text: prompt.into(),
content: None,
source: InputSource::Prompt,
}),
message(
"I'll start by loading the environment configurations for both production and staging to compare them.",
),
make_envelope(
2,
"evt-detect-drift-2",
message("I'll start by loading the environment configurations for both production and staging to compare them."),
tool_started("toolu_01", "environments/production/config.toml"),
tool_completed(
"toolu_01",
"[redis]\nhost = \"redis-prod.internal\"\nport = 6379",
),
make_envelope(
3,
"evt-detect-drift-3",
tool_started("toolu_01", "environments/production/config.toml"),
tool_started("toolu_02", "environments/staging/config.toml"),
tool_completed(
"toolu_02",
"[redis]\nhost = \"redis-staging.internal\"\nport = 6379",
),
make_envelope(
4,
"evt-detect-drift-4",
tool_completed("toolu_01", "[redis]\nhost = \"redis-prod.internal\"\nport = 6379"),
agent(CodingEvent::SkillActivated {
skill_name: "drift-triage".into(),
source: SkillActivationSource::Tool,
}),
call_started(
"mcp__github__list_issues",
"toolu_03",
serde_json::json!({ "labels": ["drift"] }),
),
make_envelope(
5,
"evt-detect-drift-5",
tool_started("toolu_02", "environments/staging/config.toml"),
call_completed("mcp__github__list_issues", "toolu_03", "[]"),
agent(CodingEvent::SubAgentSpawned {
agent_id: "sub-1".into(),
depth: 1,
task: "Check the IAM session policy in staging".into(),
generation: 1,
}),
subagent(started("anthropic", "claude-opus-4.6")),
subagent(answer(
"claude-opus-4.6",
"Staging sets iam.session_duration to 1800s; production uses 3600s.",
640,
90,
)),
agent(CodingEvent::SubAgentCompleted {
agent_id: "sub-1".into(),
depth: 1,
generation: 1,
success: true,
turns_used: 2,
}),
agent(CodingEvent::RouteFailover {
from: "anthropic/claude-opus-4.6".into(),
to: "openai/gpt-5.4".into(),
attempt: 1,
error: ErrorData::new(ErrorKind::Llm, "rate limited: retry after 30s"),
usage: TokenUsage {
input: 3_600,
output: 540,
..TokenUsage::default()
},
cost_usd_micros: None,
inference_ms: 4_200,
tool_ms: 900,
continuation: FailoverContinuation::ContinueTurn,
}),
agent(CodingEvent::CompactionCompleted {
original_turn_count: 20,
preserved_turn_count: 6,
summary_token_estimate: 500,
tracked_file_count: 2,
reason: CompactionReason::Threshold,
usage: TokenUsage {
input: 2_000,
output: 500,
..TokenUsage::default()
},
cost_usd_micros: None,
}),
call_started(
"write_file",
"toolu_04",
serde_json::json!({ "file_path": "reports/drift.md", "content": report }),
),
make_envelope(
6,
"evt-detect-drift-6",
tool_completed("toolu_02", "[redis]\nhost = \"redis-staging.internal\"\nport = 6379"),
),
make_envelope(
7,
"evt-detect-drift-7",
message("I've detected drift in 3 resources between production and staging:\n\n1. **redis.max_connections** — production has 200, staging has 100\n2. **redis.tls** — enabled in production, disabled in staging\n3. **iam.session_duration** — production uses 3600s, staging uses 1800s"),
),
]
call_completed("write_file", "toolu_04", "wrote reports/drift.md"),
agent(answer(
"gpt-5.4",
"I've detected drift in 3 resources between production and staging:\n\n1. **redis.max_connections** — production has 200, staging has 100\n2. **redis.tls** — enabled in production, disabled in staging\n3. **iam.session_duration** — production uses 3600s, staging uses 1800s\n\nThe report is in `reports/drift.md`.",
1_500,
260,
)),
agent(CodingEvent::ProcessingEnd),
];
events
.into_iter()
.enumerate()
.map(|(index, body)| {
let seq = u32::try_from(index + 1).expect("the demo stream is short");
make_envelope(seq, &format!("evt-detect-drift-{seq}"), body)
})
.collect()
}
/// The demo run's projection: each stage as `stages()` lists it, and the
/// agent stage carrying the coding agent's fold of `stage_events()`.
pub(super) fn run_state() -> fabro_types::RunProjection {
use fabro_types::{
EventBody, Graph, RunProjection, RunProvenance, RunSpec, StageTiming, WorkflowSettings,
first_event_seq,
};
use pebble_coding_agent::projection::SessionProjection;
let created_at = ts("2026-03-06T14:30:00Z");
let spec = RunSpec {
run_id: demo_run_id(1),
settings: WorkflowSettings::default(),
graph: Graph::new("drift-remediation"),
graph_source: Some(super::DEMO_GRAPH_DOT.to_string()),
workflow_slug: Some("implement".to_string()),
workflow_version_id: None,
target: None,
automation: None,
source_directory: Some("/demo/api-server".to_string()),
labels: HashMap::new(),
provenance: RunProvenance {
server: None,
client: None,
subject: DEMO_PRINCIPAL.clone(),
},
manifest_blob: None,
definition_blob: None,
spec_blob: None,
git: None,
fork_source_ref: None,
};
let mut projection = RunProjection::new(
"Detect and fix environment drift".to_string(),
spec,
created_at,
);
for (index, stage) in stages().into_iter().enumerate() {
let seq = u32::try_from(index + 1).expect("the demo has a handful of stages");
let entry =
projection.stage_entry(stage.id.node_id(), stage.id.visit(), first_event_seq(seq));
entry.handler = Some(stage.handler);
entry.state = stage.status;
entry.started_at = stage.started_at;
entry.timing = stage.wall_time_ms.map(StageTiming::wall_only);
}
let mut agent = SessionProjection::new();
for envelope in stage_events() {
if let EventBody::Agent(props) = &envelope.event.body {
agent.apply(&props.event);
}
}
let detect = projection.stage_entry("detect-drift", 1, first_event_seq(1));
detect.agent = Some(agent);
projection
}
pub(super) fn billing() -> RunBilling {

View file

@ -92,7 +92,7 @@ pub(super) fn demo_routes() -> Router<Arc<AppState>> {
.route("/runs/{id}", get(demo::get_run_status))
.route("/runs/{id}/questions", get(demo::get_questions_stub))
.route("/runs/{id}/questions/{qid}/answer", post(demo::answer_stub))
.route("/runs/{id}/state", get(not_implemented))
.route("/runs/{id}/state", get(demo::get_run_state))
.route("/runs/{id}/logs", get(not_implemented))
.route(
"/runs/{id}/events",

View file

@ -34,7 +34,7 @@ async fn demo_stage_events_default_returns_all_fixture_events_with_no_more() {
let body = get_json(&app, "/api/v1/runs/run-1/stages/detect-drift@1/events").await;
let data = body["data"].as_array().expect("data is an array");
assert_eq!(data.len(), 7, "all seven fixture events should be returned");
assert_eq!(data.len(), 24, "every fixture event should be returned");
assert_eq!(body["meta"]["has_more"], false);
}
@ -57,7 +57,7 @@ async fn demo_stage_events_limit_one_signals_has_more() {
async fn demo_stage_events_since_seq_filters_out_earlier_events() {
let app = fabro_server::test_support::build_test_router(test_app_state());
// The fixture seqs are 1..=7. since_seq=4 should skip the first three.
// The fixture seqs are 1..=24. since_seq=4 should skip the first three.
let body = get_json(
&app,
"/api/v1/runs/run-1/stages/detect-drift@1/events?since_seq=4",
@ -65,11 +65,49 @@ async fn demo_stage_events_since_seq_filters_out_earlier_events() {
.await;
let data = body["data"].as_array().expect("data is an array");
assert_eq!(data.len(), 4);
assert_eq!(data.len(), 21);
let seqs: Vec<u64> = data
.iter()
.map(|envelope| envelope["seq"].as_u64().expect("seq is a number"))
.collect();
assert_eq!(seqs, vec![4, 5, 6, 7]);
assert_eq!(seqs, (4..=24).collect::<Vec<u64>>());
assert_eq!(body["meta"]["has_more"], false);
}
#[tokio::test]
async fn demo_run_state_carries_the_agent_stages_fold() {
let app = fabro_server::test_support::build_test_router(test_app_state());
let body = get_json(&app, "/api/v1/runs/run-1/state").await;
let stage = &body["stages"]["detect-drift@1"];
assert_eq!(stage["handler"], "agent");
let agent = &stage["agent"];
assert_eq!(agent["root_session_id"], "ses_demo_detect_drift");
assert_eq!(
agent["route"]["model"], "gpt-5.4",
"the route after the failover"
);
assert_eq!(agent["activity"], "idle");
assert_eq!(
agent["mcp_servers"]["github"]["tools"]
.as_array()
.unwrap()
.len(),
2
);
assert_eq!(agent["mcp_servers"]["github"]["invoked"], true);
assert_eq!(
agent["mcp_servers"]["atlassian"]["error"],
"auth failed: the API token has expired"
);
assert_eq!(agent["skills"]["activated"][0]["name"], "drift-triage");
assert_eq!(agent["subagents"][0]["status"]["status"], "completed");
assert_eq!(agent["failovers"][0]["to"], "openai/gpt-5.4");
assert_eq!(agent["compactions"].as_array().unwrap().len(), 1);
assert_eq!(
agent["files_touched"],
serde_json::json!(["reports/drift.md"])
);
assert!(agent.get("pending_writes").is_none());
assert!(body["stages"]["apply-changes@2"]["agent"].is_null());
}