Merge pull request #634 from fabro-sh/feat/inference-observability

feat(events): make inference in-flight state observable
This commit is contained in:
Bryan Helmkamp 2026-07-24 23:02:29 -04:00 • committed by GitHub
commit ff47b40e3c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
52 changed files with 1926 additions and 243 deletions

View file

@ -0,0 +1,107 @@
import { afterEach, beforeEach, describe, expect, test } from "bun:test";
import TestRenderer, { act } from "react-test-renderer";
import { LlmOutputKind } from "@qltysh/fabro-api-client";
import type { StageInferenceProjection } from "@qltysh/fabro-api-client";
import { StageInferenceIndicator } from "./stage-inference-indicator";
const OPENED_AT = new Date(Date.now() - 12_000).toISOString();
function makeInference(
overrides: Partial<StageInferenceProjection> = {},
): StageInferenceProjection {
return {
session_id: "ses_root",
started_at: OPENED_AT,
requested_model: {
provider: "anthropic",
model_id: "claude-fable-5",
},
retries: 0,
...overrides,
};
}
function render(
inference: StageInferenceProjection | null | undefined,
settled = false,
): string {
let renderer!: TestRenderer.ReactTestRenderer;
act(() => {
renderer = TestRenderer.create(
<StageInferenceIndicator inference={inference} settled={settled} />,
);
});
const output = JSON.stringify(renderer.toJSON());
act(() => renderer.unmount());
return output;
}
describe("StageInferenceIndicator", () => {
const actGlobal = globalThis as {
IS_REACT_ACT_ENVIRONMENT?: boolean;
};
const previousActEnvironment = actGlobal.IS_REACT_ACT_ENVIRONMENT;
beforeEach(() => {
actGlobal.IS_REACT_ACT_ENVIRONMENT = true;
});
afterEach(() => {
if (previousActEnvironment === undefined) {
delete actGlobal.IS_REACT_ACT_ENVIRONMENT;
} else {
actGlobal.IS_REACT_ACT_ENVIRONMENT = previousActEnvironment;
}
});
test("renders nothing without an open bracket", () => {
expect(render(undefined)).toBe("null");
expect(render(null)).toBe("null");
});
test("names the requested model while nothing has come back", () => {
const output = render(makeInference());
expect(output).toContain("Model request");
expect(output).toContain("waiting on claude-fable-5");
expect(output).toContain('"aria-live":"polite"');
expect(output).toContain('"aria-hidden":"true"');
// No completion estimate exists, so none may be shown.
expect(output).not.toContain("%");
});
test("reports the observed first-output kind", () => {
expect(
render(makeInference({ first_output_kind: LlmOutputKind.REASONING })),
).toContain("reasoning");
expect(
render(makeInference({ first_output_kind: LlmOutputKind.TEXT })),
).toContain("writing");
expect(
render(makeInference({ first_output_kind: LlmOutputKind.TOOL_CALL })),
).toContain("calling tools");
});
test("never says thinking for non-reasoning output", () => {
for (const kind of [LlmOutputKind.TEXT, LlmOutputKind.TOOL_CALL]) {
expect(render(makeInference({ first_output_kind: kind }))).not.toContain(
"thinking",
);
}
});
test("counts retries without presenting them as failure", () => {
const output = render(makeInference({ retries: 2 }));
expect(output).toContain("retry 2");
expect(output).not.toContain("failed");
});
test("goes static once the run can no longer advance the bracket", () => {
const output = render(makeInference(), true);
// An open bracket on a settled run means we never learned how the request
// ended — animating it would claim work that may not be happening.
expect(output).toContain("never completed");
expect(output).not.toContain("animate-pulse");
});
});

View file

@ -0,0 +1,87 @@
import { LlmOutputKind } from "@qltysh/fabro-api-client";
import type { StageInferenceProjection } from "@qltysh/fabro-api-client";
import { Tooltip } from "./ui";
import { formatAbsoluteTs, formatDurationSecs } from "../lib/format";
import { elapsedSecsSince, useTickingNow } from "../lib/time";
export interface StageInferenceIndicatorProps {
/** Open inference bracket from the stage projection, if there is one. */
inference: StageInferenceProjection | null | undefined;
/**
* The run can no longer make progress on this bracket: it reached a terminal
* status, or the stall watchdog fired. An open bracket then means *we never
* learned how the request ended*, not *it is still working*, so the readout
* goes static.
*/
settled: boolean;
}
const ACTIVITY_LABEL: Record<LlmOutputKind, string> = {
[LlmOutputKind.REASONING]: "reasoning",
[LlmOutputKind.TEXT]: "writing",
[LlmOutputKind.TOOL_CALL]: "calling tools",
};
/**
* Live readout for an open model request.
*
* Says only what the event log proves. There is no progress bar, percentage,
* or ETA, because no completion estimate exists; the elapsed clock counts
* since the request opened rather than claiming the model is still working;
* and "reasoning" appears only when the provider actually sent reasoning
* output, never as a guess filling a gap in the log.
*/
export function StageInferenceIndicator({
inference,
settled,
}: StageInferenceIndicatorProps) {
// Ticking is what distinguishes "we are still hearing from this request"
// from "this is a record of one that never closed", so it stops the moment
// the bracket can no longer advance.
const now = useTickingNow(Boolean(inference) && !settled);
if (!inference) return null;
if (settled) {
return (
<p className="pb-2 text-xs text-fg-muted">
Model request opened {formatAbsoluteTs(inference.started_at)}, never
completed
</p>
);
}
const elapsedSecs = elapsedSecsSince(inference.started_at, now);
const activity = inference.first_output_kind
? ACTIVITY_LABEL[inference.first_output_kind]
: `waiting on ${inference.requested_model.model_id}`;
const statusParts = ["Model request", activity];
// A retry that later succeeds is normal, so this is a count, not a failure.
if (inference.retries > 0) {
statusParts.push(`retry ${inference.retries}`);
}
return (
<p className="pb-2 text-xs text-fg-muted">
<Tooltip
label={`Model request opened ${formatAbsoluteTs(inference.started_at)}`}
>
<span className="inline-flex items-center gap-1.5">
<span
className="size-1.5 animate-pulse rounded-full bg-teal-500"
aria-hidden="true"
/>
<span aria-live="polite">{statusParts.join(" · ")}</span>
{elapsedSecs !== null && (
<span aria-hidden="true">
{" · "}
{formatDurationSecs(elapsedSecs)}
</span>
)}
</span>
</Tooltip>
</p>
);
}

View file

