diff --git a/apps/fabro-web/app/components/stage-inference-indicator.test.tsx b/apps/fabro-web/app/components/stage-inference-indicator.test.tsx new file mode 100644 index 000000000..45da56ae4 --- /dev/null +++ b/apps/fabro-web/app/components/stage-inference-indicator.test.tsx @@ -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 { + 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( + , + ); + }); + 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"); + }); +}); diff --git a/apps/fabro-web/app/components/stage-inference-indicator.tsx b/apps/fabro-web/app/components/stage-inference-indicator.tsx new file mode 100644 index 000000000..a329f3fa0 --- /dev/null +++ b/apps/fabro-web/app/components/stage-inference-indicator.tsx @@ -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.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 ( +

+ Model request opened {formatAbsoluteTs(inference.started_at)}, never + completed +

+ ); + } + + 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 ( +

+ + + + +

+ ); +} diff --git a/apps/fabro-web/app/lib/query-keys.test.ts b/apps/fabro-web/app/lib/query-keys.test.ts index 6faa4e158..ac441b998 100644 --- a/apps/fabro-web/app/lib/query-keys.test.ts +++ b/apps/fabro-web/app/lib/query-keys.test.ts @@ -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"), + ]); }); }); diff --git a/apps/fabro-web/app/lib/run-actions.test.ts b/apps/fabro-web/app/lib/run-actions.test.ts index ab67c0e90..0547f9398 100644 --- a/apps/fabro-web/app/lib/run-actions.test.ts +++ b/apps/fabro-web/app/lib/run-actions.test.ts @@ -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); diff --git a/apps/fabro-web/app/lib/run-actions.ts b/apps/fabro-web/app/lib/run-actions.ts index 46fc1efdd..a45034088 100644 --- a/apps/fabro-web/app/lib/run-actions.ts +++ b/apps/fabro-web/app/lib/run-actions.ts @@ -41,7 +41,7 @@ const CANCELABLE_STATUSES = new Set([ "blocked", ]); -const ARCHIVABLE_STATUSES = new Set([ +const TERMINAL_RUN_STATUSES = new Set([ "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 | 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 { diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 55af2aefb..11544b14d 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -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", () => { diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index 32f8edfb5..d0d4f4b4d 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -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 ? [ diff --git a/apps/fabro-web/app/routes/run-stages.test.ts b/apps/fabro-web/app/routes/run-stages.test.ts index 51684a895..2238aaf74 100644 --- a/apps/fabro-web/app/routes/run-stages.test.ts +++ b/apps/fabro-web/app/routes/run-stages.test.ts @@ -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, { diff --git a/apps/fabro-web/app/routes/run-stages.tsx b/apps/fabro-web/app/routes/run-stages.tsx index de565aa6d..48111ce23 100644 --- a/apps/fabro-web/app/routes/run-stages.tsx +++ b/apps/fabro-web/app/routes/run-stages.tsx @@ -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(); 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(() => { return (stageEventsQuery.data ?? []).filter( @@ -2239,6 +2259,11 @@ function RunStageActivityStage({

)} + + ); diff --git a/docs/internal/events.md b/docs/internal/events.md index 0a51e4e13..196925e94 100644 --- a/docs/internal/events.md +++ b/docs/internal/events.md @@ -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": { diff --git a/docs/internal/fabro-event-schema-v2-concrete-shape.md b/docs/internal/fabro-event-schema-v2-concrete-shape.md index 78134f619..6d9a3c28d 100644 --- a/docs/internal/fabro-event-schema-v2-concrete-shape.md +++ b/docs/internal/fabro-event-schema-v2-concrete-shape.md @@ -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. diff --git a/docs/internal/fabro-event-schema-v2-proposal.md b/docs/internal/fabro-event-schema-v2-proposal.md index ee8714766..0ce5f9bd1 100644 --- a/docs/internal/fabro-event-schema-v2-proposal.md +++ b/docs/internal/fabro-event-schema-v2-proposal.md @@ -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 diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 1dcbca2e5..b028d89c0 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -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 diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs index 317bfa449..1febfc5b7 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs @@ -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 { 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 { 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 { 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(), diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index 28e62c2b2..b18683354 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -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, diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs index c7c8f8ef6..e71072fb2 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/stage_display.rs @@ -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, pub(super) compaction_bar: Option, + /// 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, } 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(); diff --git a/lib/components/fabro-agent/src/session.rs b/lib/components/fabro-agent/src/session.rs index bb7bd56d7..3854f3ece 100644 --- a/lib/components/fabro-agent/src/session.rs +++ b/lib/components/fabro-agent/src/session.rs @@ -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, 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 { + 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> = 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) -> 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:?}")); diff --git a/lib/components/fabro-agent/src/types.rs b/lib/components/fabro-agent/src/types.rs index 129cfc8a0..c5bd99deb 100644 --- a/lib/components/fabro-agent/src/types.rs +++ b/lib/components/fabro-agent/src/types.rs @@ -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 { diff --git a/lib/components/fabro-llm/src/codec/anthropic_messages/stream.rs b/lib/components/fabro-llm/src/codec/anthropic_messages/stream.rs index 96b4830a1..1f6a55fb8 100644 --- a/lib/components/fabro-llm/src/codec/anthropic_messages/stream.rs +++ b/lib/components/fabro-llm/src/codec/anthropic_messages/stream.rs @@ -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 { diff --git a/lib/components/fabro-llm/src/codec/bedrock_converse/stream.rs b/lib/components/fabro-llm/src/codec/bedrock_converse/stream.rs index 56c2db659..0facf2583 100644 --- a/lib/components/fabro-llm/src/codec/bedrock_converse/stream.rs +++ b/lib/components/fabro-llm/src/codec/bedrock_converse/stream.rs @@ -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", diff --git a/lib/components/fabro-llm/src/codec/gemini_generate/stream.rs b/lib/components/fabro-llm/src/codec/gemini_generate/stream.rs index 7e13ddb78..51c39eb99 100644 --- a/lib/components/fabro-llm/src/codec/gemini_generate/stream.rs +++ b/lib/components/fabro-llm/src/codec/gemini_generate/stream.rs @@ -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 { 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] diff --git a/lib/components/fabro-llm/src/codec/mod.rs b/lib/components/fabro-llm/src/codec/mod.rs index 5a7759b22..97bf842b9 100644 --- a/lib/components/fabro-llm/src/codec/mod.rs +++ b/lib/components/fabro-llm/src/codec/mod.rs @@ -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, Error>; /// Byte-stream-end hook. Semantics are per-decoder, not shared: diff --git a/lib/components/fabro-llm/src/codec/openai_responses/stream.rs b/lib/components/fabro-llm/src/codec/openai_responses/stream.rs index b541a47cb..9d0a06d4d 100644 --- a/lib/components/fabro-llm/src/codec/openai_responses/stream.rs +++ b/lib/components/fabro-llm/src/codec/openai_responses/stream.rs @@ -95,7 +95,6 @@ pub(super) struct SseAccumulator { message_items: Vec, usage: TokenCounts, finish_reason: FinishReason, - emitted_start: bool, emitted_text_start: bool, emitted_reasoning_start: bool, rate_limit: Option, @@ -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, 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, diff --git a/lib/components/fabro-llm/src/providers/bedrock/mod.rs b/lib/components/fabro-llm/src/providers/bedrock/mod.rs index deca69bd0..6f10dadb2 100644 --- a/lib/components/fabro-llm/src/providers/bedrock/mod.rs +++ b/lib/components/fabro-llm/src/providers/bedrock/mod.rs @@ -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, - pending: VecDeque, - done: bool, + response: fabro_http::Response, + frames: FrameDecoder, + decoder: Box, + pending: VecDeque>, + done: bool, /// `finish()` already drained. - finished: bool, - timeout: Option, + finished: bool, + /// [`StreamEvent::StreamStart`] already emitted for this stream. + stream_started: bool, + timeout: Option, } /// 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; + } } } } diff --git a/lib/components/fabro-llm/src/transport.rs b/lib/components/fabro-llm/src/transport.rs index 1ed1ee558..8fc62a32f 100644 --- a/lib/components/fabro-llm/src/transport.rs +++ b/lib/components/fabro-llm/src/transport.rs @@ -247,12 +247,14 @@ pub(crate) async fn stream_via_http( struct StreamLoop { decoder: Box, line_reader: LineReader, - /// Events decoded but not yet yielded. - pending: VecDeque, + /// Events or decoder errors not yet yielded. + pending: VecDeque>, /// 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, diff --git a/lib/components/fabro-llm/tests/it/support.rs b/lib/components/fabro-llm/tests/it/support.rs index c2d30c58c..3c1e3a62c 100644 --- a/lib/components/fabro-llm/tests/it/support.rs +++ b/lib/components/fabro-llm/tests/it/support.rs @@ -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 { diff --git a/lib/components/fabro-llm/tests/it/wire/anthropic.rs b/lib/components/fabro-llm/tests/it/wire/anthropic.rs index 9e790f2d4..ba73edaef 100644 --- a/lib/components/fabro-llm/tests/it/wire/anthropic.rs +++ b/lib/components/fabro-llm/tests/it/wire/anthropic.rs @@ -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); } diff --git a/lib/components/fabro-llm/tests/it/wire/gemini.rs b/lib/components/fabro-llm/tests/it/wire/gemini.rs index 6750fcfb0..2d1cc2e5e 100644 --- a/lib/components/fabro-llm/tests/it/wire/gemini.rs +++ b/lib/components/fabro-llm/tests/it/wire/gemini.rs @@ -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); } diff --git a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs index 90227f9c1..1294e2e8a 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs @@ -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(&[ diff --git a/lib/components/fabro-llm/tests/it/wire/openai_responses.rs b/lib/components/fabro-llm/tests/it/wire/openai_responses.rs index 1a3f2fb6a..079a05455 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_responses.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_responses.rs @@ -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(&[ diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_reasoning_deltas.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_reasoning_deltas.snap index ee27dbea7..9aaca3762 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_reasoning_deltas.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_reasoning_deltas.snap @@ -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 diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_events.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_events.snap index 0bb05a55b..34dd07b3f 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_events.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_text_happy_path_events.snap @@ -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 diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_tool_call_deltas.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_tool_call_deltas.snap index fdce5d082..e1bb0f279 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_tool_call_deltas.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_tool_call_deltas.snap @@ -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": { diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_usage_openrouter_cost.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_usage_openrouter_cost.snap index a560e0d83..d00f07703 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_usage_openrouter_cost.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_usage_openrouter_cost.snap @@ -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 diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_or_content_synthesizes_nothing.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_or_content_synthesizes_nothing.snap index 082fdf19c..b54da4946 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_or_content_synthesizes_nothing.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_or_content_synthesizes_nothing.snap @@ -2,4 +2,8 @@ source: lib/components/fabro-llm/tests/it/wire/openai_compatible.rs expression: rendered --- -[] +[ + { + "type": "stream_start" + } +] diff --git a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_synthesizes_finish_when_content_started.snap b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_synthesizes_finish_when_content_started.snap index 85127a104..18f1feb19 100644 --- a/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_synthesizes_finish_when_content_started.snap +++ b/lib/components/fabro-llm/tests/it/wire/snapshots/it__wire__openai_compatible__stream_without_done_synthesizes_finish_when_content_started.snap @@ -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 diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index da9152e00..9ab1e2b2d 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -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, +) { + 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> { + 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()); + } + } } diff --git a/lib/components/fabro-store/src/slate/projection_cache.rs b/lib/components/fabro-store/src/slate/projection_cache.rs index cd9efd482..7cec9a204 100644 --- a/lib/components/fabro-store/src/slate/projection_cache.rs +++ b/lib/components/fabro-store/src/slate/projection_cache.rs @@ -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, + parent_id: Option, + ) { + 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 { 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) } diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index f6e2de0f1..5cef4e581 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -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 { .. } diff --git a/lib/components/fabro-workflow/src/event/names.rs b/lib/components/fabro-workflow/src/event/names.rs index cd04b08b1..cc6ad1a93 100644 --- a/lib/components/fabro-workflow/src/event/names.rs +++ b/lib/components/fabro-workflow/src/event/names.rs @@ -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", diff --git a/lib/foundation/fabro-api/build.rs b/lib/foundation/fabro-api/build.rs index 03cd00ea5..fc943e99d 100644 --- a/lib/foundation/fabro-api/build.rs +++ b/lib/foundation/fabro-api/build.rs @@ -356,6 +356,12 @@ fn main() { &[], ), ("StageProjection", "fabro_types::StageProjection", &[]), + ( + "StageInferenceProjection", + "fabro_types::StageInferenceProjection", + &[], + ), + ("LlmOutputKind", "fabro_types::LlmOutputKind", &[]), ("PermissionLevel", "fabro_types::PermissionLevel", &[]), ( "AgentSessionActivatedProps", diff --git a/lib/foundation/fabro-api/src/lib.rs b/lib/foundation/fabro-api/src/lib.rs index 0aeed27a7..4312acb91 100644 --- a/lib/foundation/fabro-api/src/lib.rs +++ b/lib/foundation/fabro-api/src/lib.rs @@ -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::*; diff --git a/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs b/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs index c386ab03d..5faed1bd3 100644 --- a/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs +++ b/lib/foundation/fabro-api/tests/stage_projection_round_trip.rs @@ -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::( ); assert_same_type::(); + assert_same_type::(); + assert_same_type::(); +} + +#[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" }); diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index 1b19c9248..68f5644a2 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -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, diff --git a/lib/foundation/fabro-types/src/run_event/agent.rs b/lib/foundation/fabro-types/src/run_event/agent.rs index 604843011..8b27a3be1 100644 --- a/lib/foundation/fabro-types/src/run_event/agent.rs +++ b/lib/foundation/fabro-types/src/run_event/agent.rs @@ -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, 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, diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index b79f7758b..d4c8e55c1 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -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" diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index 00065b0d6..027cce2fc 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -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, #[serde(default, skip_serializing_if = "Option::is_none")] pub context_window: Option, + /// 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, #[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, + /// 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>, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub first_output_kind: Option, + /// 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 { + 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 diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 02add4f1d..e3a1a7c0e 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -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 diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 474bf0b2b..5cd29601d 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -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'; diff --git a/lib/packages/fabro-api-client/src/models/llm-output-kind.ts b/lib/packages/fabro-api-client/src/models/llm-output-kind.ts new file mode 100644 index 000000000..1c9663257 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/llm-output-kind.ts @@ -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]; diff --git a/lib/packages/fabro-api-client/src/models/stage-inference-projection.ts b/lib/packages/fabro-api-client/src/models/stage-inference-projection.ts new file mode 100644 index 000000000..13a08869b --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/stage-inference-projection.ts @@ -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; +} diff --git a/lib/packages/fabro-api-client/src/models/stage-projection.ts b/lib/packages/fabro-api-client/src/models/stage-projection.ts index 87c5b169a..57fd342ca 100644 --- a/lib/packages/fabro-api-client/src/models/stage-projection.ts +++ b/lib/packages/fabro-api-client/src/models/stage-projection.ts @@ -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; 'context_window'?: StageContextWindowProjection | null; + 'inference'?: StageInferenceProjection | null; /** * Whether the agent is executing normally or waiting for steering after an interrupt. */