From f22cf18c86c6daea50df4181eebefcf67c289d16 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 00:59:10 -0400 Subject: [PATCH] Render a Petri run's detail views from its projection and stream The web app reads a Petri run through the run stream (`useRunStream` pages `GET /runs/{id}/events` by `after`) and the projection, keyed on `RunSpec.engine`, beside the legacy event path. `lib/petri-stream.ts` holds the pure derivations `VIEWS.md` maps: the stage label of an item's subject (lowering nodes skipped), the interview pairs from `parsed.question` and the delivered `control.requested` with the `interview.answered` principal, the run phases from the platform lifecycle records, the edge a `route.applied` took, the fork's branches and the fan-in transcript from the projection, the stage context from the final `step.finished`, the Pebble envelopes of a step, and the debug rows the listings show. The events route lists stream rows (named `.` or by the platform kind, with the stage beside them and the raw item in the details panel) and the waterfall takes its phases from the lifecycle records. The stages route builds a Petri stage's turns from the projection's prompt and response and the step's envelopes, its debug tab from the stage's items, and hands the human, conditional, parallel and fan-in renderers the derived data instead of events. The overview lists the platform records (checkpoints with their commit, pull request, notices). The SSE subscription invalidates SWR keys from a stream item's event name or platform kind, ending on the terminal lifecycle record, and the cross-tab dedupe keys a stream item by its run and `stream_seq`. Co-Authored-By: Claude Fable 5.1 --- .../app/components/event-debug-helpers.tsx | 12 +- apps/fabro-web/app/components/event-debug.tsx | 32 +- .../app/components/platform-records-panel.tsx | 96 +++ .../app/components/run-waterfall.tsx | 17 +- .../stage-renderers/conditional-decision.tsx | 9 +- .../stage-renderers/fan-in-results.tsx | 10 +- .../app/components/stage-renderers/helpers.ts | 2 +- .../components/stage-renderers/human-qa.tsx | 8 +- .../stage-renderers/parallel-children.tsx | 10 +- .../stage-renderers/stage-summary.tsx | 18 +- apps/fabro-web/app/lib/cross-tab-sse.ts | 5 + apps/fabro-web/app/lib/petri-stream.ts | 692 ++++++++++++++++++ apps/fabro-web/app/lib/queries.ts | 69 +- apps/fabro-web/app/lib/query-keys.ts | 2 + apps/fabro-web/app/lib/run-events.ts | 163 +++++ apps/fabro-web/app/routes/run-events.tsx | 254 ++++++- .../app/routes/run-overview.test.tsx | 2 + apps/fabro-web/app/routes/run-overview.tsx | 4 +- apps/fabro-web/app/routes/run-stages.tsx | 306 +++++++- 19 files changed, 1640 insertions(+), 71 deletions(-) create mode 100644 apps/fabro-web/app/components/platform-records-panel.tsx create mode 100644 apps/fabro-web/app/lib/petri-stream.ts diff --git a/apps/fabro-web/app/components/event-debug-helpers.tsx b/apps/fabro-web/app/components/event-debug-helpers.tsx index 1a55035bd..c32fc1d61 100644 --- a/apps/fabro-web/app/components/event-debug-helpers.tsx +++ b/apps/fabro-web/app/components/event-debug-helpers.tsx @@ -5,7 +5,9 @@ export type DebugCategory = | "command" | "lifecycle" | "human" - | "system"; + | "system" + | "petri" + | "platform"; export const DEBUG_CATEGORIES: readonly DebugCategory[] = [ "agent", @@ -13,6 +15,8 @@ export const DEBUG_CATEGORIES: readonly DebugCategory[] = [ "lifecycle", "human", "system", + "petri", + "platform", ] as const; const PREFIX_TO_CATEGORY: Record = { @@ -34,6 +38,8 @@ const CATEGORY_LABEL: Record = { lifecycle: "Lifecycle", human: "Human", system: "System", + petri: "Petri", + platform: "Platform", }; const CATEGORY_TONE: Record = { @@ -42,6 +48,8 @@ const CATEGORY_TONE: Record = { lifecycle: "bg-amber/15 text-amber", human: "bg-coral/15 text-coral", system: "bg-overlay-strong text-fg-3", + petri: "bg-teal-500/15 text-teal-500", + platform: "bg-amber/15 text-amber", }; const CATEGORY_COLOR: Record = { @@ -50,6 +58,8 @@ const CATEGORY_COLOR: Record = { lifecycle: "var(--color-amber)", human: "var(--color-coral)", system: "var(--color-ice-300)", + petri: "var(--color-teal-500)", + platform: "var(--color-amber)", }; export function debugCategory(eventName: string | null | undefined): DebugCategory { diff --git a/apps/fabro-web/app/components/event-debug.tsx b/apps/fabro-web/app/components/event-debug.tsx index b50ced2f9..e23fd59d3 100644 --- a/apps/fabro-web/app/components/event-debug.tsx +++ b/apps/fabro-web/app/components/event-debug.tsx @@ -12,7 +12,6 @@ import { FunnelIcon, MagnifyingGlassIcon, } from "@heroicons/react/16/solid"; -import type { EventEnvelope } from "@qltysh/fabro-api-client"; import { Tooltip } from "./ui"; import { formatAbsoluteTs } from "../lib/format"; @@ -28,19 +27,36 @@ import { import { FloatingTooltip } from "./floating-tooltip"; import { useWindowEvent } from "../hooks/effects"; +/** + * What a debug row needs of an event: a legacy `EventEnvelope`, or a Petri + * run stream item as `debugRowsFromStream` shapes it (with its own + * category, since Petri's `.` names map to none of the + * legacy prefixes). + */ +export interface DebugRowLike { + seq: number; + event?: string | null; + ts: string; + category?: DebugCategory; +} + +export function debugRowCategory(row: DebugRowLike): DebugCategory { + return row.category ?? debugCategory(row.event); +} + export function DebugEventRow({ event, runStart, selected, onSelect, }: { - event: EventEnvelope; + event: DebugRowLike; runStart: string | undefined; selected: boolean; onSelect: () => void; }) { const eventName = event.event ?? ""; - const category = debugCategory(eventName); + const category = debugRowCategory(event); return ( + )} + + {all.length > 0 && ( + + {isFiltering + ? `${filtered.length.toLocaleString()} of ${all.length.toLocaleString()} items` + : `${all.length.toLocaleString()} items`} + + )} + + + +
+ {all.length === 0 ? ( +
+ +
+ ) : filtered.length === 0 ? ( +
+ No events match these filters. +
+ ) : ( + filtered.map((row) => ( + setOpenSeq(row.seq)} + /> + )) + )} +
+ + + setOpenSeq(null)} /> + + ); +} + +/** A debug row with the stage the item belongs to beside its name. */ +function StreamEventRow({ + row, + runStart, + selected, + onSelect, +}: { + row: DebugRow; + runStart: string | undefined; + selected: boolean; + onSelect: () => void; +}) { + return ( +
+ + {row.stageLabel && ( + + {row.stageLabel} + + )} +
+ ); +} + function errorMessage(error: unknown): string | undefined { return error instanceof Error ? error.message : undefined; } diff --git a/apps/fabro-web/app/routes/run-overview.test.tsx b/apps/fabro-web/app/routes/run-overview.test.tsx index ae5f159b9..cda806718 100644 --- a/apps/fabro-web/app/routes/run-overview.test.tsx +++ b/apps/fabro-web/app/routes/run-overview.test.tsx @@ -22,6 +22,8 @@ mock.module("../lib/queries", () => ({ }), useRunGraphSource: () => ({ data: undefined }), useRunStageEvents: () => ({ data: [] }), + useRunState: () => ({ data: undefined }), + useRunStream: () => ({ data: undefined }), })); mock.module("../components/run-summary-panel", () => ({ diff --git a/apps/fabro-web/app/routes/run-overview.tsx b/apps/fabro-web/app/routes/run-overview.tsx index 0150d19ea..3a6102145 100644 --- a/apps/fabro-web/app/routes/run-overview.tsx +++ b/apps/fabro-web/app/routes/run-overview.tsx @@ -3,6 +3,7 @@ import { useNavigate, useParams } from "react-router"; import { ApiError } from "../lib/api-client"; import { useRun, useRunGraph, useRunGraphSource, useRunStages } from "../lib/queries"; import { FloatingTooltip } from "../components/floating-tooltip"; +import { PlatformRecordsPanel } from "../components/platform-records-panel"; import { RunSummaryPanel } from "../components/run-summary-panel"; import { StagePopover } from "../components/stage-popover"; import { StageSidebar } from "../components/stage-sidebar"; @@ -150,8 +151,9 @@ export default function RunOverview() {
-
+
+
{graphSvg === undefined && graphQuery.isLoading ? (
diff --git a/apps/fabro-web/app/routes/run-stages.tsx b/apps/fabro-web/app/routes/run-stages.tsx index 1613096e8..9994760ef 100644 --- a/apps/fabro-web/app/routes/run-stages.tsx +++ b/apps/fabro-web/app/routes/run-stages.tsx @@ -23,6 +23,7 @@ import { EventSearchInput, MultiSelectFilter, ThreadDnaStrip, + debugRowCategory, threadSelectionId, threadSelectionsEqual, } from "../components/event-debug"; @@ -33,6 +34,7 @@ import { type DebugCategory, } from "../components/event-debug-helpers"; import type { + EventDisplayPayload, ThreadDnaItem, ThreadDnaSelection, } from "../components/event-debug"; @@ -51,7 +53,13 @@ import { } from "../components/ui"; import { ConditionalDecision } from "../components/stage-renderers/conditional-decision"; import { FanInResults } from "../components/stage-renderers/fan-in-results"; -import { extractStageContext } from "../components/stage-renderers/helpers"; +import { + extractStageContext, + type EdgeSelection, + type HumanInterviewPair, + type ParallelOverview, + type ReducerTranscript, +} from "../components/stage-renderers/helpers"; import { HumanQA } from "../components/stage-renderers/human-qa"; import { ParallelChildren } from "../components/stage-renderers/parallel-children"; import { @@ -71,6 +79,21 @@ import { } from "../lib/format"; import { costSourceTag, hasUsage, usageTokenBuckets } from "../lib/usage"; import { plural } from "../lib/plural"; +import { + agentEnvelopesOf, + commandOutcomeOf, + commandScriptOf, + debugRowSearchText, + debugRowsFromStream, + extractPetriStageContext, + findPetriEdgeForStage, + isPetriRun, + itemsForStage, + parallelOverviewFromProjection, + parsePetriInterviewPairs, + reducerTranscriptFromProjection, + type DebugRow, +} from "../lib/petri-stream"; import { useRun, useRunEventsList, @@ -79,6 +102,7 @@ import { useRunStageLog, useRunStages, useRunState, + useRunStream, } from "../lib/queries"; import { STAGE_ACTIVITY_EVENT_TYPES, @@ -98,6 +122,9 @@ import { import type { EventEnvelope, ReasoningOutput, + RunProjection, + RunStreamItem, + StageProjection, StageHandler, StageModelUsage, Usage, @@ -507,6 +534,169 @@ export function eventsToActivity( return buildStageActivity(events, stageId).turns; } +/** The stream and projection of a Petri run, threaded into the stage views. */ +export interface PetriRunData { + stream: RunStreamItem[]; + projection: RunProjection; +} + +/** + * The turns of a stage that ran on Petri: the prompt the projection holds, + * the Pebble envelopes the stage's step recorded (assistant messages, tool + * calls, interrupts), and, when no envelope carried the answer, the + * projection's response as the one assistant turn. A command stage is one + * command turn from its `step.started` and final `step.finished`. + */ +export function buildPetriStageActivity( + items: RunStreamItem[], + stage: StageProjection | undefined, + renderer: StageRenderer, +): StageActivity { + const turns: TurnType[] = []; + const pendingTools = new Map(); + const firstTs = items[0] + ? new Date(items[0].recorded_at).toISOString() + : (stage?.started_at ?? new Date(0).toISOString()); + const startTs = stage?.started_at ?? firstTs; + + if (renderer === "command") { + const started = items.some( + (item) => item.kind === "petri" && getString(getObject(getObject(item.item, "record"), "body"), "event") === "step.started", + ); + if (started || stage) { + const outcome = commandOutcomeOf(items); + const running = + stage?.state === "running" || stage?.state === "retrying"; + turns.push({ + kind: "command", + ts: startTs, + script: commandScriptOf(items) ?? "", + running, + exitCode: outcome.exitCode, + durationMs: outcome.durationMs || (stage?.timing?.wall_time_ms ?? 0), + outputBytes: stage?.output_bytes ?? 0, + }); + } + return { turns, pendingTools: [] }; + } + + if (stage?.prompt) { + turns.push({ kind: "system", ts: startTs, content: stage.prompt }); + } + let sawAssistantMessage = false; + for (const envelope of agentEnvelopesOf(items)) { + const { payload } = envelope; + switch (envelope.variant) { + case "AssistantMessage": { + sawAssistantMessage = true; + const tokens = getObject(getObject(payload, "usage"), "tokens") ?? getObject(payload, "usage") ?? {}; + turns.push({ + kind: "assistant", + ts: envelope.ts, + content: getString(payload, "text") ?? "", + inputTokens: getNumber(tokens, "input") ?? 0, + outputTokens: (getNumber(tokens, "output") ?? 0) + (getNumber(tokens, "reasoning") ?? 0), + toolCallCount: getNumber(payload, "tool_call_count") ?? null, + reasoning: readTurnReasoning(payload), + }); + break; + } + case "ToolCallStarted": { + const callId = getString(payload, "tool_call_id"); + if (!callId) break; + const args = payload.arguments; + pendingTools.set(callId, { + ts: envelope.ts, + toolName: getString(payload, "tool_name") ?? "", + input: typeof args === "string" ? args : JSON.stringify(args ?? ""), + }); + break; + } + case "ToolCallCompleted": { + const callId = getString(payload, "tool_call_id"); + if (!callId) break; + const started = pendingTools.get(callId); + pendingTools.delete(callId); + const output = payload.output ?? ""; + turns.push({ + kind: "tool", + ts: started?.ts ?? envelope.ts, + toolName: started?.toolName ?? getString(payload, "tool_name") ?? "", + input: started?.input ?? "", + result: typeof output === "string" ? output : JSON.stringify(output, null, 2), + isError: payload.is_error === true, + durationMs: durationBetween(started?.ts, envelope.ts), + }); + break; + } + case "SteeringInjected": { + const text = getString(payload, "text") ?? ""; + if (text) turns.push({ kind: "steer", ts: envelope.ts, content: text }); + break; + } + case "RoundInterrupted": + turns.push({ kind: "interrupt", ts: envelope.ts, content: "Interrupted — waiting for steering" }); + break; + default: + break; + } + } + if (!sawAssistantMessage && stage?.response) { + const tokens = stage.usage?.tokens; + turns.push({ + kind: "assistant", + ts: stage.completion?.timestamp ?? firstTs, + content: stage.response, + inputTokens: tokens?.input ?? 0, + outputTokens: tokens?.output ?? 0, + toolCallCount: null, + reasoning: null, + }); + } + + return { + turns, + pendingTools: Array.from(pendingTools, ([toolCallId, tool]) => ({ + toolCallId, + toolName: tool.toolName, + input: tool.input, + })), + }; +} + +/** A debug list item: a legacy event, or a Petri stream row. */ +type DebugListItem = EventEnvelope | DebugRow; + +function isDebugRow(item: DebugListItem): item is DebugRow { + return "item" in item && "category" in item; +} + +function debugItemSearchText(item: DebugListItem): string { + if (isDebugRow(item)) return debugRowSearchText(item); + return `${item.event ?? ""} ${JSON.stringify(item.properties ?? {})}`.toLowerCase(); +} + +/** What the details panel shows for a debug item. */ +function debugItemPayload(item: DebugListItem): EventDisplayPayload { + if (!isDebugRow(item)) return item; + return { + event: item.event, + stream_seq: item.seq, + kind: item.item.kind, + stage: item.stageLabel, + recorded_at: item.ts, + item: item.item.item, + }; +} + +/** What the Petri renderers show for one stage, derived once per stage. */ +interface PetriStageViews { + pairs: HumanInterviewPair[]; + edge: EdgeSelection | null; + overview: ParallelOverview; + reducer: ReducerTranscript | null; +} + type ToolTurn = Extract; type ToolGroupChild = { turn: ToolTurn; turnIndex: number }; type ToolGroupChildren = readonly [ @@ -2096,6 +2286,7 @@ function StageActivityBody({ contextData, runEvents, stages, + petri, }: { effectiveTab: EventsTab; renderer: StageRenderer; @@ -2107,15 +2298,18 @@ function StageActivityBody({ runId: string; selectedStage: Stage; commandTurn: CommandTurn | null; - debugEvents: EventEnvelope[]; - filteredDebugEvents: EventEnvelope[]; + debugEvents: DebugListItem[]; + filteredDebugEvents: DebugListItem[]; openDebugSeq: number | null; onDebugSeqChange: (seq: number | null) => void; contextData: ReturnType; runEvents: EventEnvelope[]; stages: Stage[]; + /** Set for a stage of a Petri run: the renderers read these, not events. */ + petri?: PetriStageViews; }) { const { turns, pendingTools } = activity; + const legacyEvents = petri ? [] : (debugEvents as EventEnvelope[]); return (
{effectiveTab === "chat" ? ( @@ -2169,23 +2363,29 @@ function StageActivityBody({ turn={commandTurn} /> ) : renderer === "human" ? ( - + ) : renderer === "conditional" ? ( ) : renderer === "parallel" ? ( ) : renderer === "fan_in" ? ( - + ) : renderer === "wait" ? ( ) : ( @@ -2227,6 +2427,7 @@ function RunStageActivityStage({ onKindsChange, onDebugCategoriesChange, onSearchChange, + petri, }: { runId: string; selectedStage: Stage; @@ -2240,25 +2441,54 @@ function RunStageActivityStage({ onKindsChange: (kinds: EventKind[]) => void; onDebugCategoriesChange: (categories: DebugCategory[]) => void; onSearchChange: (search: string) => void; + /** The run's stream and projection when it executes on Petri. */ + petri?: PetriRunData; }) { const selectedStageId = selectedStage.id; - const stageEventsQuery = useRunStageEvents(runId, selectedStageId); + const renderer: StageRenderer = selectStageRenderer(selectedStage.handler); + // A Petri run has no legacy stage events: its views read the stream and + // the projection, so the query stays idle. + const stageEventsQuery = useRunStageEvents(petri ? undefined : runId, selectedStageId); + const stageItems = useMemo( + () => (petri ? itemsForStage(petri.stream, selectedStageId) : []), + [petri, selectedStageId], + ); + const stageProjection: StageProjection | undefined = + petri?.projection.stages[selectedStageId]; const activity = useMemo( - () => buildStageActivity(stageEventsQuery.data ?? [], selectedStageId), - [stageEventsQuery.data, selectedStageId], + () => + petri + ? buildPetriStageActivity(stageItems, stageProjection, renderer) + : buildStageActivity(stageEventsQuery.data ?? [], selectedStageId), + [petri, stageItems, stageProjection, renderer, stageEventsQuery.data, selectedStageId], ); const { turns } = activity; - const renderer: StageRenderer = selectStageRenderer(selectedStage.handler); - const debugEvents = useMemo(() => { + const debugEvents = useMemo(() => { + if (petri) return debugRowsFromStream(stageItems); return (stageEventsQuery.data ?? []).filter( (event) => activityEventStageId(event) === selectedStageId, ); - }, [stageEventsQuery.data, selectedStageId]); + }, [petri, stageItems, stageEventsQuery.data, selectedStageId]); + const petriViews = useMemo( + () => + petri + ? { + pairs: parsePetriInterviewPairs(stageItems), + edge: findPetriEdgeForStage(petri.stream, selectedStageId), + overview: parallelOverviewFromProjection(stageProjection), + reducer: reducerTranscriptFromProjection(stageProjection), + } + : undefined, + [petri, stageItems, stageProjection, selectedStageId], + ); // The Context tab surfaces the workflow's deliberate per-visit outputs. It // only exists when the stage completed and actually wrote something. const contextData = useMemo( - () => extractStageContext(debugEvents), - [debugEvents], + () => + petri + ? extractPetriStageContext(stageItems) + : extractStageContext(debugEvents as EventEnvelope[]), + [petri, stageItems, debugEvents], ); const availableTabs = useMemo( () => @@ -2276,7 +2506,7 @@ function RunStageActivityStage({ // Some renderers need run-scoped events (e.g. conditional renders the // engine-level edge.selected event, which has no stage_id). Only fetch when // the active renderer actually needs it to keep this off the hot path. - const needsRunEvents = renderer === "conditional"; + const needsRunEvents = renderer === "conditional" && !petri; const runEventsQuery = useRunEventsList(needsRunEvents ? runId : undefined); const commandTurn = useMemo(() => { if (effectiveTab !== "primary" || renderer !== "command") return null; @@ -2347,34 +2577,27 @@ function RunStageActivityStage({ } return null; }, [isPrimaryAgent, displayItems, panelSelection]); - const openDebugEvent = useMemo( - () => - isDebug && openDebugSeq != null - ? (debugEvents.find((e) => e.seq === openDebugSeq) ?? null) - : null, - [isDebug, debugEvents, openDebugSeq], - ); + const openDebugEvent = useMemo(() => { + if (!isDebug || openDebugSeq == null) return null; + const item = debugEvents.find((e) => e.seq === openDebugSeq); + return item ? debugItemPayload(item) : null; + }, [isDebug, debugEvents, openDebugSeq]); const availableDebugCategories = useMemo(() => { if (!isDebug) return []; const set = new Set(); for (const event of debugEvents) { - if (event.event) set.add(debugCategory(event.event)); + if (event.event) set.add(debugRowCategory(event)); } return Array.from(set).sort(); }, [isDebug, debugEvents]); - const filteredDebugEvents = useMemo(() => { + const filteredDebugEvents = useMemo(() => { if (!isDebug) return []; const useCategoryFilter = selectedDebugCategories.length > 0; const cats = new Set(selectedDebugCategories); const needle = search.toLowerCase(); return debugEvents.filter((event) => { - const name = event.event ?? ""; - if (useCategoryFilter && !cats.has(debugCategory(name))) return false; - if (needle) { - const blob = - `${name} ${JSON.stringify(event.properties ?? {})}`.toLowerCase(); - if (!blob.includes(needle)) return false; - } + if (useCategoryFilter && !cats.has(debugRowCategory(event))) return false; + if (needle && !debugItemSearchText(event).includes(needle)) return false; return true; }); }, [isDebug, debugEvents, selectedDebugCategories, search]); @@ -2461,6 +2684,7 @@ function RunStageActivityStage({ contextData={contextData} runEvents={runEventsQuery.data ?? []} stages={stages} + petri={petriViews} />
@@ -2493,11 +2717,13 @@ function RunStageActivity({ selectedStage, stages, runStart, + petri, }: { runId: string; selectedStage: Stage; stages: Stage[]; runStart: string | undefined; + petri?: PetriRunData; }) { const [activityState, dispatchActivity] = useReducer( stageActivityReducer, @@ -2532,6 +2758,7 @@ function RunStageActivity({ onSearchChange={(nextSearch) => dispatchActivity({ type: "searchChanged", search: nextSearch }) } + petri={petri} /> ); } @@ -2552,10 +2779,20 @@ export default function RunStages() { selectedStage?.startedAt ?? runQuery.data?.timestamps.started_at ?? runQuery.data?.timestamps.created_at; - // Insights sidebar only renders for agent stages; fetch projection + context - // window only when the user is on one to keep the hot path lean. + // The projection says which engine ran the run (a Petri run's stage views + // read it and the run's stream) and feeds the insights sidebar of an + // agent stage; the context window is fetched only for one. const isAgentStage = selectedStage?.handler === "agent"; - const runStateQuery = useRunState(isAgentStage ? id : undefined); + const runStateQuery = useRunState(id); + const petri = isPetriRun(runStateQuery.data); + const streamQuery = useRunStream(petri ? id : undefined); + const petriData = useMemo( + () => + petri && runStateQuery.data && streamQuery.data + ? { stream: streamQuery.data, projection: runStateQuery.data } + : undefined, + [petri, runStateQuery.data, streamQuery.data], + ); const contextWindowQuery = useRunStageContextWindow( isAgentStage ? id : undefined, isAgentStage ? selectedStageId : undefined, @@ -2615,6 +2852,7 @@ export default function RunStages() { selectedStage={selectedStage} stages={stages} runStart={runStart} + petri={petriData} />
);