@ -87,7 +87,6 @@ describe("queryKeys", () => {
test("agent activity events invalidate per-stage resources", () => {
for (const event of [
"stage.prompt",
"agent.message",
"agent.tool.started",
"agent.tool.completed",
"command.started",
@ -98,9 +97,16 @@ describe("queryKeys", () => {
queryKeys.runs.stageContextWindow("run-1", "stage-1"),
]);
}
expect(queryKeysForRunEvent("run-1", "agent.message", "stage-1")).toEqual([
queryKeys.runs.state("run-1"),
queryKeys.runs.stageEvents("run-1", "stage-1"),
queryKeys.runs.stageContextWindow("run-1", "stage-1"),
]);
});
test("agent activity events without a node_id invalidate nothing", () => {
expect(queryKeysForRunEvent("run-1", "agent.message")).toEqual([]);
test("agent message without a node_id still invalidates projected state", () => {
expect(queryKeysForRunEvent("run-1", "agent.message")).toEqual([
queryKeys.runs.state("run-1"),
]);
});
});

View file

@ -20,6 +20,7 @@ import {
cancelRun,
deleteRuns,
isTerminalCancelledRun,
isTerminalRunStatus,
isCancellationPending,
isCancellationPendingState,
mapError,
@ -424,6 +425,8 @@ describe("run lifecycle actions", () => {
expect(canArchive("failed")).toBe(true);
expect(canArchive("dead")).toBe(true);
expect(canArchive("archived")).toBe(false);
expect(isTerminalRunStatus("succeeded")).toBe(true);
expect(isTerminalRunStatus("running")).toBe(false);
expect(canUnarchive("archived")).toBe(true);
expect(canUnarchive("failed")).toBe(false);

View file

@ -41,7 +41,7 @@ const CANCELABLE_STATUSES = new Set<RunStatus>([
"blocked",
]);
const ARCHIVABLE_STATUSES = new Set<RunStatus>([
const TERMINAL_RUN_STATUSES = new Set<RunStatus>([
"succeeded",
"failed",
"dead",
@ -143,7 +143,7 @@ export function canApprove(run: Run | null | undefined): boolean {
}
export function canArchive(status: string | null | undefined): boolean {
return !!status && ARCHIVABLE_STATUSES.has(status as RunStatus);
return isTerminalRunStatus(status);
}
export function canUnarchive(status: string | null | undefined): boolean {
@ -152,8 +152,13 @@ export function canUnarchive(status: string | null | undefined): boolean {
export function canRetry(run: Pick<Run, "lifecycle"> | null | undefined): boolean {
if (!run || run.lifecycle.archived) return false;
const status = run.lifecycle.status;
return status.kind === "succeeded" || status.kind === "failed" || status.kind === "dead";
return isTerminalRunStatus(run.lifecycle.status.kind);
}
export function isTerminalRunStatus(
status: string | null | undefined,
): boolean {
return !!status && TERMINAL_RUN_STATUSES.has(status as RunStatus);
}
export function canDelete(status: string | null | undefined): boolean {

View file

@ -125,6 +125,36 @@ describe("queryKeysForRunEvent", () => {
queryKeys.runs.stageEvents("run-1", "code@1"),
]);
});
test("every inference projection transition invalidates live run state", () => {
for (const event of [
"agent.llm.started",
"agent.llm.first_output",
"agent.llm.retry",
"agent.error",
]) {
expect(queryKeysForRunEvent("run-1", event, "code@1")).toEqual([
queryKeys.runs.state("run-1"),
queryKeys.runs.stageEvents("run-1", "code@1"),
]);
}
expect(
queryKeysForRunEvent("run-1", "agent.message", "code@1"),
).toEqual([
queryKeys.runs.state("run-1"),
queryKeys.runs.stageEvents("run-1", "code@1"),
queryKeys.runs.stageContextWindow("run-1", "code@1"),
]);
expect(queryKeysForRunEvent("run-1", "agent.session.ended")).toEqual([
queryKeys.runs.state("run-1"),
]);
});
test("watchdog timeout refreshes the stage events that settle inference", () => {
expect(
queryKeysForRunEvent("run-1", "watchdog.timeout", "code@1"),
).toEqual([queryKeys.runs.stageEvents("run-1", "code@1")]);
});
});
describe("subscribeToRunEvents", () => {

View file

@ -110,6 +110,14 @@ const AGENT_CONTROL_STATE_EVENTS = new Set([
"agent.steering.injected",
"agent.session.deactivated",
]);
const INFERENCE_EVENTS = new Set([
"agent.llm.started",
"agent.llm.first_output",
"agent.llm.retry",
"agent.message",
"agent.error",
"agent.session.ended",
]);
// Todo / task mutation events refresh `getRunState` consumers (so per-stage
// todo projections update live) and the run events list.
const TODO_EVENTS = new Set([
@ -183,6 +191,21 @@ export function queryKeysForRunEvent(
return keys;
}
if (INFERENCE_EVENTS.has(event)) {
const keys: Key[] = [queryKeys.runs.state(runId)];
if (stageId) {
keys.push(queryKeys.runs.stageEvents(runId, stageId));
if (event === "agent.message") {
keys.push(queryKeys.runs.stageContextWindow(runId, stageId));
}
}
return keys;
}
if (event === "watchdog.timeout") {
return stageId ? [queryKeys.runs.stageEvents(runId, stageId)] : [];
}
if (STAGE_ACTIVITY_EVENTS.has(event)) {
return stageId
? [

View file

@ -884,6 +884,19 @@ describe("buildChatItems", () => {
});
describe("buildStageActivity pending tools", () => {
test("records watchdog settlement only for the selected stage", () => {
const events: EventEnvelope[] = [
envelope(1, {
event: "watchdog.timeout",
stage_id: "plan@1",
node_id: "plan",
}),
];
expect(buildStageActivity(events, "plan@1").watchdogTimedOut).toBe(true);
expect(buildStageActivity(events, "code@1").watchdogTimedOut).toBe(false);
});
test("returns started-but-not-completed calls for the stage", () => {
const events: EventEnvelope[] = [
envelope(1, {

View file

@ -31,6 +31,7 @@ import type {
ThreadDnaSelection,
} from "../components/event-debug";
import { StageContext } from "../components/stage-context";
import { StageInferenceIndicator } from "../components/stage-inference-indicator";
import { StageInsightsSidebar } from "../components/stage-insights-sidebar";
import { StageSidebar } from "../components/stage-sidebar";
import type { Stage } from "../components/stage-sidebar";
@ -72,6 +73,7 @@ import {
useRunStages,
useRunState,
} from "../lib/queries";
import { isTerminalRunStatus } from "../lib/run-actions";
import {
STAGE_ACTIVITY_EVENT_TYPES,
type StageActivityEventType,
@ -84,6 +86,7 @@ import { getNumber, getString, type UnknownRecord } from "../lib/unknown";
import type {
EventEnvelope,
StageHandler,
StageInferenceProjection,
StageModelUsage,
} from "@qltysh/fabro-api-client";
@ -268,6 +271,7 @@ export interface PendingToolCall {
interface StageActivity {
turns: TurnType[];
pendingTools: PendingToolCall[];
watchdogTimedOut: boolean;
}
interface PendingCommand {
@ -283,11 +287,18 @@ export function buildStageActivity(
const pendingTools = new Map<string, PendingTool>();
let pendingCommand: PendingCommand | undefined;
let sawAssistantMessage = false;
let watchdogTimedOut = false;
for (const e of events) {
const eventName = e.event;
if (activityEventStageId(e) !== stageId) {
continue;
}
if (eventName === "watchdog.timeout") {
watchdogTimedOut = true;
continue;
}
if (
activityEventStageId(e) !== stageId ||
!eventName ||
!STAGE_ACTIVITY_EVENT_SET.has(eventName)
) {
@ -441,6 +452,7 @@ export function buildStageActivity(
return {
turns,
watchdogTimedOut,
pendingTools: Array.from(pendingTools, ([toolCallId, tool]) => ({
toolCallId,
toolName: tool.toolName,
@ -2063,6 +2075,8 @@ function RunStageActivityStage({
selectedStage,
stages,
runStart,
inference,
runSettled,
tab,
selectedKinds,
selectedDebugCategories,
@ -2076,6 +2090,8 @@ function RunStageActivityStage({
selectedStage: Stage;
stages: Stage[];
runStart: string | undefined;
inference: StageInferenceProjection | null | undefined;
runSettled: boolean;
tab: EventsTab;
selectedKinds: EventKind[];
selectedDebugCategories: DebugCategory[];
@ -2087,11 +2103,15 @@ function RunStageActivityStage({
}) {
const selectedStageId = selectedStage.id;
const stageEventsQuery = useRunStageEvents(runId, selectedStageId);
// An open bracket on a run that can no longer advance means we never learned
// how the request ended, not that it is still working. The watchdog stays
// the authority on "actually stuck", so its timeout settles the readout too.
const activity = useMemo(
() => buildStageActivity(stageEventsQuery.data ?? [], selectedStageId),
[stageEventsQuery.data, selectedStageId],
);
const { turns } = activity;
const inferenceSettled = runSettled || activity.watchdogTimedOut;
const renderer: StageRenderer = selectStageRenderer(selectedStage.handler);
const debugEvents = useMemo<EventEnvelope[]>(() => {
return (stageEventsQuery.data ?? []).filter(
@ -2239,6 +2259,11 @@ function RunStageActivityStage({
</Link>
</p>
)}
<StageInferenceIndicator
inference={inference}
settled={inferenceSettled}
/>
<EventsToolbar
tab={effectiveTab}
renderer={renderer}
@ -2336,11 +2361,15 @@ function RunStageActivity({
selectedStage,
stages,
runStart,
inference,
runSettled,
}: {
runId: string;
selectedStage: Stage;
stages: Stage[];
runStart: string | undefined;
inference: StageInferenceProjection | null | undefined;
runSettled: boolean;
}) {
const [activityState, dispatchActivity] = useReducer(
stageActivityReducer,
@ -2356,6 +2385,8 @@ function RunStageActivity({
selectedStage={selectedStage}
stages={stages}
runStart={runStart}
inference={inference}
runSettled={runSettled}
tab={tab}
selectedKinds={selectedKinds}
selectedDebugCategories={selectedDebugCategories}
@ -2407,6 +2438,8 @@ export default function RunStages() {
isAgentStage && selectedStageId
? runStateQuery.data?.stages[selectedStageId]
: undefined;
const runStatusKind = runQuery.data?.lifecycle.status.kind;
const runSettled = isTerminalRunStatus(runStatusKind);
if (!id || !selectedStage) {
return (
@ -2458,6 +2491,8 @@ export default function RunStages() {
selectedStage={selectedStage}
stages={stages}
runStart={runStart}
inference={stageProjection?.inference}
runSettled={runSettled}
/>
</div>
);

View file

@ -968,43 +968,62 @@ No properties.
|----------|------|-------------|
| `text` | string | User input text |
### `agent.output.start`
### `agent.llm.started`
Signals the beginning of assistant text output.
An inference request is about to be dispatched for this round. Emitted once
per round, after the request is built and compaction has run, immediately
before the stream is opened.
`requested_model` is the canonical requested target, including an optional
speed tier. Failover can re-target mid-stage, so `agent.message` remains
authoritative for what actually answered. No usage or cost fields: neither
exists yet at this point.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.output.start",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {}
}
```
No properties.
### `agent.output.replace`
Replaces the current in-progress assistant output buffers.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.output.replace",
"event": "agent.llm.started",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"text": "I'll fix the login bug by...",
"reasoning": "The user wants..."
"requested_model": {
"provider": "anthropic",
"model_id": "claude-fable-5",
"speed": "fast"
},
"visit": 1
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `text` | string | Replacement assistant text |
| `reasoning` | string? | Replacement reasoning text |
| `requested_model` | object | Requested provider, model ID, and optional speed tier |
| `visit` | number | Graph visit |
### `agent.llm.first_output`
The provider produced its first output for the current attempt. Edge-triggered
once per stream attempt; the latch re-arms when a broken or finish-less stream
replays the turn.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.llm.first_output",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"kind": "reasoning",
"visit": 1
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `kind` | string | `reasoning`, `text`, or `tool_call` — observed, not inferred |
| `visit` | number | Graph visit |
### `agent.message`
@ -1047,46 +1066,6 @@ Emitted when the assistant produces a complete message.
| `usage.raw` | object? | Raw provider-specific usage |
| `tool_call_count` | number | Number of tool calls in this turn |
### `agent.text.delta`
Streaming text chunk from the assistant.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.text.delta",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"delta": "I'll start by reading"
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `delta` | string | Text chunk |
### `agent.reasoning.delta`
Streaming reasoning/thinking chunk from the assistant.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.reasoning.delta",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"delta": "The user needs me to..."
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `delta` | string | Reasoning text chunk |
### `agent.tool.started`
Emitted when the agent begins a tool call.
@ -1111,26 +1090,6 @@ Emitted when the agent begins a tool call.
| `tool_call_id` | string | Unique tool call id |
| `arguments` | object | Tool call arguments |
### `agent.tool.output.delta`
Streaming tool output chunk.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.tool.output.delta",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"delta": "fn login(user: &str)..."
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `delta` | string | Output text chunk |
### `agent.tool.completed`
Emitted when a tool call finishes.
@ -1254,24 +1213,6 @@ Emitted when the agent detects a tool-use loop.
No properties.
### `agent.skill.expanded`
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.skill.expanded",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"properties": {
"skill_name": "read_file"
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `skill_name` | string | Expanded skill name |
### `agent.steering.injected`
```json
@ -1336,7 +1277,9 @@ No properties.
### `agent.llm.retry`
Emitted when an LLM API call is retried.
Emitted when an attempt fails to open **or sustain** a stream and the turn is
replayed. The finish-less-stream case carries a synthetic `Stream` error and a
zero delay: the turn restarts even though no error was reported.
```json
{
@ -1349,6 +1292,7 @@ Emitted when an LLM API call is retried.
"model": "claude-sonnet-4-20250514",
"attempt": 2,
"delay_secs": 1.5,
"phase": "open",
"error": { ... }
}
}
@ -1358,8 +1302,9 @@ Emitted when an LLM API call is retried.
|----------|------|-------------|
| `provider` | string | LLM provider name |
| `model` | string | Model identifier |
| `attempt` | number | Retry attempt number |
| `attempt` | number | Retry attempt number, 0-based within the loop named by `phase` |
| `delay_secs` | number | Delay before retry in seconds |
| `phase` | string? | `open` (stream failed to open) or `consume` (stream broke or ended without a finish event). Absent on events stored before the discriminator existed |
| `error` | object | SdkError (serialized) |
### `agent.sub.spawned`
@ -1615,10 +1560,10 @@ Emitted whenever a skill is activated in the running session. Sources:
| `source` | string | `"slash"` for `/skill-name` expansion, `"tool"` for `use_skill` activations |
| `visit` | number | Stage visit count |
> `agent.skill.expanded` is no longer surfaced as a durable run event. The
> internal `AgentEvent::SkillExpanded` variant remains classified as streaming
> noise and is not persisted; slash-skill expansion is reported through
> `agent.skill.activated` with `source == "slash"` instead.
> `agent.skill.expanded` does not exist. The `AgentEvent::SkillExpanded`
> variant this note once described has since been removed from the code
> entirely; slash-skill expansion is reported through `agent.skill.activated`
> with `source == "slash"` instead.
### `agent.failover`
@ -1648,6 +1593,25 @@ Emitted when the agent fails over to a different LLM provider/model.
| `to_model` | string | Failover model |
| `error` | string | Error that triggered failover |
### Agent events that are never serialized
`AgentEvent` also has variants that exist only on the agent session's
in-process broadcast channel. `is_streaming_noise()` filters them out before
the workflow emitter builds a `RunEvent`, so they never reach the run store,
SSE, `fabro events`, or a JSONL sink — they have no envelope, and no external
consumer can observe them:
- `AssistantOutputReplace` — clears in-progress output buffers when a turn is
replayed
- `TextDelta`, `ReasoningDelta` — streaming assistant chunks
- `ToolCallOutputDelta` — streaming tool output chunks
They were previously documented here as though they were durable events, with
full envelope examples. If any of them ever needs to be durable, it belongs in
a separate transient stream rather than the canonical persisted contract —
long autonomous runs would generate orders of magnitude more delta traffic
than the interactive sessions surface handles.
---
## Subgraph events
@ -2222,14 +2186,14 @@ These legacy events may appear in older run logs. Current CLI backend runs do no
|----------|------|-------------|
| `error` | string | Error message |
## Asset events
## Artifact events
### `asset.captured`
### `artifact.captured`
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "asset.captured",
"event": "artifact.captured",
"node_id": "code",
"node_label": "code",
"properties": {

View file

@ -377,6 +377,8 @@ V2 keeps the current durable family surface broadly intact.
- `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`
@ -417,14 +419,14 @@ The current boundary that keeps live token/delta noise out of `RunEvent` should
These stay outside the durable persisted contract:
- `agent.output.start`
- `agent.output.replace`
- `agent.text.delta`
- `agent.reasoning.delta`
- `agent.tool.output.delta`
- `agent.skill.expanded`
`agent.skill.expanded` stays in this non-durable bucket because it is display-oriented expansion metadata, not a durable workflow fact.
(`agent.skill.expanded` was previously listed here. No such event exists — the
`AgentEvent::SkillExpanded` variant was removed, and slash-skill expansion is
reported through the durable `agent.skill.activated` with `source == "slash"`.)
If Fabro needs those for UI, they belong in a separate transient stream, not in the canonical persisted Rust event contract.

View file

@ -283,7 +283,6 @@ Current style is already decent, but V2 should be stricter.
Examples:
- `agent.output.start` -> `message.part.started`
- `agent.text.delta` -> `message.part.delta`
- `agent.tool.output.delta` -> `tool.output.delta`
- `agent.processing.end` -> `turn.completed` or `session.idle`, depending on actual semantics

View file

@ -10719,6 +10719,14 @@ components:
- $ref: "#/components/schemas/StageContextWindowProjection"
- type: "null"
description: Latest content-free context-window snapshot for this agent stage.
inference:
oneOf:
- $ref: "#/components/schemas/StageInferenceProjection"
- type: "null"
description: >
Open inference bracket, if the event log contains one. Present means
a model request was dispatched and no closing event has been seen —
not that the model is computing right now.
agent_control:
$ref: "#/components/schemas/AgentControlState"
description: Whether the agent is executing normally or waiting for steering after an interrupt.
@ -10726,6 +10734,58 @@ components:
$ref: "#/components/schemas/StageState"
description: Lifecycle state of the stage projection.
StageInferenceProjection:
description: >
One open inference bracket: a dispatched LLM request that has not yet
produced a message, error, or interrupt. Carries no usage or cost — none
exists until the turn completes.
type: object
required:
- session_id
- started_at
- requested_model
- retries
properties:
session_id:
type: string
description: >
Agent session that opened the bracket, copied from the event
envelope. Transitions are gated on it so sub-agent rounds cannot
overwrite the root session's bracket.
started_at:
type: string
format: date-time
description: When the request was dispatched.
requested_model:
$ref: "#/components/schemas/BillingModelRef"
description: >
Provider and model the request was sent to. Failover can re-target,
so `StageProjection.model` stays authoritative for what answered.
first_output_at:
type: ["string", "null"]
format: date-time
description: When the provider produced its first output, if it has.
first_output_kind:
oneOf:
- $ref: "#/components/schemas/LlmOutputKind"
- type: "null"
description: Kind of the first output observed for the current attempt.
retries:
type: integer
format: uint32
minimum: 0
description: Attempts that failed and restarted within this bracket.
LlmOutputKind:
description: >
Kind of output a provider produced first for an inference attempt.
Observed, never inferred.
type: string
enum:
- reasoning
- text
- tool_call
SubAgentProjection:
description: Current projected state for one subagent spawned by an agent stage.
type: object

View file

@ -2,7 +2,7 @@ use std::convert::TryFrom;
use chrono::{DateTime, Utc};
use fabro_agent::Error as AgentError;
use fabro_types::{BilledModelUsage, EventBody, RunEvent};
use fabro_types::{BilledModelUsage, EventBody, LlmOutputKind, RunEvent};
use fabro_util::error;
use fabro_workflow::event::RunNoticeLevel;
use serde_json::Value;
@ -137,6 +137,7 @@ pub(super) enum ProgressEvent {
AssistantMessage {
stage_node_id: String,
model: String,
root_session: bool,
},
ToolCallStarted {
stage_node_id: String,
@ -168,6 +169,15 @@ pub(super) enum ProgressEvent {
CompactionFailed {
stage_node_id: String,
error: String,
root_session: bool,
},
LlmRequestStarted {
stage_node_id: String,
model: String,
},
LlmFirstOutput {
stage_node_id: String,
kind: LlmOutputKind,
},
LlmRetry {
stage_node_id: String,
@ -176,6 +186,9 @@ pub(super) enum ProgressEvent {
delay_ms: u64,
error: String,
},
LlmRequestFinished {
stage_node_id: String,
},
SubagentSpawned {
stage_node_id: String,
agent_id: String,
@ -329,6 +342,7 @@ pub(super) fn from_run_event(stored: &RunEvent) -> Option<ProgressEvent> {
EventBody::AgentMessage(props) => Some(ProgressEvent::AssistantMessage {
stage_node_id: node_id,
model: props.model.model_id.to_string(),
root_session: stored.parent_session_id.is_none(),
}),
EventBody::AgentToolStarted(props) => Some(ProgressEvent::ToolCallStarted {
stage_node_id: node_id,
@ -365,13 +379,30 @@ pub(super) fn from_run_event(stored: &RunEvent) -> Option<ProgressEvent> {
preserved_turn_count: props.preserved_turn_count as u64,
tracked_file_count: props.tracked_file_count as u64,
}),
EventBody::AgentError(props) => {
display_compaction_error(&props.error).map(|error| ProgressEvent::CompactionFailed {
EventBody::AgentError(props) => match display_compaction_error(&props.error) {
Some(error) => Some(ProgressEvent::CompactionFailed {
stage_node_id: node_id,
error,
root_session: stored.parent_session_id.is_none(),
}),
None if stored.parent_session_id.is_none() => Some(ProgressEvent::LlmRequestFinished {
stage_node_id: node_id,
}),
None => None,
},
EventBody::AgentLlmStarted(props) if stored.parent_session_id.is_none() => {
Some(ProgressEvent::LlmRequestStarted {
stage_node_id: node_id,
model: props.requested_model.model_id.to_string(),
})
}
EventBody::AgentLlmRetry(props) => {
EventBody::AgentLlmFirstOutput(props) if stored.parent_session_id.is_none() => {
Some(ProgressEvent::LlmFirstOutput {
stage_node_id: node_id,
kind: props.kind,
})
}
EventBody::AgentLlmRetry(props) if stored.parent_session_id.is_none() => {
#[allow(
clippy::cast_possible_truncation,
clippy::cast_sign_loss,
@ -386,6 +417,11 @@ pub(super) fn from_run_event(stored: &RunEvent) -> Option<ProgressEvent> {
error: display_value(&props.error).unwrap_or_else(|| "unknown error".to_string()),
})
}
EventBody::AgentRoundInterrupted(_) if stored.parent_session_id.is_none() => {
Some(ProgressEvent::LlmRequestFinished {
stage_node_id: node_id,
})
}
EventBody::AgentSubSpawned(props) => Some(ProgressEvent::SubagentSpawned {
stage_node_id: node_id,
agent_id: props.agent_id.clone(),

View file

@ -256,7 +256,11 @@ impl ProgressUI {
ProgressEvent::AssistantMessage {
stage_node_id,
model,
root_session,
} => {
if root_session {
self.stage.on_llm_request_finished(&stage_node_id);
}
self.stage
.on_assistant_message(renderer, &stage_node_id, &model);
}
@ -319,10 +323,27 @@ impl ProgressUI {
ProgressEvent::CompactionFailed {
stage_node_id,
error,
root_session,
} => {
if root_session {
self.stage.on_llm_request_finished(&stage_node_id);
}
self.stage
.on_compaction_failed(renderer, &stage_node_id, &error);
}
ProgressEvent::LlmRequestStarted {
stage_node_id,
model,
} => {
self.stage
.on_llm_request_started(renderer, &stage_node_id, &model);
}
ProgressEvent::LlmFirstOutput {
stage_node_id,
kind,
} => {
self.stage.on_llm_first_output(&stage_node_id, kind);
}
ProgressEvent::LlmRetry {
stage_node_id,
model,
@ -339,6 +360,9 @@ impl ProgressUI {
&error,
);
}
ProgressEvent::LlmRequestFinished { stage_node_id } => {
self.stage.on_llm_request_finished(&stage_node_id);
}
ProgressEvent::SubagentSpawned {
stage_node_id,
agent_id,
@ -515,6 +539,17 @@ mod tests {
}
}
fn child_agent_event(stage: &str, event: AgentEvent) -> Event {
Event::Agent {
stage: stage.into(),
visit: 1,
event,
session_id: Some("ses_child".into()),
parent_session_id: Some("ses_root".into()),
tool_call_id: None,
}
}
fn stage_started(node_id: &str, name: &str) -> Event {
Event::StageStarted {
graph_visit: None,
@ -528,9 +563,9 @@ mod tests {
}
}
fn assistant_message(stage: &str, model: &str) -> Event {
agent_event(stage, AgentEvent::AssistantMessage {
text: "done".into(),
fn assistant_event(model: &str, text: &str) -> AgentEvent {
AgentEvent::AssistantMessage {
text: text.into(),
model: ModelRef {
provider: ProviderId::openai(),
model_id: model.into(),
@ -542,6 +577,24 @@ mod tests {
tool_call_count: 0,
context_window: None,
reasoning: None,
}
}
fn assistant_message(stage: &str, model: &str) -> Event {
agent_event(stage, assistant_event(model, "done"))
}
fn child_assistant_message(stage: &str, model: &str) -> Event {
child_agent_event(stage, assistant_event(model, "child done"))
}
fn llm_request_started(stage: &str, model: &str) -> Event {
agent_event(stage, AgentEvent::LlmRequestStarted {
requested_model: ModelRef {
provider: ProviderId::anthropic(),
model_id: model.into(),
speed: None,
},
})
}
@ -673,6 +726,8 @@ mod tests {
}),
);
assert!(ui.stage.active_stages["s1"].compaction_bar.is_some());
emit(&mut ui, llm_request_started("s1", "claude-fable-5"));
assert!(ui.stage.active_stages["s1"].inference_bar.is_some());
emit(
&mut ui,
@ -710,6 +765,7 @@ mod tests {
);
assert!(ui.stage.active_stages["s1"].compaction_bar.is_none());
assert!(ui.stage.active_stages["s1"].inference_bar.is_none());
}
#[test]
@ -728,6 +784,118 @@ mod tests {
insta::assert_snapshot!(rendered(&buffer), @" ✗ compaction failed: generated summary was empty after trimming; refused to replace 14 turns and left history intact");
}
#[test]
fn inference_bracket_sets_updates_and_clears_bar() {
let mut ui = ProgressUI::new(true, false);
emit(&mut ui, stage_started("s1", "Build"));
assert!(ui.stage.active_stages["s1"].inference_bar.is_none());
emit(&mut ui, llm_request_started("s1", "claude-fable-5"));
let message = ui.stage.active_stages["s1"]
.inference_bar
.as_ref()
.expect("bracket should open a live line")
.message();
assert!(
message.contains("waiting on claude-fable-5"),
"expected the requested model, got: {message:?}"
);
emit(
&mut ui,
agent_event("s1", AgentEvent::LlmFirstOutput {
kind: fabro_types::LlmOutputKind::ToolCall,
}),
);
let message = ui.stage.active_stages["s1"]
.inference_bar
.as_ref()
.expect("the line stays open until the round ends")
.message();
assert!(
message.contains("calling tools"),
"expected the observed output kind, got: {message:?}"
);
emit(&mut ui, assistant_message("s1", "claude-fable-5"));
assert!(ui.stage.active_stages["s1"].inference_bar.is_none());
}
#[test]
fn inference_retry_resets_the_live_line_before_verbose_output() {
let mut ui = ProgressUI::new(true, false);
emit(&mut ui, stage_started("s1", "Build"));
emit(&mut ui, llm_request_started("s1", "claude-fable-5"));
emit(
&mut ui,
agent_event("s1", AgentEvent::LlmFirstOutput {
kind: fabro_types::LlmOutputKind::Text,
}),
);
emit(
&mut ui,
agent_event("s1", AgentEvent::LlmRetry {
provider: "anthropic".into(),
model: "claude-fable-5".into(),
attempt: 1,
delay_secs: 0.1,
phase: fabro_types::LlmRetryPhase::Consume,
error: fabro_llm::Error::Configuration {
message: "retry".into(),
source: None,
},
}),
);
let message = ui.stage.active_stages["s1"]
.inference_bar
.as_ref()
.expect("retry keeps the bracket open")
.message();
assert!(message.contains("waiting on claude-fable-5"));
}
#[test]
fn inference_interrupt_clears_the_live_line() {
let mut ui = ProgressUI::new(true, false);
emit(&mut ui, stage_started("s1", "Build"));
emit(&mut ui, llm_request_started("s1", "claude-fable-5"));
emit(
&mut ui,
agent_event("s1", AgentEvent::RoundInterrupted { generation: 1 }),
);
assert!(ui.stage.active_stages["s1"].inference_bar.is_none());
}
#[test]
fn child_session_events_do_not_mutate_the_root_inference_line() {
let mut ui = ProgressUI::new(true, false);
emit(&mut ui, stage_started("s1", "Build"));
emit(&mut ui, llm_request_started("s1", "claude-fable-5"));
emit(
&mut ui,
child_agent_event("s1", AgentEvent::LlmFirstOutput {
kind: fabro_types::LlmOutputKind::ToolCall,
}),
);
emit(&mut ui, child_assistant_message("s1", "child-model"));
let message = ui.stage.active_stages["s1"]
.inference_bar
.as_ref()
.expect("child output must not close the root bracket")
.message();
assert!(message.contains("waiting on claude-fable-5"));
emit(&mut ui, assistant_message("s1", "claude-fable-5"));
assert!(ui.stage.active_stages["s1"].inference_bar.is_none());
}
#[test]
fn handle_json_line_ignores_invalid_json() {
let (mut ui, buffer) = capture_ui(false);
@ -792,6 +960,7 @@ mod tests {
model: "gpt-5-mini".into(),
attempt: 2,
delay_secs: 1.5,
phase: fabro_types::LlmRetryPhase::Open,
error: fabro_llm::Error::Configuration {
message: "busy".into(),
source: None,
@ -1150,6 +1319,7 @@ mod tests {
model: "gpt-5-mini".into(),
attempt: 2,
delay_secs: 1.5,
phase: fabro_types::LlmRetryPhase::Open,
error: fabro_llm::Error::Configuration {
message: "busy".into(),
source: None,

View file

@ -3,6 +3,7 @@ use std::convert::TryFrom;
use std::time::Duration;
use chrono::{DateTime, Utc};
use fabro_types::LlmOutputKind;
use fabro_workflow::outcome::{StageOutcome, format_cost};
use indicatif::ProgressBar;
@ -37,6 +38,9 @@ pub(super) struct ActiveStage {
pub(super) spinner: ProgressBar,
pub(super) tool_calls: VecDeque<ToolCallEntry>,
pub(super) compaction_bar: Option<ProgressBar>,
/// Live line for the open inference bracket. Cleared when the round ends,
/// so a stale "waiting on model" never outlives the request.
pub(super) inference_bar: Option<ProgressBar>,
}
impl ActiveStage {
@ -77,6 +81,9 @@ impl StageDisplay {
if let Some(bar) = stage.compaction_bar {
bar.finish_and_clear();
}
if let Some(bar) = stage.inference_bar {
bar.finish_and_clear();
}
for entry in &stage.tool_calls {
if entry.is_branch || self.verbose {
entry.bar.abandon();
@ -123,6 +130,7 @@ impl StageDisplay {
spinner: bar,
tool_calls: VecDeque::new(),
compaction_bar: None,
inference_bar: None,
});
}
@ -501,6 +509,69 @@ impl StageDisplay {
}
}
/// Open the live line for an inference request.
///
/// It says only what is provable: a request is open and nothing has come
/// back yet. No percentage, no ETA — the elapsed time ticks from the
/// steady tick, and it is time since the request opened, not a claim
/// about how much longer it will take.
pub(super) fn on_llm_request_started(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
model: &str,
) {
if !renderer.is_tty() {
return;
}
let Some(stage) = self.active_stages.get_mut(stage_node_id) else {
return;
};
if let Some(old) = stage.inference_bar.take() {
old.finish_and_clear();
}
let bar = renderer.insert_after(stage.last_bar());
bar.set_style(styles::style_tool_running());
bar.set_message(format!(
"\u{27f3} model request: waiting on {model}\u{2026}"
));
bar.enable_steady_tick(Duration::from_millis(100));
stage.inference_bar = Some(bar);
}
/// Replace the waiting message once the provider produces output.
///
/// "thinking" is used only for reasoning, because that is the only case
/// where the provider said so.
pub(super) fn on_llm_first_output(&mut self, stage_node_id: &str, kind: LlmOutputKind) {
let Some(bar) = self
.active_stages
.get(stage_node_id)
.and_then(|stage| stage.inference_bar.as_ref())
else {
return;
};
let activity = match kind {
LlmOutputKind::Reasoning => "reasoning",
LlmOutputKind::Text => "writing",
LlmOutputKind::ToolCall => "calling tools",
};
bar.set_message(format!("\u{27f3} model request: {activity}\u{2026}"));
}
/// Close the live line for a round that produced a message.
pub(super) fn on_llm_request_finished(&mut self, stage_node_id: &str) {
if let Some(bar) = self
.active_stages
.get_mut(stage_node_id)
.and_then(|stage| stage.inference_bar.take())
{
bar.finish_and_clear();
}
}
pub(super) fn on_llm_retry(
&mut self,
renderer: &ProgressRenderer,
@ -510,6 +581,16 @@ impl StageDisplay {
delay_ms: u64,
error: &str,
) {
if let Some(bar) = self
.active_stages
.get(stage_node_id)
.and_then(|stage| stage.inference_bar.as_ref())
{
bar.set_message(format!(
"\u{27f3} model request: waiting on {model}\u{2026}"
));
}
if !self.verbose {
return;
}
@ -592,6 +673,9 @@ impl StageDisplay {
if let Some(bar) = stage.compaction_bar {
bar.finish_and_clear();
}
if let Some(bar) = stage.inference_bar {
bar.finish_and_clear();
}
for entry in &stage.tool_calls {
if entry.is_branch || self.verbose {
entry.bar.abandon();

View file

@ -15,10 +15,10 @@ use fabro_llm::{Error as LlmError, retry};
use fabro_mcp::config::{McpServerSettings, McpTransport};
use fabro_mcp::connection_manager::McpConnectionManager;
use fabro_mcp::http_transport;
use fabro_model::{AgentProfileKind, Catalog, ModelRef, Speed, UsdMicros};
use fabro_model::{AgentProfileKind, Catalog, ModelId, ModelRef, Speed, UsdMicros};
use fabro_types::{
AgentToolSummary, PermissionLevel, Principal, SessionMessage, SessionRecord,
StageContextWindowProjection, SteeringMessage,
AgentToolSummary, LlmOutputKind, LlmRetryPhase, PermissionLevel, Principal, SessionMessage,
SessionRecord, StageContextWindowProjection, SteeringMessage,
};
use futures::StreamExt;
use tokio::sync::{Notify, broadcast};
@ -95,6 +95,25 @@ fn record_elapsed(start: &mut Option<Instant>, total: &mut Duration) {
}
}
/// Classify a stream event as the first unit of provider output, or `None`
/// when it carries no output.
///
/// `StreamStart` is deliberately excluded because it proves only that the
/// provider responded, not what kind of output followed. The start/delta/end
/// events below identify the first observed content kind.
fn first_output_kind(event: &StreamEvent) -> Option<LlmOutputKind> {
match event {
StreamEvent::ReasoningStart | StreamEvent::ReasoningDelta { .. } => {
Some(LlmOutputKind::Reasoning)
}
StreamEvent::TextStart { .. } | StreamEvent::TextDelta { .. } => Some(LlmOutputKind::Text),
StreamEvent::ToolCallStart { .. }
| StreamEvent::ToolCallDelta { .. }
| StreamEvent::ToolCallEnd { .. } => Some(LlmOutputKind::ToolCall),
_ => None,
}
}
impl SteeringItem {
#[must_use]
pub fn actor(&self) -> Option<&Principal> {
@ -1491,15 +1510,26 @@ impl Session {
let local_context_window = built_request.context_window.clone();
let request = built_request.request;
// Emit AssistantTextStart before LLM call
let requested_model = ModelRef {
provider: self.provider_profile.provider_id(),
model_id: ModelId::new(self.provider_profile.model()),
speed: self.config.speed,
};
// Open the inference bracket for this round. The request is built
// and compaction has run, so this is the last point before the
// provider is contacted at which we still know nothing about the
// response.
self.event_emitter
.emit(self.id.clone(), AgentEvent::AssistantTextStart);
.emit(self.id.clone(), AgentEvent::LlmRequestStarted {
requested_model: requested_model.clone(),
});
// Call LLM (streaming) with retry for transient errors
let retry_emitter = self.event_emitter.clone();
let retry_session_id = self.id.clone();
let retry_provider = self.provider_profile.provider_id().to_string();
let retry_model = self.provider_profile.model().to_string();
let retry_provider = requested_model.provider.to_string();
let retry_model = requested_model.model_id.to_string();
let retry_policy = RetryPolicy {
max_retries: 3,
on_retry: Some(std::sync::Arc::new(move |err, attempt, delay| {
@ -1509,6 +1539,7 @@ impl Session {
attempt: attempt as usize,
delay_secs: delay.as_secs_f64(),
error: err.clone(),
phase: LlmRetryPhase::Open,
});
})),
..Default::default()
@ -1554,6 +1585,10 @@ impl Session {
let mut accumulator = StreamAccumulator::new();
let mut attempt_emitted_output = false;
let mut stream_error = None;
// Re-armed per attempt: a replayed turn discards everything
// the previous attempt produced, so its first output is a new
// observation rather than a continuation.
let mut first_output_emitted = false;
loop {
let chunk = tokio::select! {
@ -1572,6 +1607,13 @@ impl Session {
};
match event_result {
Ok(event) => {
if !first_output_emitted {
if let Some(kind) = first_output_kind(&event) {
first_output_emitted = true;
self.event_emitter
.emit(self.id.clone(), AgentEvent::LlmFirstOutput { kind });
}
}
match &event {
StreamEvent::TextDelta { ref delta, .. } => {
attempt_emitted_output = true;
@ -1650,9 +1692,19 @@ impl Session {
);
visible_output_present = false;
}
if let Some(ref on_retry) = retry_policy.on_retry {
on_retry(&err, retry_attempt, delay);
}
// Emitted directly rather than through
// `retry_policy.on_retry` so the event can name the
// consume loop as the source of `attempt`; the policy
// callback only ever runs for stream-open failures.
self.event_emitter
.emit(self.id.clone(), AgentEvent::LlmRetry {
provider: requested_model.provider.to_string(),
model: requested_model.model_id.to_string(),
attempt: stream_attempt,
delay_secs: delay.as_secs_f64(),
error: err,
phase: LlmRetryPhase::Consume,
});
let delay_outcome = tokio::select! {
biased;
@ -1719,6 +1771,21 @@ impl Session {
);
visible_output_present = false;
}
// The only mid-turn restart that reaches no error handler:
// without this the round replays and discards its output
// with nothing on the durable stream to show for it.
self.event_emitter
.emit(self.id.clone(), AgentEvent::LlmRetry {
provider: requested_model.provider.to_string(),
model: requested_model.model_id.to_string(),
attempt: stream_attempt,
delay_secs: 0.0,
error: LlmError::Stream {
message: "Stream ended without a finish event".to_string(),
source: None,
},
phase: LlmRetryPhase::Consume,
});
let cancel_token_for_select = self.cancel_token.clone();
let retry_outcome: Option<Result<StreamEventStream, Error>> = tokio::select! {
biased;
@ -4227,6 +4294,125 @@ mod tests {
assert_eq!(tool_completed_count, 0);
}
/// Drain the receiver into `(label, detail)` pairs for the inference
/// bracket events, ignoring everything else.
fn collect_bracket_events(rx: &mut broadcast::Receiver<SessionEvent>) -> Vec<(String, String)> {
let mut observed = Vec::new();
while let Ok(event) = rx.try_recv() {
match event.event {
AgentEvent::LlmRequestStarted { requested_model } => {
observed.push((
"started".to_string(),
format!("{}/{}", requested_model.provider, requested_model.model_id),
));
}
AgentEvent::LlmFirstOutput { kind } => {
observed.push(("first_output".to_string(), kind.to_string()));
}
AgentEvent::AssistantMessage { text, .. } => {
observed.push(("message".to_string(), text));
}
_ => {}
}
}
observed
}
#[tokio::test]
async fn inference_bracket_wraps_a_text_first_turn() {
let provider = Arc::new(ScriptedStreamProvider::new(vec![
ScriptedStreamCall::Response(Box::new(text_response("Hello"))),
]));
let mut session = make_session_with_provider(provider).await;
let mut rx = session.subscribe();
session.process_input("Hi").await.unwrap();
// `started` carries the requested provider/model, and precedes any
// knowledge of what the response will contain.
assert_eq!(collect_bracket_events(&mut rx), vec![
("started".to_string(), "anthropic/mock-model".to_string()),
("first_output".to_string(), "text".to_string()),
("message".to_string(), "Hello".to_string()),
]);
}
#[tokio::test]
async fn first_output_reports_reasoning_when_reasoning_arrives_first() {
let response = text_response("Hello");
let provider = Arc::new(ScriptedStreamProvider::new(vec![
ScriptedStreamCall::Events(vec![
Ok(StreamEvent::ReasoningDelta {
delta: "weighing options".to_string(),
}),
Ok(StreamEvent::text_delta("Hello", None)),
Ok(StreamEvent::finish(
response.finish_reason.clone(),
response.usage.clone(),
response,
)),
]),
]));
let mut session = make_session_with_provider(provider).await;
let mut rx = session.subscribe();
session.process_input("Hi").await.unwrap();
// Edge-triggered: the later text delta does not re-fire the latch.
assert_eq!(collect_bracket_events(&mut rx), vec![
("started".to_string(), "anthropic/mock-model".to_string()),
("first_output".to_string(), "reasoning".to_string()),
("message".to_string(), "Hello".to_string()),
]);
}
#[tokio::test]
async fn first_output_reports_tool_call_for_a_turn_with_no_text_or_reasoning() {
let tool_call = ToolCall::new("call_1", "nonexistent_tool", serde_json::json!({}));
let mut response = tool_call_response("nonexistent_tool", "call_1", serde_json::json!({}));
// Strip the visible text so the turn produces neither a text nor a
// reasoning delta — the case a latch keyed on those two would miss
// entirely, leaving tool-heavy rounds silent.
response.message.content = vec![ContentPart::ToolCall(tool_call.clone())];
let provider = Arc::new(ScriptedStreamProvider::new(vec![
ScriptedStreamCall::Events(vec![
Ok(StreamEvent::ToolCallStart {
tool_call: tool_call.clone(),
}),
Ok(StreamEvent::ToolCallEnd {
tool_call: tool_call.clone(),
}),
Ok(StreamEvent::finish(
response.finish_reason.clone(),
response.usage.clone(),
response,
)),
]),
ScriptedStreamCall::Response(Box::new(text_response("Done"))),
]));
let mut session = make_session_with_provider(provider).await;
let mut rx = session.subscribe();
session.process_input("Use the tool").await.unwrap();
let observed = collect_bracket_events(&mut rx);
let kinds: Vec<&str> = observed
.iter()
.filter(|(label, _)| label == "first_output")
.map(|(_, kind)| kind.as_str())
.collect();
assert_eq!(kinds, vec!["tool_call", "text"]);
// One bracket per round: the tool round and the round that follows it.
assert_eq!(
observed
.iter()
.filter(|(label, _)| label == "started")
.count(),
2
);
}
#[tokio::test]
async fn stream_retries_when_stream_ends_without_finish_before_any_deltas() {
let provider = Arc::new(ScriptedStreamProvider::new(vec![
@ -4245,24 +4431,33 @@ mod tests {
Some(Message::Assistant { content, .. }) if content == "Recovered"
));
let mut assistant_text_start_count = 0;
let mut request_started_count = 0;
let mut replace_count = 0;
let mut deltas = Vec::new();
let mut assistant_messages = Vec::new();
let mut consume_retries = Vec::new();
while let Ok(event) = rx.try_recv() {
match event.event {
AgentEvent::AssistantTextStart => assistant_text_start_count += 1,
AgentEvent::LlmRequestStarted { .. } => request_started_count += 1,
AgentEvent::AssistantOutputReplace { .. } => replace_count += 1,
AgentEvent::TextDelta { delta } => deltas.push(delta),
AgentEvent::AssistantMessage { text, .. } => assistant_messages.push(text),
AgentEvent::LlmRetry { attempt, phase, .. } => {
consume_retries.push((attempt, phase));
}
_ => {}
}
}
assert_eq!(assistant_text_start_count, 1);
// One round, so one bracket open — the finish-less stream is replayed
// inside the round rather than starting a new one.
assert_eq!(request_started_count, 1);
assert_eq!(replace_count, 0);
assert_eq!(deltas, vec!["Recovered".to_string()]);
assert_eq!(assistant_messages, vec!["Recovered".to_string()]);
// The finish-less restart is the one mid-turn path with no error to
// report; without this event it would be invisible downstream.
assert_eq!(consume_retries, vec![(0, LlmRetryPhase::Consume)]);
}
#[tokio::test]
@ -4286,11 +4481,15 @@ mod tests {
let mut observed = Vec::new();
while let Ok(event) = rx.try_recv() {
match event.event {
AgentEvent::AssistantTextStart => observed.push("start".to_string()),
AgentEvent::LlmRequestStarted { .. } => observed.push("start".to_string()),
AgentEvent::LlmFirstOutput { kind } => observed.push(format!("first:{kind}")),
AgentEvent::TextDelta { delta } => observed.push(format!("delta:{delta}")),
AgentEvent::AssistantOutputReplace { text, reasoning } => {
observed.push(format!("replace:{text}:{reasoning:?}"));
}
AgentEvent::LlmRetry { phase, .. } => {
observed.push(format!("retry:{phase}"));
}
AgentEvent::AssistantMessage { text, .. } => {
observed.push(format!("message:{text}"));
}
@ -4298,10 +4497,15 @@ mod tests {
}
}
// The latch re-arms on restart: the replayed attempt's first delta is
// a fresh observation, not a continuation of the discarded one.
assert_eq!(observed, vec![
"start".to_string(),
"first:text".to_string(),
"delta:Hel".to_string(),
"replace::None".to_string(),
"retry:consume".to_string(),
"first:text".to_string(),
"delta:Hello".to_string(),
"message:Hello".to_string(),
]);
@ -4339,7 +4543,7 @@ mod tests {
let mut found_auth_error_event = false;
while let Ok(event) = rx.try_recv() {
match event.event {
AgentEvent::AssistantTextStart => observed.push("start".to_string()),
AgentEvent::LlmRequestStarted { .. } => observed.push("start".to_string()),
AgentEvent::TextDelta { delta } => observed.push(format!("delta:{delta}")),
AgentEvent::AssistantOutputReplace { text, reasoning } => {
observed.push(format!("replace:{text}:{reasoning:?}"));

View file

@ -5,8 +5,8 @@ use fabro_llm::Error as LlmError;
use fabro_llm::types::{ContentPart, ThinkingData, TokenCounts, ToolCall, ToolResult};
use fabro_model::{CostSource, ModelRef};
use fabro_types::{
CommandTermination, ExecOutputTail, ReasoningOutput, SessionMessage,
StageContextWindowProjection,
CommandTermination, ExecOutputTail, LlmOutputKind, LlmRetryPhase, ReasoningOutput,
SessionMessage, StageContextWindowProjection,
};
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
@ -237,7 +237,20 @@ pub enum AgentEvent {
UserInput {
text: String,
},
AssistantTextStart,
/// An inference request is about to be dispatched for this round. Emitted
/// after the request is built and compaction has run, immediately before
/// the stream is opened. `provider` and `model` are the requested target;
/// failover can re-target, so `AssistantMessage` stays authoritative for
/// what actually answered.
LlmRequestStarted {
requested_model: ModelRef,
},
/// The provider produced its first output for the current attempt.
/// Edge-triggered: emitted once per stream attempt, re-armed when a
/// broken or finish-less stream restarts the turn.
LlmFirstOutput {
kind: LlmOutputKind,
},
/// Replaces the current in-progress assistant output buffers.
AssistantOutputReplace {
text: String,
@ -328,12 +341,15 @@ pub enum AgentEvent {
summary_token_estimate: usize,
tracked_file_count: usize,
},
/// An attempt failed to open **or sustain** a stream and the turn is
/// being replayed. `phase` names which retry loop `attempt` counts.
LlmRetry {
provider: String,
model: String,
attempt: usize,
delay_secs: f64,
error: LlmError,
phase: LlmRetryPhase,
},
SubAgentSpawned {
agent_id: String,
@ -396,8 +412,7 @@ impl AgentEvent {
pub fn is_streaming_noise(&self) -> bool {
matches!(
self,
Self::AssistantTextStart
| Self::AssistantOutputReplace { .. }
Self::AssistantOutputReplace { .. }
| Self::TextDelta { .. }
| Self::ReasoningDelta { .. }
| Self::ToolCallOutputDelta { .. }
@ -424,8 +439,17 @@ impl AgentEvent {
Self::UserInput { text } => {
debug!(session_id, text_len = text.len(), "User input received");
}
Self::AssistantTextStart => {
debug!(session_id, "Assistant response started");
Self::LlmRequestStarted { requested_model } => {
debug!(
session_id,
provider = %requested_model.provider,
model = %requested_model.model_id,
speed = requested_model.speed.map_or("", <&'static str>::from),
"LLM request started"
);
}
Self::LlmFirstOutput { kind } => {
debug!(session_id, kind = %kind, "LLM produced first output");
}
Self::AssistantMessage {
model,
@ -545,6 +569,7 @@ impl AgentEvent {
attempt,
delay_secs,
error,
phase,
} => {
warn!(
session_id,
@ -552,6 +577,7 @@ impl AgentEvent {
model,
attempt,
delay_secs,
phase = %phase,
error = %error,
"LLM request failed, retrying"
);
@ -1001,6 +1027,7 @@ mod tests {
model: "gpt-4".into(),
attempt: 1,
delay_secs: 2.0,
phase: LlmRetryPhase::Open,
error: LlmError::Provider {
kind: ProviderErrorKind::RateLimit,
detail: Box::new(ProviderErrorDetail {

View file

@ -121,7 +121,8 @@ impl SseAccumulator {
.unwrap_or(0);
}
}
vec![StreamEvent::StreamStart]
// `StreamStart` is the driver's; this handler only captures metadata.
vec![]
}
fn handle_content_block_start(&mut self, data: &serde_json::Value) -> Vec<StreamEvent> {

View file

@ -271,7 +271,6 @@ impl StreamDecoder for ConverseStreamDecoder {
.map_err(|e| Error::stream_error(format!("converse stream event json: {e}"), e))?;
Ok(match event_type {
"messageStart" => vec![StreamEvent::StreamStart],
"contentBlockStart" => self.on_block_start(&payload),
"contentBlockDelta" => self.on_block_delta(&payload),
"contentBlockStop" => self.on_block_stop(&payload),
@ -284,7 +283,9 @@ impl StreamDecoder for ConverseStreamDecoder {
self.usage = token_counts_from_usage(payload.get("usage"));
vec![self.finish_event()]
}
// Tolerate unknown event types — the union grows.
// `messageStart` carries nothing this decoder needs — the driving
// loop owns `StreamStart` — and unknown event types are tolerated
// because the union grows.
_ => Vec::new(),
})
}
@ -347,10 +348,8 @@ mod tests {
#[test]
fn text_happy_path_finishes_on_metadata() {
let mut d = decoder();
assert!(matches!(
feed(&mut d, "messageStart", r#"{"role":"assistant"}"#)[0],
StreamEvent::StreamStart
));
// `StreamStart` belongs to the driving loop, not the decoder.
assert!(feed(&mut d, "messageStart", r#"{"role":"assistant"}"#).is_empty());
let events = feed(
&mut d,
"contentBlockDelta",

View file

@ -21,8 +21,6 @@ pub(super) struct SseAccumulator {
model: String,
/// Configured provider name stamped into the final `Response.provider`.
provider: String,
/// Whether we have emitted a `StreamStart` event.
stream_started: bool,
/// Whether we have emitted a `TextStart` event.
text_started: bool,
/// Whether we are currently inside a reasoning (thought) segment.
@ -50,7 +48,6 @@ impl SseAccumulator {
Self {
model: ctx.request.model.clone(),
provider: ctx.provider_name.to_string(),
stream_started: false,
text_started: false,
reasoning_started: false,
accumulated_thinking: String::new(),
@ -68,11 +65,6 @@ impl SseAccumulator {
fn process_chunk(&mut self, chunk: &ApiResponse) -> Vec<StreamEvent> {
let mut events = Vec::new();
if !self.stream_started {
self.stream_started = true;
events.push(StreamEvent::StreamStart);
}
let parts = chunk
.candidates
.as_ref()
@ -254,13 +246,11 @@ mod tests {
/// Build an accumulator without threading a `CodecCtx`/`Request`: the test
/// module sees the private fields, so the few that matter are set
/// directly. `stream_started` is true so event assertions don't see the
/// initial `StreamStart`.
/// directly.
fn empty_accumulator() -> SseAccumulator {
SseAccumulator {
model: "gemini-2.0-flash".to_string(),
provider: "gemini".to_string(),
stream_started: true,
text_started: false,
reasoning_started: false,
accumulated_thinking: String::new(),
@ -279,9 +269,8 @@ mod tests {
}
#[test]
fn first_chunk_emits_stream_start() {
fn first_chunk_opens_text_without_a_decoder_level_stream_start() {
let mut acc = empty_accumulator();
acc.stream_started = false;
let events = on_data(
&mut acc,
@ -289,8 +278,10 @@ mod tests {
)
.expect("chunk should parse");
assert!(matches!(events[0], StreamEvent::StreamStart));
assert!(matches!(events[1], StreamEvent::TextStart { .. }));
// `StreamStart` is the driving loop's, so the decoder's first event
// is the content itself.
assert!(matches!(events[0], StreamEvent::TextStart { .. }));
assert!(matches!(events[1], StreamEvent::TextDelta { .. }));
}
#[test]

View file

@ -222,6 +222,14 @@ pub(crate) trait StreamDecoder: Send + 'static {
/// One framed event → zero or more canonical `StreamEvent`s. Returns
/// `Err` for dialect error events (anthropic `error`, openai
/// `response.failed`), which the transport yields as a stream error.
///
/// Decoders must **not** emit [`StreamEvent::StreamStart`]. The driving
/// loop emits exactly one, immediately before handing over the first
/// framed event, so `StreamStart` means the same thing for every
/// provider: the provider is responding, whatever it turns out to say.
/// Leaving it to decoders made it depend on each dialect's opening frame
/// — anthropic and bedrock keyed it on `message_start`/`messageStart`,
/// and `openai_compatible` had no such frame and so emitted it never.
fn on_event(&mut self, ev: RawEvent<'_>) -> Result<Vec<StreamEvent>, Error>;
/// Byte-stream-end hook. Semantics are per-decoder, not shared:

View file

@ -95,7 +95,6 @@ pub(super) struct SseAccumulator {
message_items: Vec<serde_json::Value>,
usage: TokenCounts,
finish_reason: FinishReason,
emitted_start: bool,
emitted_text_start: bool,
emitted_reasoning_start: bool,
rate_limit: Option<RateLimitInfo>,
@ -114,7 +113,6 @@ impl SseAccumulator {
message_items: Vec::new(),
usage: TokenCounts::default(),
finish_reason: FinishReason::Stop,
emitted_start: false,
emitted_text_start: false,
emitted_reasoning_start: false,
rate_limit,
@ -130,11 +128,6 @@ impl SseAccumulator {
) -> Result<Vec<StreamEvent>, Error> {
let mut events = Vec::new();
if !self.emitted_start {
self.emitted_start = true;
events.push(StreamEvent::StreamStart);
}
let json: serde_json::Value = match serde_json::from_str(data) {
Ok(v) => v,
Err(_) => return Ok(events),
@ -446,8 +439,7 @@ mod tests {
/// Build an accumulator without threading a `CodecCtx`/`Request`: the test
/// module sees the private fields, so the few that matter are set
/// directly. `emitted_start` is true so event assertions don't see the
/// initial `StreamStart`.
/// directly.
fn empty_accumulator() -> SseAccumulator {
SseAccumulator {
model: String::new(),
@ -460,7 +452,6 @@ mod tests {
message_items: Vec::new(),
usage: TokenCounts::default(),
finish_reason: FinishReason::Stop,
emitted_start: true,
emitted_text_start: false,
emitted_reasoning_start: false,
rate_limit: None,

View file

@ -309,14 +309,16 @@ impl ProviderAdapter for Adapter {
/// State driving the event-stream byte loop: the codec's decoder plus the
/// frame decoder, with a buffer that flattens batched events.
struct EventStreamLoop {
response: fabro_http::Response,
frames: FrameDecoder,
decoder: Box<dyn StreamDecoder>,
pending: VecDeque<StreamEvent>,
done: bool,
response: fabro_http::Response,
frames: FrameDecoder,
decoder: Box<dyn StreamDecoder>,
pending: VecDeque<Result<StreamEvent, Error>>,
done: bool,
/// `finish()` already drained.
finished: bool,
timeout: Option<Duration>,
finished: bool,
/// [`StreamEvent::StreamStart`] already emitted for this stream.
stream_started: bool,
timeout: Option<Duration>,
}
/// Drive `decoder` over the AWS event-stream byte stream of `response`: the
@ -335,12 +337,13 @@ fn decode_eventstream(
pending: VecDeque::new(),
done: false,
finished: false,
stream_started: false,
timeout,
},
move |mut state| async move {
loop {
if let Some(event) = state.pending.pop_front() {
return Some((Ok(event), state));
return Some((event, state));
}
if state.done {
@ -348,7 +351,9 @@ fn decode_eventstream(
return None;
}
state.finished = true;
state.pending.extend(state.decoder.finish());
state
.pending
.extend(state.decoder.finish().into_iter().map(Ok));
if state.pending.is_empty() {
return None;
}
@ -370,9 +375,19 @@ fn decode_eventstream(
event: Some(frame.event_type.as_str()),
data: frame.payload.as_str(),
};
// Mirrors the SSE loop: the first decoded frame is
// the liveness edge, independent of which event
// type the provider happens to open with.
if !state.stream_started {
state.stream_started = true;
state.pending.push_back(Ok(StreamEvent::StreamStart));
}
match state.decoder.on_event(raw) {
Ok(events) => state.pending.extend(events),
Err(e) => return Some((Err(e), state)),
Ok(events) => state.pending.extend(events.into_iter().map(Ok)),
Err(error) => {
state.pending.push_back(Err(error));
break;
}
}
}
}

View file

@ -247,12 +247,14 @@ pub(crate) async fn stream_via_http(
struct StreamLoop {
decoder: Box<dyn StreamDecoder>,
line_reader: LineReader,
/// Events decoded but not yet yielded.
pending: VecDeque<StreamEvent>,
/// Events or decoder errors not yet yielded.
pending: VecDeque<Result<StreamEvent, Error>>,
/// Byte stream exhausted.
done: bool,
/// `finish()` already drained.
finished_emitted: bool,
/// [`StreamEvent::StreamStart`] already emitted for this stream.
stream_started: bool,
}
/// Drive `decoder` over the SSE byte stream of `response`: frame each chunk,
@ -271,11 +273,12 @@ fn decode_sse_stream(
pending: VecDeque::new(),
done: false,
finished_emitted: false,
stream_started: false,
},
move |mut state| async move {
loop {
if let Some(event) = state.pending.pop_front() {
return Some((Ok(event), state));
return Some((event, state));
}
if state.done {
@ -283,7 +286,9 @@ fn decode_sse_stream(
return None;
}
state.finished_emitted = true;
state.pending.extend(state.decoder.finish());
state
.pending
.extend(state.decoder.finish().into_iter().map(Ok));
if state.pending.is_empty() {
return None;
}
@ -295,9 +300,18 @@ fn decode_sse_stream(
let Some((event, data)) = frame_sse_chunk(framing, &chunk) else {
continue;
};
// Provider-independent liveness edge: the first framed
// event proves the provider is responding, whatever it
// turns out to contain. Owned here rather than in each
// decoder so it cannot depend on a provider sending a
// particular opening frame.
if !state.stream_started {
state.stream_started = true;
state.pending.push_back(Ok(StreamEvent::StreamStart));
}
match state.decoder.on_event(RawEvent { event, data: &data }) {
Ok(events) => state.pending.extend(events),
Err(e) => return Some((Err(e), state)),
Ok(events) => state.pending.extend(events.into_iter().map(Ok)),
Err(error) => state.pending.push_back(Err(error)),
}
}
Ok(None) => state.done = true,

View file

@ -136,6 +136,18 @@ pub(crate) async fn collect_stream_events(
events
}
/// Pin the transport-level liveness contract independently of snapshots.
pub(crate) fn assert_stream_starts(events: &[serde_json::Value]) {
assert_eq!(
events
.first()
.and_then(|event| event.get("type"))
.and_then(serde_json::Value::as_str),
Some("stream_start"),
"the first decoded provider frame must open with stream_start"
);
}
/// Builds a catalog from inline TOML (same `LlmCatalogSettings` schema as the
/// shipped catalog files).
pub(crate) fn catalog_from_toml(source: &str) -> Arc<Catalog> {

View file

@ -648,6 +648,7 @@ async fn stream_text_happy_path_request() {
#[tokio::test]
async fn stream_text_happy_path_events() {
let (_, events) = stream_text_happy_path_capture().await;
support::assert_stream_starts(&events);
fabro_test::fabro_json_snapshot!(events);
}

View file

@ -448,6 +448,7 @@ async fn stream_text_happy_path_request() {
#[tokio::test]
async fn stream_text_happy_path_events() {
let (_, events) = stream_text_happy_path_capture().await;
support::assert_stream_starts(&events);
fabro_test::fabro_json_snapshot!(events);
}

View file

@ -823,6 +823,7 @@ async fn stream_text_happy_path_request() {
#[tokio::test]
async fn stream_text_happy_path_events() {
let (_, events) = stream_text_happy_path_capture().await;
support::assert_stream_starts(&events);
fabro_test::fabro_json_snapshot!(events);
}
@ -881,7 +882,9 @@ async fn stream_without_done_synthesizes_finish_when_content_started() {
}
/// The other half of the minimax contract: no content started and no
/// `[DONE]` — nothing is synthesized.
/// `[DONE]` — nothing is synthesized. `StreamStart` is not synthesis: the
/// provider did send a chunk, so the liveness edge is a fact about this
/// stream even though nothing usable followed.
#[tokio::test]
async fn stream_without_done_or_content_synthesizes_nothing() {
let sse = support::sse_data_transcript(&[

View file

@ -542,9 +542,27 @@ async fn stream_text_happy_path_request() {
#[tokio::test]
async fn stream_text_happy_path_events() {
let (_, events) = stream_text_happy_path_capture().await;
support::assert_stream_starts(&events);
fabro_test::fabro_json_snapshot!(events);
}
#[tokio::test]
async fn stream_first_frame_error_still_opens_with_stream_start() {
let sse = support::sse_data_transcript(&[
r#"{"type":"response.failed","response":{"id":"resp_stream","error":{"code":"server_error","message":"boom"}}}"#,
]);
let (_capture, events) = stream_capture(adapter(), &base_request(MODEL), &sse).await;
support::assert_stream_starts(&events);
assert!(
events
.get(1)
.and_then(|event| event.get("stream_item_error"))
.is_some(),
"the decoder error should follow stream_start: {events:?}"
);
}
#[tokio::test]
async fn stream_tool_call_deltas() {
let sse = support::sse_data_transcript(&[

View file

@ -3,6 +3,9 @@ source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs
expression: rendered
---
[
{
"type": "stream_start"
},
{
"type": "text_start",
"text_id": null

View file

@ -3,6 +3,9 @@ source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs
expression: rendered
---
[
{
"type": "stream_start"
},
{
"type": "text_start",
"text_id": null

View file

@ -3,6 +3,9 @@ source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs
expression: rendered
---
[
{
"type": "stream_start"
},
{
"type": "tool_call_start",
"tool_call": {

View file

@ -3,6 +3,9 @@ source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs
expression: rendered
---
[
{
"type": "stream_start"
},
{
"type": "text_start",
"text_id": null

View file

@ -2,4 +2,8 @@
source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs
expression: rendered
---
[]
[
{
"type": "stream_start"
}
]

View file

@ -3,6 +3,9 @@ source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs
expression: rendered
---
[
{
"type": "stream_start"
},
{
"type": "text_start",
"text_id": null

View file

@ -4,8 +4,8 @@ use std::sync::Arc;
use chrono::{DateTime, Utc};
use fabro_types::run_event::{
CheckpointCompletedProps, RunCompletedProps, RunFailedProps, StageCompletedProps,
TodoCreatedProps, TodoDeletedProps, TodoUpdatedProps,
AgentLlmStartedProps, CheckpointCompletedProps, RunCompletedProps, RunFailedProps,
StageCompletedProps, TodoCreatedProps, TodoDeletedProps, TodoUpdatedProps,
};
use fabro_types::settings::run::{EnvironmentProvider, RunEnvironmentSettings};
use fabro_types::{
@ -16,9 +16,9 @@ use fabro_types::{
RunBillingSummary, RunControlAction, RunDiff, RunEvent, RunId, RunLifecycle, RunLinks,
RunModel, RunOrigin, RunProjection, RunSandbox, RunSandboxFailure, RunSandboxInstance,
RunSandboxPlan, RunSandboxRuntime, RunSize, RunSpec, RunStatus, RunTimestamps,
SandboxProviderKind, StageCompletion, StageHandler, StageId, StageModelUsage, StageOutcome,
StageProjection, StageState, StartRecord, SubAgentProjection, SubAgentStatus, TodoListKind,
TodoListProjection, TodoProjection, WorkflowRef, first_event_seq,
SandboxProviderKind, StageCompletion, StageHandler, StageId, StageInferenceProjection,
StageModelUsage, StageOutcome, StageProjection, StageState, StartRecord, SubAgentProjection,
SubAgentStatus, TodoListKind, TodoListProjection, TodoProjection, WorkflowRef, first_event_seq,
};
use fabro_util::error::render_compact_with_causes;
@ -449,6 +449,39 @@ impl RunProjectionReducer for RunProjection {
context_window.event_seq = Some(event.seq);
stage.context_window = Some(context_window);
}
close_inference_bracket(self, stored, props.visit, event.seq);
}
EventBody::AgentLlmStarted(props) => {
open_inference_bracket(self, stored, props, event.seq, ts);
}
EventBody::AgentLlmFirstOutput(props) => {
let Some(inference) =
matching_inference_bracket(self, stored, props.visit, event.seq)
else {
return Ok(());
};
inference.first_output_at = Some(ts);
inference.first_output_kind = Some(props.kind);
}
EventBody::AgentLlmRetry(props) => {
let Some(inference) =
matching_inference_bracket(self, stored, props.visit, event.seq)
else {
return Ok(());
};
inference.retries = inference.retries.saturating_add(1);
// A retry discards whatever the failed attempt produced.
// Replay is driven purely by events, so resetting the
// in-process latch is not enough: without this the projection
// keeps asserting output the agent already threw away.
inference.first_output_at = None;
inference.first_output_kind = None;
}
EventBody::AgentError(props) => {
close_inference_bracket(self, stored, props.visit, event.seq);
}
EventBody::AgentSessionEnded(_) => {
close_inference_brackets_for_session(self, stored);
}
EventBody::AgentSessionActivated(props) => {
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
@ -464,6 +497,7 @@ impl RunProjectionReducer for RunProjection {
return Ok(());
};
stage.agent_control = AgentControlState::WaitingForSteer;
close_inference_bracket(self, stored, props.visit, event.seq);
}
EventBody::AgentSteeringInjected(props) => {
let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq)
@ -966,6 +1000,112 @@ fn stage_at_stored_or_visit<'a>(
stage_at_visit(state, stored, visit, seq)
}
/// Open the inference bracket for the stage this event names.
///
/// Only root-session events open a bracket: forwarded child events carry a
/// parent session and must not overwrite the parent's bracket. The session id
/// is copied from the envelope so every later transition can be gated on it.
fn open_inference_bracket(
state: &mut RunProjection,
stored: &RunEvent,
props: &AgentLlmStartedProps,
seq: u32,
ts: DateTime<Utc>,
) {
if stored.parent_session_id.is_some() {
return;
}
let Some(session_id) = stored.session_id.clone() else {
return;
};
let Some(stage) = stage_at_stored_or_visit(state, stored, props.visit, seq) else {
return;
};
stage.inference = Some(StageInferenceProjection {
session_id,
started_at: ts,
requested_model: props.requested_model.clone(),
first_output_at: None,
first_output_kind: None,
retries: 0,
});
}
/// Resolve the stage's `inference` slot when it holds a bracket this event is
/// allowed to mutate.
///
/// `None` when the event came from a child session, when the stage has no
/// open bracket, or when the bracket belongs to a different session — the
/// last case matters after failover, which discards the session and builds a
/// new one within a single visit.
fn matching_inference_slot<'a>(
state: &'a mut RunProjection,
stored: &RunEvent,
visit: u32,
seq: u32,
) -> Option<&'a mut Option<StageInferenceProjection>> {
if stored.parent_session_id.is_some() {
return None;
}
let session_id = stored.session_id.as_deref()?;
let stage = stage_at_stored_or_visit(state, stored, visit, seq)?;
let opened_here = stage
.inference
.as_ref()
.is_some_and(|inference| inference.session_id == session_id);
opened_here.then_some(&mut stage.inference)
}
/// Resolve the open inference bracket this event is allowed to mutate.
fn matching_inference_bracket<'a>(
state: &'a mut RunProjection,
stored: &RunEvent,
visit: u32,
seq: u32,
) -> Option<&'a mut StageInferenceProjection> {
matching_inference_slot(state, stored, visit, seq)?.as_mut()
}
/// Close the bracket on a stage-addressed terminal event.
fn close_inference_bracket(state: &mut RunProjection, stored: &RunEvent, visit: u32, seq: u32) {
if let Some(slot) = matching_inference_slot(state, stored, visit, seq) {
*slot = None;
}
}
/// Close every bracket opened by the session that just ended.
///
/// `agent.session.ended` is the only ordering-safe backstop for terminal
/// cancel and wall-clock timeout, which tear the session down through
/// `discard_session` without emitting a message, error, or interrupt. It is
/// emitted after the forwarder drains queued agent events, unlike
/// `agent.session.deactivated`, which is emitted before the drain and so can
/// be followed by a queued `agent.llm.started` that would re-open the bracket.
///
/// The tradeoff is that it carries no stage identity — its props are empty and
/// its envelope has only session ids. So the close takes ordering from the
/// event and identity from the projection, scanning for brackets this session
/// opened. Implemented as a normal stage lookup it would find no target and
/// silently no-op, leaving the bracket open forever on exactly the path it
/// exists to cover.
fn close_inference_brackets_for_session(state: &mut RunProjection, stored: &RunEvent) {
if stored.parent_session_id.is_some() {
return;
}
let Some(session_id) = stored.session_id.as_deref() else {
return;
};
for (_, stage) in state.iter_stages_unordered_mut() {
let opened_here = stage
.inference
.as_ref()
.is_some_and(|inference| inference.session_id == session_id);
if opened_here {
stage.inference = None;
}
}
}
fn stage_at_stored_or_current_visit<'a>(
state: &'a mut RunProjection,
stored: &RunEvent,
@ -6141,4 +6281,260 @@ mod tests {
}
}
}
mod inference_bracket_reducer {
use fabro_types::run_event::{
AgentErrorProps, AgentLlmFirstOutputProps, AgentLlmRetryProps, AgentLlmStartedProps,
};
use fabro_types::{
LlmOutputKind, LlmRetryPhase, ModelRef, Speed, StageInferenceProjection,
};
use super::*;
const ROOT: &str = "ses_root";
fn stage_id() -> StageId {
StageId::new("code", 1)
}
/// Stage-addressed event attributed to a root session.
fn root_event(seq: u32, body: EventBody) -> EventEnvelope {
let mut event = test_stage_event(seq, body, stage_id());
event.event.session_id = Some(ROOT.to_string());
event
}
/// Forwarded child-session event: same stage, but with a parent link.
fn child_event(seq: u32, body: EventBody) -> EventEnvelope {
let mut event = test_stage_event(seq, body, stage_id());
event.event.session_id = Some("ses_child".to_string());
event.event.parent_session_id = Some(ROOT.to_string());
event
}
/// `agent.session.ended` as it is actually stored: session ids only,
/// no `node_id` and no `stage_id`.
fn session_ended_event(seq: u32, session_id: &str) -> EventEnvelope {
let mut event = test_event(
seq,
EventBody::AgentSessionEnded(AgentSessionEndedProps {}),
None,
);
event.event.session_id = Some(session_id.to_string());
event
}
fn started() -> EventBody {
EventBody::AgentLlmStarted(AgentLlmStartedProps {
requested_model: ModelRef {
provider: "anthropic".parse().unwrap(),
model_id: "claude-fable-5".into(),
speed: Some(Speed::Fast),
},
visit: 1,
})
}
fn first_output(kind: LlmOutputKind) -> EventBody {
EventBody::AgentLlmFirstOutput(AgentLlmFirstOutputProps { kind, visit: 1 })
}
fn retry(phase: LlmRetryPhase) -> EventBody {
EventBody::AgentLlmRetry(AgentLlmRetryProps {
provider: "anthropic".to_string(),
model: "claude-fable-5".to_string(),
attempt: 0,
delay_secs: 0.0,
error: json!({ "kind": "stream" }),
phase: Some(phase),
visit: 1,
})
}
fn open_bracket(state: &RunProjection) -> Option<&StageInferenceProjection> {
state.stage(&stage_id()).unwrap().inference.as_ref()
}
#[test]
fn started_opens_a_bracket_carrying_the_requested_model() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
let inference = open_bracket(&state).expect("bracket should be open");
assert_eq!(inference.session_id, ROOT);
assert_eq!(inference.requested_model.provider.as_str(), "anthropic");
assert_eq!(
inference.requested_model.model_id.as_str(),
"claude-fable-5"
);
assert_eq!(inference.requested_model.speed, Some(Speed::Fast));
assert_eq!(inference.first_output_at, None);
assert_eq!(inference.first_output_kind, None);
assert_eq!(inference.retries, 0);
}
#[test]
fn first_output_records_the_observed_kind() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
state
.apply_event(&root_event(2, first_output(LlmOutputKind::ToolCall)))
.unwrap();
let inference = open_bracket(&state).unwrap();
assert!(inference.first_output_at.is_some());
assert_eq!(inference.first_output_kind, Some(LlmOutputKind::ToolCall));
}
#[test]
fn message_closes_the_bracket() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
state
.apply_event(&root_event(2, first_output(LlmOutputKind::Text)))
.unwrap();
state
.apply_event(&root_event(
3,
EventBody::AgentMessage(live_agent_message_props(live_counts(10, 5))),
))
.unwrap();
assert!(open_bracket(&state).is_none());
// The close must not undo the rest of the message's work.
assert_eq!(state.stage(&stage_id()).unwrap().usage.input_tokens, 10);
}
#[test]
fn error_and_round_interrupt_close_the_bracket() {
for close in [
EventBody::AgentError(AgentErrorProps {
error: json!({ "message": "boom" }),
visit: 1,
}),
EventBody::AgentRoundInterrupted(AgentRoundInterruptedProps {
generation: 1,
visit: 1,
}),
] {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
state.apply_event(&root_event(2, close)).unwrap();
assert!(open_bracket(&state).is_none());
}
}
#[test]
fn retry_counts_the_attempt_and_clears_observed_output() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
state
.apply_event(&root_event(2, first_output(LlmOutputKind::Text)))
.unwrap();
state
.apply_event(&root_event(3, retry(LlmRetryPhase::Consume)))
.unwrap();
// Replay is event-driven, so the projection must forget output the
// agent discarded rather than keep asserting it.
let inference = open_bracket(&state).unwrap();
assert_eq!(inference.retries, 1);
assert_eq!(inference.first_output_at, None);
assert_eq!(inference.first_output_kind, None);
state
.apply_event(&root_event(4, first_output(LlmOutputKind::Reasoning)))
.unwrap();
state
.apply_event(&root_event(5, retry(LlmRetryPhase::Open)))
.unwrap();
let inference = open_bracket(&state).unwrap();
assert_eq!(inference.retries, 2);
assert_eq!(inference.first_output_kind, None);
}
#[test]
fn session_ended_closes_a_bracket_it_cannot_address_by_stage() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
assert!(open_bracket(&state).is_some());
// Terminal cancel path: no message, error, or interrupt is ever
// emitted. The envelope carries no node_id or stage_id, so a
// normal stage lookup would silently no-op and leave the bracket
// open forever.
state.apply_event(&session_ended_event(2, ROOT)).unwrap();
assert!(open_bracket(&state).is_none());
}
#[test]
fn session_ended_from_another_session_leaves_the_bracket_open() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
state
.apply_event(&session_ended_event(2, "ses_other"))
.unwrap();
assert!(open_bracket(&state).is_some());
}
#[test]
fn a_start_queued_behind_deactivation_does_not_leave_a_phantom_bracket() {
let mut state = initialized_projection();
// `lease.release()` emits deactivation *before* the forwarder
// drains queued agent events, so a queued start legitimately
// arrives after it. Keying the close on deactivation would clear
// the bracket and then immediately re-open it.
state
.apply_event(&root_event(
1,
EventBody::AgentSessionDeactivated(AgentSessionDeactivatedProps { visit: 1 }),
))
.unwrap();
state.apply_event(&root_event(2, started())).unwrap();
state.apply_event(&session_ended_event(3, ROOT)).unwrap();
assert!(open_bracket(&state).is_none());
}
#[test]
fn child_session_events_do_not_touch_the_root_bracket() {
let mut state = initialized_projection();
state.apply_event(&root_event(1, started())).unwrap();
state
.apply_event(&root_event(2, first_output(LlmOutputKind::Text)))
.unwrap();
// A sub-agent runs its own rounds on the same stage. None of them
// may open, advance, or close the root session's bracket.
state.apply_event(&child_event(3, started())).unwrap();
state
.apply_event(&child_event(4, first_output(LlmOutputKind::ToolCall)))
.unwrap();
state
.apply_event(&child_event(5, retry(LlmRetryPhase::Open)))
.unwrap();
state
.apply_event(&session_ended_event(6, "ses_child"))
.unwrap();
let inference = open_bracket(&state).expect("root bracket should survive");
assert_eq!(inference.session_id, ROOT);
assert_eq!(inference.first_output_kind, Some(LlmOutputKind::Text));
assert_eq!(inference.retries, 0);
}
#[test]
fn transitions_without_an_open_bracket_are_ignored() {
let mut state = initialized_projection();
state
.apply_event(&root_event(1, first_output(LlmOutputKind::Text)))
.unwrap();
state
.apply_event(&root_event(2, retry(LlmRetryPhase::Open)))
.unwrap();
assert!(open_bracket(&state).is_none());
}
}
}

View file

@ -72,12 +72,36 @@ impl RunProjectionCacheState {
let Some(parent_id) = entry.summary.parent_id else {
return;
};
let Some(children) = self.children_by_parent.get_mut(&parent_id) else {
self.remove_parent_link(&parent_id, &entry.run_id);
}
fn remove_parent_link(&mut self, parent_id: &RunId, run_id: &RunId) {
let Some(children) = self.children_by_parent.get_mut(parent_id) else {
return;
};
children.remove(&entry.run_id);
children.remove(run_id);
if children.is_empty() {
self.children_by_parent.remove(&parent_id);
self.children_by_parent.remove(parent_id);
}
}
fn update_parent_index(
&mut self,
run_id: RunId,
previous_parent_id: Option<RunId>,
parent_id: Option<RunId>,
) {
if previous_parent_id == parent_id {
return;
}
if let Some(previous_parent_id) = previous_parent_id {
self.remove_parent_link(&previous_parent_id, &run_id);
}
if let Some(parent_id) = parent_id {
self.children_by_parent
.entry(parent_id)
.or_default()
.insert(run_id);
}
}
@ -204,7 +228,7 @@ impl RunProjectionCache {
event: &EventEnvelope,
) -> Result<CachedRunProjection> {
let mut state = self.state.lock().await;
let Some(entry) = state.entries.get(run_id).cloned() else {
let Some(entry) = state.entries.get(run_id) else {
if event.seq == 1 {
let projection = RunProjection::apply_events(std::slice::from_ref(event))?;
let entry = CachedRunProjection::from_projection(*run_id, projection, event.seq);
@ -217,20 +241,29 @@ impl RunProjectionCache {
)));
};
if event.seq <= entry.last_seq {
return Ok(entry);
let last_seq = entry.last_seq;
if event.seq <= last_seq {
return Ok(entry.clone());
}
if event.seq != entry.last_seq.saturating_add(1) {
if event.seq != last_seq.saturating_add(1) {
return Err(Error::Other(format!(
"projection cache sequence gap for run {run_id}: last_seq={}, event_seq={}",
entry.last_seq, event.seq
last_seq, event.seq
)));
}
let mut projection = (*entry.projection).clone();
projection.apply_event(event)?;
let entry = CachedRunProjection::from_projection(*run_id, projection, event.seq);
state.insert(entry.clone());
let (previous_parent_id, parent_id, entry) = {
let entry = state
.entries
.get_mut(run_id)
.expect("entry was read from the same locked map");
let previous_parent_id = entry.summary.parent_id;
Arc::make_mut(&mut entry.projection).apply_event(event)?;
entry.summary = build_summary(&entry.projection, run_id);
entry.last_seq = event.seq;
(previous_parent_id, entry.summary.parent_id, entry.clone())
};
state.update_parent_index(*run_id, previous_parent_id, parent_id);
Ok(entry)
}

View file

@ -720,18 +720,32 @@ fn event_body_from_event(event: &Event) -> EventBody {
tracked_file_count: *tracked_file_count,
visit: *visit,
}),
AgentEvent::LlmRequestStarted { requested_model } => {
EventBody::AgentLlmStarted(fabro_types::AgentLlmStartedProps {
requested_model: requested_model.clone(),
visit: *visit,
})
}
AgentEvent::LlmFirstOutput { kind } => {
EventBody::AgentLlmFirstOutput(fabro_types::AgentLlmFirstOutputProps {
kind: *kind,
visit: *visit,
})
}
AgentEvent::LlmRetry {
provider,
model,
attempt,
delay_secs,
error,
phase,
} => EventBody::AgentLlmRetry(fabro_types::AgentLlmRetryProps {
provider: provider.clone(),
model: model.clone(),
attempt: *attempt,
delay_secs: *delay_secs,
error: serde_json::to_value(error).expect("LLM SDK error derives Serialize with no custom logic that can fail"),
phase: Some(*phase),
visit: *visit,
}),
AgentEvent::SubAgentSpawned {
@ -849,8 +863,7 @@ fn event_body_from_event(event: &Event) -> EventBody {
AgentEvent::TodoCreated(props) => EventBody::TodoCreated(props.clone()),
AgentEvent::TodoUpdated(props) => EventBody::TodoUpdated(props.clone()),
AgentEvent::TodoDeleted(props) => EventBody::TodoDeleted(props.clone()),
AgentEvent::AssistantTextStart
| AgentEvent::AssistantOutputReplace { .. }
AgentEvent::AssistantOutputReplace { .. }
| AgentEvent::TextDelta { .. }
| AgentEvent::ReasoningDelta { .. }
| AgentEvent::ToolCallOutputDelta { .. }

View file

@ -67,7 +67,8 @@ pub fn event_name(event: &Event) -> &'static str {
AgentEvent::SessionEnded => "agent.session.ended",
AgentEvent::ProcessingEnd => "agent.processing.end",
AgentEvent::UserInput { .. } => "agent.input",
AgentEvent::AssistantTextStart => "agent.output.start",
AgentEvent::LlmRequestStarted { .. } => "agent.llm.started",
AgentEvent::LlmFirstOutput { .. } => "agent.llm.first_output",
AgentEvent::AssistantOutputReplace { .. } => "agent.output.replace",
AgentEvent::AssistantMessage { .. } => "agent.message",
AgentEvent::TextDelta { .. } => "agent.text.delta",

View file

@ -356,6 +356,12 @@ fn main() {
&[],
),
("StageProjection", "fabro_types::StageProjection", &[]),
(
"StageInferenceProjection",
"fabro_types::StageInferenceProjection",
&[],
),
("LlmOutputKind", "fabro_types::LlmOutputKind", &[]),
("PermissionLevel", "fabro_types::PermissionLevel", &[]),
(
"AgentSessionActivatedProps",

View file

@ -47,7 +47,7 @@ pub mod types {
DirtyStatus, EventEnvelope, ExecOutputTail, FailureCategory, FailureDetail,
FailureSignature, GitContext, IdpIdentity, IntegrationConnectionKind,
IntegrationConnectionState, IntegrationConnectionStatus, IntegrationProvider,
IntegrationStatus, InterviewOption, InterviewQuestionRecord,
IntegrationStatus, InterviewOption, InterviewQuestionRecord, LlmOutputKind,
McpServerDraft as CreateMcpServerRequest, McpServerProjection,
McpServerReplace as ReplaceMcpServerRequest, McpServerStatus, McpServerView as McpServer,
McpTransportView, Message, PairId, PairMessageId, PairMessageRecord, PairMessageRequest,
@ -68,10 +68,11 @@ pub mod types {
SkillsProjection, StageCompletion, StageContextWindow, StageContextWindowBreakdownItem,
StageContextWindowCategory, StageContextWindowCountMethod, StageContextWindowProjection,
StageContextWindowStaleness, StageContextWindowUnavailableReason,
StageContextWindowWarning, StageHandler, StageId, StageModelUsage, StageOutcome,
StageProjection, StageState, SubAgentProjection, SubAgentStatus, SystemActorKind,
SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, TurnId,
UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowSettings,
StageContextWindowWarning, StageHandler, StageId, StageInferenceProjection,
StageModelUsage, StageOutcome, StageProjection, StageState, SubAgentProjection,
SubAgentStatus, SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse,
TodoListProjection, TurnId, UpdateVariableRequest, UserPrincipal, Variable,
VariableListResponse, WorkflowSettings,
};
pub use crate::generated::types::*;

View file

@ -6,7 +6,7 @@ use fabro_api::types::{
AgentSkillActivationSource as ApiAgentSkillActivationSource,
AgentSkillSummary as ApiAgentSkillSummary, AgentToolCategory as ApiAgentToolCategory,
AgentToolSource as ApiAgentToolSource, AgentToolSummary as ApiAgentToolSummary,
AgentToolsAvailableProps as ApiAgentToolsAvailableProps,
AgentToolsAvailableProps as ApiAgentToolsAvailableProps, LlmOutputKind as ApiLlmOutputKind,
McpServerProjection as ApiMcpServerProjection, McpServerStatus as ApiMcpServerStatus,
ParallelBranchResult as ApiParallelBranchResult, PermissionLevel as ApiPermissionLevel,
SkillsProjection as ApiSkillsProjection, StageContextWindow as ApiStageContextWindow,
@ -17,17 +17,20 @@ use fabro_api::types::{
StageContextWindowStaleness as ApiStageContextWindowStaleness,
StageContextWindowUnavailableReason as ApiStageContextWindowUnavailableReason,
StageContextWindowWarning as ApiStageContextWindowWarning,
StageProjection as ApiStageProjection, SubAgentProjection as ApiSubAgentProjection,
SubAgentStatus as ApiSubAgentStatus, TodoListProjection as ApiTodoListProjection,
StageInferenceProjection as ApiStageInferenceProjection, StageProjection as ApiStageProjection,
SubAgentProjection as ApiSubAgentProjection, SubAgentStatus as ApiSubAgentStatus,
TodoListProjection as ApiTodoListProjection,
};
use fabro_model::{ModelId, ModelRef, ProviderId, Speed};
use fabro_types::{
ActivatedSkill, AgentControlState, AgentMcpToolSummary, AgentSkillActivationSource,
AgentSkillSummary, AgentToolCategory, AgentToolSource, AgentToolSummary,
AgentToolsAvailableProps, McpServerProjection, McpServerStatus, ParallelBranchResult,
PermissionLevel, SkillsProjection, StageContextWindow, StageContextWindowBreakdownItem,
StageContextWindowCategory, StageContextWindowCountMethod, StageContextWindowProjection,
StageContextWindowStaleness, StageContextWindowUnavailableReason, StageContextWindowWarning,
StageProjection, SubAgentProjection, SubAgentStatus, TodoListKind, TodoListProjection,
AgentToolsAvailableProps, LlmOutputKind, McpServerProjection, McpServerStatus,
ParallelBranchResult, PermissionLevel, SkillsProjection, StageContextWindow,
StageContextWindowBreakdownItem, StageContextWindowCategory, StageContextWindowCountMethod,
StageContextWindowProjection, StageContextWindowStaleness, StageContextWindowUnavailableReason,
StageContextWindowWarning, StageInferenceProjection, StageProjection, SubAgentProjection,
SubAgentStatus, TodoListKind, TodoListProjection,
};
use serde_json::json;
@ -64,6 +67,88 @@ fn stage_projection_reuses_nested_agent_state_types() {
assert_same_type::<ApiStageContextWindowUnavailableReason, StageContextWindowUnavailableReason>(
);
assert_same_type::<ApiStageContextWindowWarning, StageContextWindowWarning>();
assert_same_type::<ApiStageInferenceProjection, StageInferenceProjection>();
assert_same_type::<ApiLlmOutputKind, LlmOutputKind>();
}
#[test]
fn stage_inference_projection_matches_openapi_json_shape() {
let inference = StageInferenceProjection {
session_id: "ses_root".to_string(),
started_at: "2026-04-29T12:34:00Z".parse().unwrap(),
requested_model: ModelRef {
provider: ProviderId::new("anthropic"),
model_id: ModelId::new("claude-fable-5"),
speed: Some(Speed::Fast),
},
first_output_at: Some("2026-04-29T12:34:07Z".parse().unwrap()),
first_output_kind: Some(LlmOutputKind::Reasoning),
retries: 1,
};
let value = serde_json::to_value(&inference).unwrap();
assert_eq!(
value,
json!({
"session_id": "ses_root",
"started_at": "2026-04-29T12:34:00Z",
"requested_model": {
"provider": "anthropic",
"model_id": "claude-fable-5",
"speed": "fast"
},
"first_output_at": "2026-04-29T12:34:07Z",
"first_output_kind": "reasoning",
"retries": 1
})
);
let api_inference: ApiStageInferenceProjection = serde_json::from_value(value).unwrap();
assert_eq!(api_inference, inference);
}
#[test]
fn llm_enums_match_openapi_json_shape() {
for (kind, wire) in [
(LlmOutputKind::Reasoning, "reasoning"),
(LlmOutputKind::Text, "text"),
(LlmOutputKind::ToolCall, "tool_call"),
] {
let value = serde_json::to_value(kind).unwrap();
assert_eq!(value, json!(wire));
let api_kind: ApiLlmOutputKind = serde_json::from_value(value).unwrap();
assert_eq!(api_kind, kind);
}
}
/// A stage projection written before `inference` existed must still
/// deserialize, and must not gain a phantom open bracket.
#[test]
fn stage_projection_without_inference_round_trips() {
let value = json!({
"first_event_seq": 1,
"prompt": null,
"response": null,
"completion": null,
"provider_used": null,
"diff": null,
"script_invocation": null,
"script_timing": null,
"parallel_results": null,
"output": null,
"usage": {
"input_tokens": 0,
"output_tokens": 0,
"total_tokens": 0,
"reasoning_tokens": 0,
"cache_read_tokens": 0,
"cache_write_tokens": 0
},
"agent_control": "running",
"state": "running"
});
let stage: StageProjection = serde_json::from_value(value.clone()).unwrap();
assert!(stage.inference.is_none());
assert_eq!(serde_json::to_value(stage).unwrap(), value);
}
#[test]
@ -215,6 +300,17 @@ fn stage_projection_round_trips_representative_json() {
],
"warnings": []
},
"inference": {
"session_id": "ses_root",
"started_at": "2026-04-29T12:34:00Z",
"requested_model": {
"provider": "anthropic",
"model_id": "claude-fable-5"
},
"first_output_at": "2026-04-29T12:34:07Z",
"first_output_kind": "text",
"retries": 0
},
"agent_control": "running",
"state": "succeeded"
});

View file

@ -112,9 +112,10 @@ pub use run_blob_id::RunBlobId;
pub use run_event::{
AgentMcpToolSummary, AgentMemoryFileProps, AgentSkillActivationSource, AgentSkillSummary,
AgentToolCategory, AgentToolSource, AgentToolSummary, AgentToolsAvailableProps, EventBody,
ExecOutputTail, InterviewOption, MetadataSnapshotFailureKind, MetadataSnapshotPhase, RunEvent,
RunNoticeCode, RunNoticeLevel, RunPairEndedReason, RunPairFailedReason, RunRunnableSource,
SessionCapability, TodoCreatedProps, TodoDeletedProps, TodoUpdatedProps,
ExecOutputTail, InterviewOption, LlmOutputKind, LlmRetryPhase, MetadataSnapshotFailureKind,
MetadataSnapshotPhase, RunEvent, RunNoticeCode, RunNoticeLevel, RunPairEndedReason,
RunPairFailedReason, RunRunnableSource, SessionCapability, TodoCreatedProps, TodoDeletedProps,
TodoUpdatedProps,
};
pub use run_failure::RunFailure;
pub use run_id::{RunId, fixtures};
@ -123,8 +124,8 @@ pub use run_projection::{
PendingInterviewRecord, RunProjection, SkillsProjection, StageContextWindow,
StageContextWindowBreakdownItem, StageContextWindowCategory, StageContextWindowCountMethod,
StageContextWindowProjection, StageContextWindowStaleness, StageContextWindowUnavailableReason,
StageContextWindowWarning, StageModelUsage, StageProjection, SubAgentProjection,
SubAgentStatus, first_event_seq,
StageContextWindowWarning, StageInferenceProjection, StageModelUsage, StageProjection,
SubAgentProjection, SubAgentStatus, first_event_seq,
};
pub use run_sandbox::{
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,

View file

@ -290,6 +290,33 @@ pub struct AgentCompactionCompletedProps {
pub visit: u32,
}
/// Which loop produced the `attempt` index on an `agent.llm.retry` event.
///
/// `attempt` is a 0-based counter fed by two independent loops: the retry
/// policy inside `open_stream_with_retry` (`Open`) and the stream-consume
/// loop that replays a turn whose stream broke or ended without a finish
/// event (`Consume`). Without this discriminator a reader cannot tell which
/// counter an index belongs to.
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Hash,
Serialize,
Deserialize,
Display,
EnumString,
IntoStaticStr,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum LlmRetryPhase {
Open,
Consume,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentLlmRetryProps {
pub provider: String,
@ -297,9 +324,58 @@ pub struct AgentLlmRetryProps {
pub attempt: usize,
pub delay_secs: f64,
pub error: Value,
/// Which retry loop `attempt` counts. Absent on events stored before the
/// discriminator existed.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub phase: Option<LlmRetryPhase>,
pub visit: u32,
}
/// Kind of output a provider produced first for an inference attempt.
///
/// Observed, never inferred: a turn that opens with a tool call emits no text
/// or reasoning delta, so all three variants are required for the first-output
/// edge to fire on every turn.
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Hash,
Serialize,
Deserialize,
Display,
EnumString,
IntoStaticStr,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum LlmOutputKind {
Reasoning,
Text,
ToolCall,
}
/// An inference request is about to be dispatched for this round.
///
/// `requested_model` is the requested target from the session's provider
/// profile. Failover can re-target mid-stage, so `agent.message` remains
/// authoritative for what actually answered.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentLlmStartedProps {
pub requested_model: ModelRef,
pub visit: u32,
}
/// The provider produced its first output for the current attempt.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentLlmFirstOutputProps {
/// Which kind of output arrived first — observed, not inferred.
pub kind: LlmOutputKind,
pub visit: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentSubSpawnedProps {
pub agent_id: String,

View file

@ -234,6 +234,10 @@ pub enum EventBody {
AgentCompactionStarted(AgentCompactionStartedProps),
#[serde(rename = "agent.compaction.completed")]
AgentCompactionCompleted(AgentCompactionCompletedProps),
#[serde(rename = "agent.llm.started")]
AgentLlmStarted(AgentLlmStartedProps),
#[serde(rename = "agent.llm.first_output")]
AgentLlmFirstOutput(AgentLlmFirstOutputProps),
#[serde(rename = "agent.llm.retry")]
AgentLlmRetry(AgentLlmRetryProps),
#[serde(rename = "agent.sub.spawned")]
@ -502,6 +506,8 @@ impl EventBody {
Self::AgentSteerDropped(_) => "agent.steer.dropped",
Self::AgentCompactionStarted(_) => "agent.compaction.started",
Self::AgentCompactionCompleted(_) => "agent.compaction.completed",
Self::AgentLlmStarted(_) => "agent.llm.started",
Self::AgentLlmFirstOutput(_) => "agent.llm.first_output",
Self::AgentLlmRetry(_) => "agent.llm.retry",
Self::AgentSubSpawned(_) => "agent.sub.spawned",
Self::AgentSubCompleted(_) => "agent.sub.completed",
@ -672,6 +678,8 @@ fn is_known_event_name(event: &str) -> bool {
| "agent.steer.dropped"
| "agent.compaction.started"
| "agent.compaction.completed"
| "agent.llm.started"
| "agent.llm.first_output"
| "agent.llm.retry"
| "agent.sub.spawned"
| "agent.sub.completed"

View file

@ -10,9 +10,9 @@ use crate::run_event::{AgentSessionActivatedProps, StagePromptProps};
use crate::{
AgentBackend, AgentMcpToolSummary, AgentSkillActivationSource, AgentSkillSummary,
AgentToolSummary, BilledTokenCounts, Checkpoint, Conclusion, InterviewQuestionRecord,
InvalidTransition, ModelRef, PermissionLevel, PullRequestLink, RunApproval, RunControlAction,
RunDiff, RunId, RunSandbox, RunSpec, RunStatus, RunTiming, StageCompletion, StageHandler,
StageId, StageState, StageTiming, StartRecord, TodoListProjection,
InvalidTransition, LlmOutputKind, ModelRef, PermissionLevel, PullRequestLink, RunApproval,
RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus, RunTiming, StageCompletion,
StageHandler, StageId, StageState, StageTiming, StartRecord, TodoListProjection,
};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
@ -381,11 +381,41 @@ pub struct StageProjection {
pub mcp_servers: Vec<McpServerProjection>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context_window: Option<StageContextWindowProjection>,
/// Open inference bracket for this stage, if the event log contains one.
///
/// `Some` means exactly *"an `agent.llm.started` was recorded and no
/// closing event has been seen"* — not "the model is computing right
/// now". A worker killed mid-turn leaves the bracket open, which is the
/// truthful statement of what the log knows. `watchdog.timeout` remains
/// the authority on whether a run is actually stuck.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub inference: Option<StageInferenceProjection>,
#[serde(default)]
pub agent_control: AgentControlState,
pub state: StageState,
}
/// One open inference bracket: a dispatched LLM request that has not yet
/// produced a message, error, or interrupt.
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct StageInferenceProjection {
/// Copied from the envelope `RunEvent.session_id` when the bracket opens.
/// Every later transition is gated on it so a sub-agent's rounds cannot
/// overwrite the root session's bracket.
pub session_id: String,
pub started_at: DateTime<Utc>,
/// Provider and model the request was *sent to*. Failover can re-target,
/// so `StageProjection::model` stays authoritative for what answered.
pub requested_model: ModelRef,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub first_output_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub first_output_kind: Option<LlmOutputKind>,
/// Attempts that failed and restarted within this bracket.
#[serde(default)]
pub retries: u32,
}
#[derive(
Debug,
Clone,
@ -485,6 +515,7 @@ impl StageProjection {
agent_tools: Vec::new(),
mcp_servers: Vec::new(),
context_window: None,
inference: None,
agent_control: AgentControlState::default(),
provider_used: None,
diff: None,
@ -601,6 +632,16 @@ impl RunProjection {
self.stages.iter()
}
/// Mutable counterpart of [`Self::iter_stages_unordered`].
///
/// Use this only for order-independent mutation. Presentation and
/// serialization callers should use [`Self::iter_stages_mut`] instead.
pub fn iter_stages_unordered_mut(
&mut self,
) -> impl Iterator<Item = (&StageId, &mut StageProjection)> {
self.stages.iter_mut()
}
/// Iterate stages in `first_event_seq` order (the chronological order in
/// which each stage's first lifecycle event was recorded). Internal
/// storage is a `HashMap`, so presentation callers sort through this

View file

@ -197,6 +197,7 @@ models/interview-option.ts
models/interview-provider-settings.ts
models/interview-question-record.ts
models/link-run-pull-request-request.ts
models/llm-output-kind.ts
models/log-destination.ts
models/manifest-args.ts
models/manifest-config.ts
@ -479,6 +480,7 @@ models/stage-context-window-unavailable-reason.ts
models/stage-context-window-warning.ts
models/stage-context-window.ts
models/stage-handler.ts
models/stage-inference-projection.ts
models/stage-model-usage.ts
models/stage-outcome.ts
models/stage-projection.ts

View file

@ -167,6 +167,7 @@ export * from './interview-option';
export * from './interview-provider-settings';
export * from './interview-question-record';
export * from './link-run-pull-request-request';
export * from './llm-output-kind';
export * from './log-destination';
export * from './manifest-args';
export * from './manifest-config';
@ -449,6 +450,7 @@ export * from './stage-context-window-staleness';
export * from './stage-context-window-unavailable-reason';
export * from './stage-context-window-warning';
export * from './stage-handler';
export * from './stage-inference-projection';
export * from './stage-model-usage';
export * from './stage-outcome';
export * from './stage-projection';

View file

@ -0,0 +1,27 @@
/* tslint:disable */
/* eslint-disable */
/**
* Fabro Run API
* HTTP API for managing Fabro workflow run executions.
*
* The version of the OpenAPI document: 0.1.0
*
*
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
* https://openapi-generator.tech
* Do not edit the class manually.
*/
/**
* Kind of output a provider produced first for an inference attempt. Observed, never inferred.
*/
export const LlmOutputKind = {
REASONING: 'reasoning',
TEXT: 'text',
TOOL_CALL: 'tool_call'
} as const;
export type LlmOutputKind = typeof LlmOutputKind[keyof typeof LlmOutputKind];

View file

@ -0,0 +1,48 @@
/* tslint:disable */
/* eslint-disable */
/**
* Fabro Run API
* HTTP API for managing Fabro workflow run executions.
*
* The version of the OpenAPI document: 0.1.0
*
*
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
* https://openapi-generator.tech
* Do not edit the class manually.
*/
// May contain unused imports in some cases
// @ts-ignore
import type { BillingModelRef } from './billing-model-ref';
// May contain unused imports in some cases
// @ts-ignore
import type { LlmOutputKind } from './llm-output-kind';
/**
* One open inference bracket: a dispatched LLM request that has not yet produced a message, error, or interrupt. Carries no usage or cost — none exists until the turn completes.
*/
export interface StageInferenceProjection {
/**
* Agent session that opened the bracket, copied from the event envelope. Transitions are gated on it so sub-agent rounds cannot overwrite the root session\'s bracket.
*/
'session_id': string;
/**
* When the request was dispatched.
*/
'started_at': string;
/**
* Provider and model the request was sent to. Failover can re-target, so `StageProjection.model` stays authoritative for what answered.
*/
'requested_model': BillingModelRef;
/**
* When the provider produced its first output, if it has.
*/
'first_output_at'?: string | null;
'first_output_kind'?: LlmOutputKind | null;
/**
* Attempts that failed and restarted within this bracket.
*/
'retries': number;
}

View file

@ -48,6 +48,9 @@ import type { StageCompletion } from './stage-completion';
import type { StageContextWindowProjection } from './stage-context-window-projection';
// May contain unused imports in some cases
// @ts-ignore
import type { StageInferenceProjection } from './stage-inference-projection';
// May contain unused imports in some cases
// @ts-ignore
import type { StageModelUsage } from './stage-model-usage';
// May contain unused imports in some cases
// @ts-ignore
@ -114,6 +117,7 @@ export interface StageProjection {
*/
'mcp_servers'?: Array<McpServerProjection>;
'context_window'?: StageContextWindowProjection | null;
'inference'?: StageInferenceProjection | null;
/**
* Whether the agent is executing normally or waiting for steering after an interrupt.
*/