mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
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 `<subject>.<verb>` 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 <noreply@anthropic.com>
This commit is contained in:
parent
523f831a7a
commit
f22cf18c86
19 changed files with 1640 additions and 71 deletions
|
|
@ -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<string, DebugCategory> = {
|
||||
|
|
@ -34,6 +38,8 @@ const CATEGORY_LABEL: Record<DebugCategory, string> = {
|
|||
lifecycle: "Lifecycle",
|
||||
human: "Human",
|
||||
system: "System",
|
||||
petri: "Petri",
|
||||
platform: "Platform",
|
||||
};
|
||||
|
||||
const CATEGORY_TONE: Record<DebugCategory, string> = {
|
||||
|
|
@ -42,6 +48,8 @@ const CATEGORY_TONE: Record<DebugCategory, string> = {
|
|||
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<DebugCategory, string> = {
|
||||
|
|
@ -50,6 +58,8 @@ const CATEGORY_COLOR: Record<DebugCategory, string> = {
|
|||
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 {
|
||||
|
|
|
|||
|
|
@ -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 `<subject>.<verb>` 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 (
|
||||
<button
|
||||
type="button"
|
||||
|
|
@ -291,7 +307,7 @@ export function DebugDnaStrip({
|
|||
onSelect,
|
||||
runStart,
|
||||
}: {
|
||||
events: EventEnvelope[];
|
||||
events: DebugRowLike[];
|
||||
selectedSeq: number | null;
|
||||
onSelect: (seq: number) => void;
|
||||
runStart: string | undefined;
|
||||
|
|
@ -351,7 +367,7 @@ export function DebugDnaStrip({
|
|||
const ms = Date.parse(event.ts);
|
||||
if (Number.isNaN(ms)) return null;
|
||||
const pct = ((ms - range.start) / range.duration) * 100;
|
||||
const category = debugCategory(event.event);
|
||||
const category = debugRowCategory(event);
|
||||
const color = debugCategoryColor(category);
|
||||
const isSelected = event.seq === selectedSeq;
|
||||
const isHovered = hover?.seq === event.seq;
|
||||
|
|
@ -411,14 +427,14 @@ function DnaPopover({
|
|||
anchorRect,
|
||||
runStart,
|
||||
}: {
|
||||
event: EventEnvelope;
|
||||
event: DebugRowLike;
|
||||
anchorRect: DOMRect;
|
||||
runStart: string | undefined;
|
||||
}) {
|
||||
const category = debugCategory(event.event);
|
||||
const category = debugRowCategory(event);
|
||||
return (
|
||||
<FloatingTooltip rect={anchorRect} placement="top">
|
||||
{`${debugCategoryLabel(category)} · ${friendlyEventName(event.event)} · ${formatElapsed(event.ts, runStart)}`}
|
||||
{`${debugCategoryLabel(category)} · ${friendlyEventName(event.event ?? "")} · ${formatElapsed(event.ts, runStart)}`}
|
||||
</FloatingTooltip>
|
||||
);
|
||||
}
|
||||
|
|
|
|||
96
apps/fabro-web/app/components/platform-records-panel.tsx
Normal file
96
apps/fabro-web/app/components/platform-records-panel.tsx
Normal file
|
|
@ -0,0 +1,96 @@
|
|||
import { useMemo } from "react";
|
||||
import type { RunProjection } from "@qltysh/fabro-api-client";
|
||||
|
||||
import { formatAbsoluteTs } from "../lib/format";
|
||||
import {
|
||||
isPetriRun,
|
||||
platformRecordsOf,
|
||||
type PlatformRecordEntry,
|
||||
} from "../lib/petri-stream";
|
||||
import { useRunState, useRunStream } from "../lib/queries";
|
||||
|
||||
const KIND_LABEL: Record<string, string> = {
|
||||
"checkpoint": "Checkpoint",
|
||||
"pull_request.created": "Pull request",
|
||||
"run.notice": "Notice",
|
||||
"run.title": "Title",
|
||||
"run.branch": "Run branch",
|
||||
};
|
||||
|
||||
function kindLabel(kind: string): string {
|
||||
return KIND_LABEL[kind] ?? kind;
|
||||
}
|
||||
|
||||
/**
|
||||
* The platform records of a Petri run: Fabro's own facts beside the engine's
|
||||
* events (a checkpoint with its commit, a pull request, a notice), each with
|
||||
* the stage it belongs to when it belongs to one.
|
||||
*/
|
||||
export function PlatformRecordsPanelView({
|
||||
records,
|
||||
projection,
|
||||
}: {
|
||||
records: PlatformRecordEntry[];
|
||||
projection: RunProjection | null | undefined;
|
||||
}) {
|
||||
const pullRequest = projection?.pull_request ?? null;
|
||||
if (records.length === 0 && !pullRequest) return null;
|
||||
return (
|
||||
<section
|
||||
aria-label="Platform records"
|
||||
className="rounded-md border border-line bg-panel/60 px-6 py-4"
|
||||
>
|
||||
<h3 className="text-[10px] font-medium uppercase tracking-[0.08em] text-fg-muted">
|
||||
Platform records
|
||||
</h3>
|
||||
<ul className="mt-2 space-y-1 text-sm">
|
||||
{pullRequest && (
|
||||
<li className="flex items-baseline gap-3">
|
||||
<span className="w-28 shrink-0 text-fg-muted">Pull request</span>
|
||||
<a
|
||||
href={pullRequest.html_url}
|
||||
target="_blank"
|
||||
rel="noreferrer"
|
||||
className="truncate font-mono text-teal-500 hover:text-teal-300"
|
||||
>
|
||||
#{pullRequest.number}
|
||||
</a>
|
||||
</li>
|
||||
)}
|
||||
{records.map((record) => (
|
||||
<li
|
||||
key={record.streamSeq}
|
||||
data-kind={record.kind}
|
||||
className="flex items-baseline gap-3"
|
||||
>
|
||||
<span className="w-28 shrink-0 text-fg-muted">{kindLabel(record.kind)}</span>
|
||||
<span className="min-w-0 flex-1 truncate font-mono text-fg-2">
|
||||
{record.detail ?? "—"}
|
||||
</span>
|
||||
{record.stageKey && (
|
||||
<span className="shrink-0 font-mono text-xs text-fg-muted">
|
||||
stage {record.stageKey}
|
||||
</span>
|
||||
)}
|
||||
<span className="shrink-0 font-mono text-xs tabular-nums text-fg-muted">
|
||||
{formatAbsoluteTs(record.ts)}
|
||||
</span>
|
||||
</li>
|
||||
))}
|
||||
</ul>
|
||||
</section>
|
||||
);
|
||||
}
|
||||
|
||||
/** The panel for a run, shown only when the run executes on Petri. */
|
||||
export function PlatformRecordsPanel({ runId }: { runId: string }) {
|
||||
const runStateQuery = useRunState(runId);
|
||||
const petri = isPetriRun(runStateQuery.data);
|
||||
const streamQuery = useRunStream(petri ? runId : undefined);
|
||||
const records = useMemo(
|
||||
() => (streamQuery.data ? platformRecordsOf(streamQuery.data) : []),
|
||||
[streamQuery.data],
|
||||
);
|
||||
if (!petri) return null;
|
||||
return <PlatformRecordsPanelView records={records} projection={runStateQuery.data} />;
|
||||
}
|
||||
|
|
@ -17,6 +17,12 @@ import type { EventEnvelope } from "@qltysh/fabro-api-client";
|
|||
interface WaterfallProps {
|
||||
runId: string;
|
||||
events: EventEnvelope[];
|
||||
/**
|
||||
* The run's phases when the caller derives them itself: a Petri run's
|
||||
* come from its platform lifecycle records (`deriveRunPhasesFromStream`),
|
||||
* not from legacy events.
|
||||
*/
|
||||
phases?: RunPhase[];
|
||||
stages: RunStage[];
|
||||
createdAtIso: string;
|
||||
completedAtIso: string | null;
|
||||
|
|
@ -158,17 +164,21 @@ function stageRow(runId: string, stage: RunStage, nowMs: number): Row | null {
|
|||
function buildRows({
|
||||
runId,
|
||||
events,
|
||||
phases: givenPhases,
|
||||
stages,
|
||||
createdAtIso,
|
||||
nowMs,
|
||||
}: {
|
||||
runId: string;
|
||||
events: EventEnvelope[];
|
||||
phases?: RunPhase[];
|
||||
stages: RunStage[];
|
||||
createdAtIso: string;
|
||||
nowMs: number;
|
||||
}): Row[] {
|
||||
const phases = deriveRunPhases(events, createdAtIso).map((p) => phaseRow(p, nowMs));
|
||||
const phases = (givenPhases ?? deriveRunPhases(events, createdAtIso)).map((p) =>
|
||||
phaseRow(p, nowMs),
|
||||
);
|
||||
const stageRows: Row[] = [];
|
||||
for (const stage of stages) {
|
||||
if (!isVisibleStage(stage.node_id)) continue;
|
||||
|
|
@ -182,14 +192,15 @@ function buildRows({
|
|||
export function RunWaterfall({
|
||||
runId,
|
||||
events,
|
||||
phases,
|
||||
stages,
|
||||
createdAtIso,
|
||||
completedAtIso,
|
||||
}: WaterfallProps) {
|
||||
const nowMs = useTickingNow(true, 1000);
|
||||
const rows = useMemo(
|
||||
() => buildRows({ runId, events, stages, createdAtIso, nowMs }),
|
||||
[runId, events, stages, createdAtIso, nowMs],
|
||||
() => buildRows({ runId, events, phases, stages, createdAtIso, nowMs }),
|
||||
[runId, events, phases, stages, createdAtIso, nowMs],
|
||||
);
|
||||
|
||||
const createdMs = Date.parse(createdAtIso);
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import type { EventEnvelope } from "@qltysh/fabro-api-client";
|
|||
|
||||
import type { Stage } from "../stage-sidebar";
|
||||
import { StageMetaBar } from "./meta-bar";
|
||||
import { findEdgeForNode } from "./helpers";
|
||||
import { findEdgeForNode, type EdgeSelection } from "./helpers";
|
||||
|
||||
const REASON_LABEL: Record<string, string> = {
|
||||
condition: "Matched condition",
|
||||
|
|
@ -24,17 +24,20 @@ function reasonLabel(reason: string): string {
|
|||
export function ConditionalDecision({
|
||||
stage,
|
||||
runEvents,
|
||||
edge: givenEdge,
|
||||
allStages,
|
||||
runId,
|
||||
}: {
|
||||
stage: Stage;
|
||||
runEvents: EventEnvelope[];
|
||||
/** The edge when the caller derived it (a Petri run's `route.applied`). */
|
||||
edge?: EdgeSelection | null;
|
||||
allStages: Stage[];
|
||||
runId: string;
|
||||
}) {
|
||||
const edge = useMemo(
|
||||
() => findEdgeForNode(runEvents, stage.nodeId),
|
||||
[runEvents, stage.nodeId],
|
||||
() => (givenEdge !== undefined ? givenEdge : findEdgeForNode(runEvents, stage.nodeId)),
|
||||
[givenEdge, runEvents, stage.nodeId],
|
||||
);
|
||||
const targetStage = useMemo(() => {
|
||||
if (!edge) return null;
|
||||
|
|
|
|||
|
|
@ -9,16 +9,22 @@ import type { Stage } from "../stage-sidebar";
|
|||
import { formatTokenCount } from "../../lib/format";
|
||||
import { Markdown } from "./primitives";
|
||||
import { StageMetaBar } from "./meta-bar";
|
||||
import { parseReducerTranscript } from "./helpers";
|
||||
import { parseReducerTranscript, type ReducerTranscript } from "./helpers";
|
||||
|
||||
export function FanInResults({
|
||||
stage,
|
||||
events,
|
||||
reducer: givenReducer,
|
||||
}: {
|
||||
stage: Stage;
|
||||
events: EventEnvelope[];
|
||||
/** The transcript when the caller derived it (a Petri run's projection). */
|
||||
reducer?: ReducerTranscript | null;
|
||||
}) {
|
||||
const reducer = useMemo(() => parseReducerTranscript(events), [events]);
|
||||
const reducer = useMemo(
|
||||
() => (givenReducer !== undefined ? givenReducer : parseReducerTranscript(events)),
|
||||
[givenReducer, events],
|
||||
);
|
||||
|
||||
return (
|
||||
<div className="space-y-6 pl-3 pr-4 sm:pr-6 lg:pr-8">
|
||||
|
|
|
|||
|
|
@ -45,7 +45,7 @@ export interface HumanInterviewPair {
|
|||
resolution: HumanResolution | null;
|
||||
}
|
||||
|
||||
function principalLabel(actor: unknown): string | null {
|
||||
export function principalLabel(actor: unknown): string | null {
|
||||
if (!actor || typeof actor !== "object") return null;
|
||||
const record = actor as UnknownRecord;
|
||||
const kind = getString(record, "kind") ?? "";
|
||||
|
|
|
|||
|
|
@ -225,11 +225,17 @@ function QuestionBlock({
|
|||
export function HumanQA({
|
||||
stage,
|
||||
events,
|
||||
pairs: givenPairs,
|
||||
}: {
|
||||
stage: Stage;
|
||||
events: EventEnvelope[];
|
||||
/** The pairs when the caller derived them (a Petri run's stream). */
|
||||
pairs?: HumanInterviewPair[];
|
||||
}) {
|
||||
const pairs = useMemo(() => parseHumanInterviewPairs(events), [events]);
|
||||
const pairs = useMemo(
|
||||
() => givenPairs ?? parseHumanInterviewPairs(events),
|
||||
[givenPairs, events],
|
||||
);
|
||||
const stageActive = ACTIVE_STAGE_STATES.has(stage.status);
|
||||
const pendingCount = pairs.filter((p) => p.resolution == null).length;
|
||||
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import type { Stage } from "../stage-sidebar";
|
|||
import { formatStageLabel, stageStatusLabel, stageStatusTone } from "../../lib/stage-sidebar";
|
||||
import { StageMetaBar } from "./meta-bar";
|
||||
import { parseParallelOverview } from "./helpers";
|
||||
import type { ParallelBranchSummary } from "./helpers";
|
||||
import type { ParallelBranchSummary, ParallelOverview } from "./helpers";
|
||||
|
||||
/** Branch row view state sourced from a live branch stage or completed result. */
|
||||
interface BranchRow {
|
||||
|
|
@ -99,15 +99,21 @@ function ChildRow({
|
|||
export function ParallelChildren({
|
||||
stage,
|
||||
events,
|
||||
overview: givenOverview,
|
||||
runId,
|
||||
allStages,
|
||||
}: {
|
||||
stage: Stage;
|
||||
events: EventEnvelope[];
|
||||
/** The overview when the caller derived it (a Petri run's projection). */
|
||||
overview?: ParallelOverview;
|
||||
runId: string;
|
||||
allStages: Stage[];
|
||||
}) {
|
||||
const overview = useMemo(() => parseParallelOverview(events), [events]);
|
||||
const overview = useMemo(
|
||||
() => givenOverview ?? parseParallelOverview(events),
|
||||
[givenOverview, events],
|
||||
);
|
||||
|
||||
const stagesByBranchIndex = useMemo(() => {
|
||||
const byIndex = new Map<number, Stage>();
|
||||
|
|
|
|||
|
|
@ -1,5 +1,3 @@
|
|||
import type { EventEnvelope } from "@qltysh/fabro-api-client";
|
||||
|
||||
import type { Stage } from "../stage-sidebar";
|
||||
import {
|
||||
debugCategory,
|
||||
|
|
@ -9,16 +7,22 @@ import {
|
|||
} from "../event-debug-helpers";
|
||||
import { StageMetaBar } from "./meta-bar";
|
||||
|
||||
interface CategoryCount {
|
||||
export interface CategoryCount {
|
||||
category: DebugCategory;
|
||||
count: number;
|
||||
}
|
||||
|
||||
function summarizeEventCategories(events: EventEnvelope[]): CategoryCount[] {
|
||||
/** What the summary counts: a legacy event, or a Petri stream row with its own category. */
|
||||
export interface CategorizedEvent {
|
||||
event?: string | null;
|
||||
category?: DebugCategory;
|
||||
}
|
||||
|
||||
export function summarizeEventCategories(events: CategorizedEvent[]): CategoryCount[] {
|
||||
const counts = new Map<DebugCategory, number>();
|
||||
for (const event of events) {
|
||||
if (!event.event) continue;
|
||||
const cat = debugCategory(event.event);
|
||||
const cat = event.category ?? (event.event ? debugCategory(event.event) : null);
|
||||
if (!cat) continue;
|
||||
counts.set(cat, (counts.get(cat) ?? 0) + 1);
|
||||
}
|
||||
return Array.from(counts.entries())
|
||||
|
|
@ -31,7 +35,7 @@ export function StageSummary({
|
|||
events,
|
||||
}: {
|
||||
stage: Stage;
|
||||
events: EventEnvelope[];
|
||||
events: CategorizedEvent[];
|
||||
}) {
|
||||
const categories = summarizeEventCategories(events);
|
||||
|
||||
|
|
|
|||
|
|
@ -1163,6 +1163,11 @@ function candidateKey(candidate: CandidateMessage): string {
|
|||
}
|
||||
|
||||
export function eventDedupeKey(payload: EventPayload): string | undefined {
|
||||
// A run stream item's `id` is the item's own identity within its run (a
|
||||
// Petri `EventId` or a platform record seq), so two runs share ids.
|
||||
if (typeof payload.stream_seq === "number" && typeof payload.run_id === "string") {
|
||||
return `${payload.run_id}:stream:${payload.stream_seq}`;
|
||||
}
|
||||
if (typeof payload.id === "string" && payload.id.length > 0) {
|
||||
return payload.id;
|
||||
}
|
||||
|
|
|
|||
692
apps/fabro-web/app/lib/petri-stream.ts
Normal file
692
apps/fabro-web/app/lib/petri-stream.ts
Normal file
|
|
@ -0,0 +1,692 @@
|
|||
/**
|
||||
* Pure helpers over a Petri run's stream: the `RunStreamItem`s
|
||||
* `GET /runs/{id}/events` serves for a run that executes on Petri. Each item
|
||||
* is a Petri `RunEvent` (Petri's event contract, passed through as JSON) or a
|
||||
* stored platform record (Fabro's own fact about the run), in one envelope
|
||||
* keyed by `stream_seq`. The mapping from items to views follows
|
||||
* `lib/components/fabro-petri/VIEWS.md`.
|
||||
*/
|
||||
import { StageOutcome, StageState } from "@qltysh/fabro-api-client";
|
||||
import type {
|
||||
RunProjection,
|
||||
RunStreamItem,
|
||||
StageProjection,
|
||||
} from "@qltysh/fabro-api-client";
|
||||
|
||||
import type {
|
||||
EdgeSelection,
|
||||
HumanInterviewPair,
|
||||
HumanResolution,
|
||||
InterviewOption,
|
||||
ParallelOverview,
|
||||
ReducerTranscript,
|
||||
StageContextData,
|
||||
} from "../components/stage-renderers/helpers";
|
||||
import { principalLabel } from "../components/stage-renderers/helpers";
|
||||
import type { Stage } from "./stage-sidebar";
|
||||
import type { RunPhase, RunPhaseKind } from "./run-phases";
|
||||
import { formatDurationMs } from "./format";
|
||||
import {
|
||||
getArray,
|
||||
getBool,
|
||||
getNumber,
|
||||
getObject,
|
||||
getString,
|
||||
isRecord,
|
||||
type UnknownRecord,
|
||||
} from "./unknown";
|
||||
|
||||
export type PetriStream = ReadonlyArray<RunStreamItem>;
|
||||
|
||||
export function isPetriItem(item: RunStreamItem): boolean {
|
||||
return item.kind === "petri";
|
||||
}
|
||||
|
||||
export function isPlatformItem(item: RunStreamItem): boolean {
|
||||
return item.kind === "platform";
|
||||
}
|
||||
|
||||
/** Whether the projection is of a run that executes on Petri. */
|
||||
export function isPetriRun(
|
||||
projection: RunProjection | null | undefined,
|
||||
): boolean {
|
||||
const engine = projection?.spec?.engine;
|
||||
return isRecord(engine) && getString(engine, "kind") === "petri";
|
||||
}
|
||||
|
||||
/** Whether an SSE payload is a run stream item rather than a legacy event. */
|
||||
export function isStreamItemPayload(
|
||||
payload: unknown,
|
||||
): payload is RunStreamItem {
|
||||
return (
|
||||
isRecord(payload) &&
|
||||
typeof payload.stream_seq === "number" &&
|
||||
(payload.kind === "petri" || payload.kind === "platform")
|
||||
);
|
||||
}
|
||||
|
||||
function record(item: RunStreamItem): UnknownRecord | undefined {
|
||||
return getObject(item.item, "record");
|
||||
}
|
||||
|
||||
function derived(item: RunStreamItem): UnknownRecord | undefined {
|
||||
return getObject(item.item, "derived");
|
||||
}
|
||||
|
||||
/** The stored platform record's `kind`, for a platform item. */
|
||||
export function platformRecordKind(item: RunStreamItem): string | undefined {
|
||||
if (!isPlatformItem(item)) return undefined;
|
||||
return getString(record(item), "kind");
|
||||
}
|
||||
|
||||
/**
|
||||
* The `<subject>.<verb>` name of a Petri event: the recorded body's `event`
|
||||
* tag, or a view event's tag under `derived`.
|
||||
*/
|
||||
export function petriEventName(item: RunStreamItem): string | undefined {
|
||||
if (!isPetriItem(item)) return undefined;
|
||||
return (
|
||||
getString(getObject(record(item), "body"), "event") ??
|
||||
getString(derived(item), "event")
|
||||
);
|
||||
}
|
||||
|
||||
/** The name a listing shows: the Petri event name or the platform kind. */
|
||||
export function streamItemName(item: RunStreamItem): string {
|
||||
return petriEventName(item) ?? platformRecordKind(item) ?? item.kind;
|
||||
}
|
||||
|
||||
/** When the item's record was appended, as an ISO timestamp. */
|
||||
export function streamItemTs(item: RunStreamItem): string {
|
||||
return new Date(item.recorded_at).toISOString();
|
||||
}
|
||||
|
||||
/** Petri's reading of a `step.progress.recorded` payload (`derived.parsed`). */
|
||||
export function petriParsed(item: RunStreamItem): UnknownRecord | undefined {
|
||||
return getObject(derived(item), "parsed");
|
||||
}
|
||||
|
||||
/** The recorded event body of a Petri item (`record.body`). */
|
||||
export function petriBody(item: RunStreamItem): UnknownRecord | undefined {
|
||||
return getObject(record(item), "body");
|
||||
}
|
||||
|
||||
function subject(item: RunStreamItem): UnknownRecord | undefined {
|
||||
return getObject(item.item, "subject");
|
||||
}
|
||||
|
||||
function subjectNode(item: RunStreamItem): UnknownRecord | undefined {
|
||||
return getObject(subject(item), "node");
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether the subject's node is a stage of its own. A lowering node (the
|
||||
* `parallel.branch` delegate the fork's execution holds for each branch, a
|
||||
* synthetic fan-in placeholder) shares a name with a real stage and is not
|
||||
* one.
|
||||
*/
|
||||
function isShownNode(node: UnknownRecord | undefined): boolean {
|
||||
if (!node) return false;
|
||||
const meta = getObject(node, "meta");
|
||||
if (getBool(meta, "synthetic") === true) return false;
|
||||
return getString(meta, "kind") !== "parallel.branch";
|
||||
}
|
||||
|
||||
/**
|
||||
* The stage label (`node@visit`) of a Petri item's subject, or `undefined`
|
||||
* for an item with no subject or one whose node is a lowering node.
|
||||
*/
|
||||
export function petriStageLabel(item: RunStreamItem): string | undefined {
|
||||
const node = subjectNode(item);
|
||||
if (!isShownNode(node)) return undefined;
|
||||
const name = getString(node, "name");
|
||||
if (!name) return undefined;
|
||||
const visit = getNumber(subject(item), "visit") ?? 1;
|
||||
return `${name}@${visit}`;
|
||||
}
|
||||
|
||||
/** The subject's stage key, `(execution, firing)`, for an item under a firing. */
|
||||
export function petriStageKey(item: RunStreamItem): string | undefined {
|
||||
const firing = getNumber(subject(item), "firing");
|
||||
const execution = getNumber(getObject(item.item, "context"), "execution");
|
||||
if (firing === undefined || execution === undefined) return undefined;
|
||||
return `${execution}:${firing}`;
|
||||
}
|
||||
|
||||
/** The items whose subject is the stage with this label, in stream order. */
|
||||
export function itemsForStage(
|
||||
stream: PetriStream,
|
||||
stageLabel: string,
|
||||
): RunStreamItem[] {
|
||||
return stream.filter((item) => petriStageLabel(item) === stageLabel);
|
||||
}
|
||||
|
||||
// ── Questions ───────────────────────────────────────────────────────────
|
||||
|
||||
function parseOptions(value: unknown): InterviewOption[] {
|
||||
if (!Array.isArray(value)) return [];
|
||||
const out: InterviewOption[] = [];
|
||||
for (const entry of value) {
|
||||
const key = getString(entry, "key");
|
||||
const label = getString(entry, "label");
|
||||
if (!key || !label) continue;
|
||||
const option: InterviewOption = { key, label };
|
||||
const description = getString(entry, "description");
|
||||
const preview = getString(entry, "preview");
|
||||
if (description !== undefined) option.description = description;
|
||||
if (preview !== undefined) option.preview = preview;
|
||||
out.push(option);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
function answerText(answer: UnknownRecord): string {
|
||||
const choice = getString(answer, "choice");
|
||||
if (choice) return choice;
|
||||
const text = getString(answer, "text");
|
||||
if (text) return text;
|
||||
const choices = getArray(answer, "choices");
|
||||
if (choices) return choices.filter((c): c is string => typeof c === "string").join(", ");
|
||||
if (getBool(answer, "cancelled") === true) return "";
|
||||
if (getBool(answer, "confirmed") !== undefined) {
|
||||
return getBool(answer, "confirmed") ? "yes" : "no";
|
||||
}
|
||||
for (const value of Object.values(answer)) {
|
||||
if (typeof value === "string") return value;
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
/**
|
||||
* Pair each question a stage asked (`step.progress.recorded` with
|
||||
* `derived.parsed.kind === "question"`) with what resolved it: the delivered
|
||||
* `control.requested` answer, a `question_expired` reading, or a cancelled
|
||||
* answer. The answering principal comes from the `interview.answered`
|
||||
* platform record keyed on Petri's question id.
|
||||
*/
|
||||
export function parsePetriInterviewPairs(stream: PetriStream): HumanInterviewPair[] {
|
||||
const pairs = new Map<string, HumanInterviewPair>();
|
||||
const actors = new Map<string, string | null>();
|
||||
const askedAt = new Map<string, number>();
|
||||
|
||||
for (const item of stream) {
|
||||
if (isPlatformItem(item)) {
|
||||
const rec = record(item);
|
||||
if (getString(rec, "kind") === "interview.answered") {
|
||||
const question = getString(rec, "question");
|
||||
if (question) actors.set(question, principalLabel(rec?.principal));
|
||||
}
|
||||
continue;
|
||||
}
|
||||
const name = petriEventName(item);
|
||||
const parsed = petriParsed(item);
|
||||
if (name === "step.progress.recorded" && getString(parsed, "kind") === "question") {
|
||||
const question = getObject(parsed, "question");
|
||||
const id = getString(question, "id");
|
||||
if (!id) continue;
|
||||
const timeoutMs = getNumber(question, "timeout_ms");
|
||||
askedAt.set(id, item.recorded_at);
|
||||
pairs.set(id, {
|
||||
question: {
|
||||
ts: streamItemTs(item),
|
||||
questionId: id,
|
||||
question: getString(question, "text") ?? "",
|
||||
questionType: getString(question, "kind") ?? "freeform",
|
||||
options: parseOptions(question?.options),
|
||||
allowFreeform: getBool(question, "freeform") === true,
|
||||
timeoutSeconds: timeoutMs !== undefined ? Math.round(timeoutMs / 1000) : null,
|
||||
contextDisplay: getString(question, "context") ?? null,
|
||||
reviewTarget: null,
|
||||
},
|
||||
resolution: null,
|
||||
});
|
||||
continue;
|
||||
}
|
||||
if (name === "step.progress.recorded" && getString(parsed, "kind") === "question_expired") {
|
||||
const expired = getObject(parsed, "expired");
|
||||
const id = getString(expired, "question");
|
||||
const pair = id ? pairs.get(id) : undefined;
|
||||
if (!pair || !id) continue;
|
||||
pair.resolution = {
|
||||
kind: "timeout",
|
||||
ts: streamItemTs(item),
|
||||
durationMs: getNumber(expired, "waited_ms") ?? item.recorded_at - (askedAt.get(id) ?? item.recorded_at),
|
||||
};
|
||||
continue;
|
||||
}
|
||||
if (name === "control.requested") {
|
||||
const d = derived(item);
|
||||
const answer = getObject(d, "answer");
|
||||
const id = getString(answer, "question");
|
||||
const pair = id ? pairs.get(id) : undefined;
|
||||
if (!pair || !id || !answer) continue;
|
||||
if (getBool(d, "deliverable") === false) continue;
|
||||
const durationMs = item.recorded_at - (askedAt.get(id) ?? item.recorded_at);
|
||||
const resolution: HumanResolution =
|
||||
getBool(answer, "cancelled") === true
|
||||
? {
|
||||
kind: "interrupted",
|
||||
ts: streamItemTs(item),
|
||||
reason: "cancelled",
|
||||
durationMs,
|
||||
actor: null,
|
||||
}
|
||||
: {
|
||||
kind: "answered",
|
||||
ts: streamItemTs(item),
|
||||
answer: answerText(answer),
|
||||
durationMs,
|
||||
actor: null,
|
||||
};
|
||||
pair.resolution = resolution;
|
||||
}
|
||||
}
|
||||
|
||||
for (const pair of pairs.values()) {
|
||||
const resolution = pair.resolution;
|
||||
if (resolution && resolution.kind !== "timeout") {
|
||||
resolution.actor = actors.get(pair.question.questionId) ?? null;
|
||||
}
|
||||
}
|
||||
|
||||
return Array.from(pairs.values()).sort((a, b) => a.question.ts.localeCompare(b.question.ts));
|
||||
}
|
||||
|
||||
// ── Run phases ──────────────────────────────────────────────────────────
|
||||
|
||||
const PHASE_LABEL: Record<RunPhaseKind, string> = {
|
||||
submitted: "Submitted",
|
||||
pending: "Pending",
|
||||
runnable: "Runnable",
|
||||
initializing: "Initializing",
|
||||
};
|
||||
|
||||
const TERMINAL_TRANSITIONS: ReadonlySet<string> = new Set(["succeeded", "failed", "dead"]);
|
||||
|
||||
/** Whether a platform item is the run's terminal lifecycle record. */
|
||||
export function isTerminalLifecycleItem(item: RunStreamItem): boolean {
|
||||
if (platformRecordKind(item) !== "run.lifecycle") return false;
|
||||
const transition = getString(record(item), "transition");
|
||||
return transition !== undefined && TERMINAL_TRANSITIONS.has(transition);
|
||||
}
|
||||
|
||||
/**
|
||||
* The run's phases before its stages own the timeline, from the platform
|
||||
* `run.lifecycle` records: the same slices `deriveRunPhases` cuts from the
|
||||
* legacy lifecycle events.
|
||||
*/
|
||||
export function deriveRunPhasesFromStream(
|
||||
stream: PetriStream,
|
||||
createdAtIso: string,
|
||||
): RunPhase[] {
|
||||
const createdMs = Date.parse(createdAtIso);
|
||||
if (Number.isNaN(createdMs)) return [];
|
||||
|
||||
let startRequestedMs: number | null = null;
|
||||
let pendingMs: number | null = null;
|
||||
let runnableMs: number | null = null;
|
||||
let startingMs: number | null = null;
|
||||
let runningMs: number | null = null;
|
||||
let terminalMs: number | null = null;
|
||||
|
||||
for (const item of stream) {
|
||||
if (platformRecordKind(item) !== "run.lifecycle") continue;
|
||||
const transition = getString(record(item), "transition");
|
||||
const ms = item.recorded_at;
|
||||
switch (transition) {
|
||||
case "start_requested":
|
||||
startRequestedMs ??= ms;
|
||||
break;
|
||||
case "pending":
|
||||
pendingMs ??= ms;
|
||||
break;
|
||||
case "runnable":
|
||||
runnableMs ??= ms;
|
||||
break;
|
||||
case "starting":
|
||||
startingMs ??= ms;
|
||||
break;
|
||||
case "running":
|
||||
runningMs ??= ms;
|
||||
break;
|
||||
case "succeeded":
|
||||
case "failed":
|
||||
case "dead":
|
||||
terminalMs ??= ms;
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
const phases: RunPhase[] = [];
|
||||
phases.push({
|
||||
kind: "submitted",
|
||||
label: PHASE_LABEL.submitted,
|
||||
startMs: createdMs,
|
||||
endMs: startRequestedMs ?? pendingMs ?? runnableMs ?? startingMs ?? runningMs ?? terminalMs,
|
||||
});
|
||||
if (pendingMs != null) {
|
||||
phases.push({
|
||||
kind: "pending",
|
||||
label: PHASE_LABEL.pending,
|
||||
startMs: pendingMs,
|
||||
endMs: runnableMs ?? startingMs ?? runningMs ?? terminalMs,
|
||||
});
|
||||
}
|
||||
if (runnableMs != null) {
|
||||
phases.push({
|
||||
kind: "runnable",
|
||||
label: PHASE_LABEL.runnable,
|
||||
startMs: runnableMs,
|
||||
endMs: startingMs ?? runningMs ?? terminalMs,
|
||||
});
|
||||
}
|
||||
if (startingMs != null) {
|
||||
phases.push({
|
||||
kind: "initializing",
|
||||
label: PHASE_LABEL.initializing,
|
||||
startMs: startingMs,
|
||||
endMs: runningMs ?? terminalMs,
|
||||
});
|
||||
}
|
||||
return phases;
|
||||
}
|
||||
|
||||
// ── Platform records ────────────────────────────────────────────────────
|
||||
|
||||
export interface PlatformRecordEntry {
|
||||
streamSeq: number;
|
||||
kind: string;
|
||||
ts: string;
|
||||
/** The stage the record belongs to, as `execution:firing`, if any. */
|
||||
stageKey: string | null;
|
||||
/** A one-line summary: the commit sha, the pull request url, the notice. */
|
||||
detail: string | null;
|
||||
}
|
||||
|
||||
/** The platform records on the stream that name a Fabro fact worth a row. */
|
||||
export function platformRecordsOf(stream: PetriStream): PlatformRecordEntry[] {
|
||||
const out: PlatformRecordEntry[] = [];
|
||||
for (const item of stream) {
|
||||
const kind = platformRecordKind(item);
|
||||
if (!kind) continue;
|
||||
const rec = record(item) ?? {};
|
||||
const position = getObject(item.item, "position");
|
||||
const execution = getNumber(position, "execution") ?? getNumber(rec, "execution");
|
||||
const firing = getNumber(position, "firing") ?? getNumber(rec, "firing");
|
||||
const stageKey =
|
||||
execution !== undefined && firing !== undefined ? `${execution}:${firing}` : null;
|
||||
let detail: string | null = null;
|
||||
switch (kind) {
|
||||
case "checkpoint":
|
||||
detail = getString(rec, "git_commit_sha")?.slice(0, 12) ?? null;
|
||||
break;
|
||||
case "pull_request.created":
|
||||
detail = getString(rec, "html_url") ?? getString(rec, "url") ?? null;
|
||||
break;
|
||||
case "run.notice":
|
||||
detail = getString(rec, "message") ?? getString(rec, "code") ?? null;
|
||||
break;
|
||||
case "run.title":
|
||||
detail = getString(rec, "title") ?? null;
|
||||
break;
|
||||
case "run.branch":
|
||||
detail = getString(rec, "run_branch") ?? null;
|
||||
break;
|
||||
default:
|
||||
continue;
|
||||
}
|
||||
out.push({ streamSeq: item.stream_seq, kind, ts: streamItemTs(item), stageKey, detail });
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// ── Debug rows ──────────────────────────────────────────────────────────
|
||||
|
||||
/** A stream item as the events listing and the stage debug tab show it. */
|
||||
export interface DebugRow {
|
||||
/** The `stream_seq`: the row's key and the cursor. */
|
||||
seq: number;
|
||||
/** The `<subject>.<verb>` name or the platform record kind. */
|
||||
event: string;
|
||||
ts: string;
|
||||
category: "petri" | "platform";
|
||||
stageLabel: string | null;
|
||||
/** The raw item, for the details panel. */
|
||||
item: RunStreamItem;
|
||||
}
|
||||
|
||||
export function debugRowsFromStream(stream: PetriStream): DebugRow[] {
|
||||
return stream.map((item) => ({
|
||||
seq: item.stream_seq,
|
||||
event: streamItemName(item),
|
||||
ts: streamItemTs(item),
|
||||
category: isPlatformItem(item) ? "platform" : "petri",
|
||||
stageLabel: petriStageLabel(item) ?? null,
|
||||
item,
|
||||
}));
|
||||
}
|
||||
|
||||
/** The text a search box matches a row against. */
|
||||
export function debugRowSearchText(row: DebugRow): string {
|
||||
const body = isPlatformItem(row.item)
|
||||
? record(row.item)
|
||||
: { ...(petriBody(row.item) ?? {}), derived: derived(row.item) ?? {} };
|
||||
return `${row.event} ${row.stageLabel ?? ""} ${JSON.stringify(body ?? {})}`.toLowerCase();
|
||||
}
|
||||
|
||||
// ── Stage renderers ─────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* The edge a stage's firing took, from its `route.applied` record:
|
||||
* `derived.target` is the node Petri resolved, `kind` says whether the edge
|
||||
* was followed or jumped to.
|
||||
*/
|
||||
export function findPetriEdgeForStage(
|
||||
stream: PetriStream,
|
||||
stageLabel: string,
|
||||
): EdgeSelection | null {
|
||||
let latest: EdgeSelection | null = null;
|
||||
for (const item of stream) {
|
||||
if (petriEventName(item) !== "route.applied") continue;
|
||||
if (petriStageLabel(item) !== stageLabel) continue;
|
||||
const target = getString(getObject(derived(item), "target"), "name");
|
||||
if (!target) continue;
|
||||
const kind = getString(petriBody(item), "kind") ?? "edge";
|
||||
latest = {
|
||||
fromNode: getString(subjectNode(item), "name") ?? stageLabel,
|
||||
toNode: target,
|
||||
reason: kind === "jump" ? "jump" : "condition",
|
||||
condition: null,
|
||||
isJump: kind === "jump",
|
||||
};
|
||||
}
|
||||
return latest;
|
||||
}
|
||||
|
||||
const STAGE_OUTCOMES: ReadonlySet<string> = new Set(Object.values(StageOutcome));
|
||||
|
||||
/** The fork's branches as the projection carries them (`parallel_results`). */
|
||||
export function parallelOverviewFromProjection(
|
||||
stage: StageProjection | undefined,
|
||||
): ParallelOverview {
|
||||
const results = (stage?.parallel_results ?? [])
|
||||
.map((result) => {
|
||||
const status = STAGE_OUTCOMES.has(result.status) ? (result.status as StageOutcome) : null;
|
||||
if (!status) return null;
|
||||
return {
|
||||
id: result.id,
|
||||
index: result.index ?? null,
|
||||
itemLabel: result.item_label ?? null,
|
||||
status,
|
||||
};
|
||||
})
|
||||
.filter((r): r is NonNullable<typeof r> => r != null);
|
||||
return { branchCount: results.length > 0 ? results.length : null, results };
|
||||
}
|
||||
|
||||
/** The fan-in's reducer prompt and response, from the projection. */
|
||||
export function reducerTranscriptFromProjection(
|
||||
stage: StageProjection | undefined,
|
||||
): ReducerTranscript | null {
|
||||
if (!stage?.prompt && !stage?.response) return null;
|
||||
const tokens = stage.usage?.tokens;
|
||||
return {
|
||||
prompt: stage.prompt ?? "",
|
||||
response: stage.response ?? "",
|
||||
model: stage.provider_used?.model ?? stage.model?.model_id ?? null,
|
||||
inputTokens: tokens?.input ?? 0,
|
||||
outputTokens: tokens?.output ?? 0,
|
||||
};
|
||||
}
|
||||
|
||||
const ENGINE_CONTEXT_KEYS = new Set(["last_stage", "last_response", "command.output"]);
|
||||
const ENGINE_CONTEXT_PREFIXES = ["response.", "internal.", "current.", "human.gate.", "parallel."];
|
||||
|
||||
function isEngineContextKey(key: string): boolean {
|
||||
if (ENGINE_CONTEXT_KEYS.has(key)) return true;
|
||||
return ENGINE_CONTEXT_PREFIXES.some((prefix) => key.startsWith(prefix));
|
||||
}
|
||||
|
||||
/**
|
||||
* The workflow's deliberate outputs from the stage's final `step.finished`:
|
||||
* its `outcome.context_updates` minus the engine's keys.
|
||||
*/
|
||||
export function extractPetriStageContext(items: PetriStream): StageContextData | null {
|
||||
let latest: StageContextData | null = null;
|
||||
for (const item of items) {
|
||||
if (petriEventName(item) !== "step.finished") continue;
|
||||
if (getBool(derived(item), "final") === false) continue;
|
||||
const outcome = getObject(petriBody(item), "outcome");
|
||||
const rawUpdates = getObject(outcome, "context_updates") ?? {};
|
||||
const updates: Record<string, unknown> = {};
|
||||
for (const [key, value] of Object.entries(rawUpdates)) {
|
||||
if (!isEngineContextKey(key)) updates[key] = value;
|
||||
}
|
||||
if (Object.keys(updates).length === 0) {
|
||||
latest = null;
|
||||
continue;
|
||||
}
|
||||
latest = { routing: { preferredLabel: null, suggestedNextIds: [] }, updates };
|
||||
}
|
||||
return latest;
|
||||
}
|
||||
|
||||
/** A Pebble `CodingAgentEvent` envelope a stage's step recorded. */
|
||||
export interface PetriAgentEnvelope {
|
||||
ts: string;
|
||||
streamSeq: number;
|
||||
/** The Pebble variant name, e.g. `AssistantMessage`. */
|
||||
variant: string;
|
||||
/** The variant's fields. */
|
||||
payload: UnknownRecord;
|
||||
sessionId: string | null;
|
||||
parentSessionId: string | null;
|
||||
}
|
||||
|
||||
/**
|
||||
* The backend envelopes among a stage's items: a `step.progress.recorded`
|
||||
* whose custom payload carries a string `kind` (the backend) and an `event`
|
||||
* object (Pebble's externally tagged `CodingAgentEvent`).
|
||||
*/
|
||||
export function agentEnvelopesOf(items: PetriStream): PetriAgentEnvelope[] {
|
||||
const out: PetriAgentEnvelope[] = [];
|
||||
for (const item of items) {
|
||||
if (petriEventName(item) !== "step.progress.recorded") continue;
|
||||
const custom = getObject(getObject(petriBody(item), "ev"), "custom");
|
||||
const event = getObject(custom, "event");
|
||||
if (!custom || !event || !getString(custom, "kind")) continue;
|
||||
let variant: string | null = null;
|
||||
let payload: UnknownRecord = {};
|
||||
for (const [key, value] of Object.entries(event)) {
|
||||
variant = key;
|
||||
payload = isRecord(value) ? value : {};
|
||||
break;
|
||||
}
|
||||
if (!variant) continue;
|
||||
out.push({
|
||||
ts: streamItemTs(item),
|
||||
streamSeq: item.stream_seq,
|
||||
variant,
|
||||
payload,
|
||||
sessionId: getString(custom, "session_id") ?? null,
|
||||
parentSessionId: getString(custom, "parent_session_id") ?? null,
|
||||
});
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** A command stage's script, from its `step.started` record, if recorded. */
|
||||
export function commandScriptOf(items: PetriStream): string | null {
|
||||
for (const item of items) {
|
||||
if (petriEventName(item) !== "step.started") continue;
|
||||
const script =
|
||||
getString(getObject(petriBody(item), "config"), "script") ??
|
||||
getString(getObject(getObject(subjectNode(item), "meta"), "config"), "script");
|
||||
if (script) return script;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/** The exit code and duration of the stage's final `step.finished`. */
|
||||
export function commandOutcomeOf(items: PetriStream): {
|
||||
exitCode: number | null;
|
||||
durationMs: number;
|
||||
} {
|
||||
let exitCode: number | null = null;
|
||||
let durationMs = 0;
|
||||
for (const item of items) {
|
||||
if (petriEventName(item) !== "step.finished") continue;
|
||||
const metrics = getObject(getObject(petriBody(item), "outcome"), "metrics");
|
||||
exitCode = getNumber(metrics, "exit_code") ?? exitCode;
|
||||
durationMs = getNumber(metrics, "duration_ms") ?? durationMs;
|
||||
}
|
||||
return { exitCode, durationMs };
|
||||
}
|
||||
|
||||
// ── Stages from the projection ──────────────────────────────────────────
|
||||
|
||||
const STAGE_STATES: ReadonlySet<string> = new Set(Object.values(StageState));
|
||||
|
||||
/**
|
||||
* The sidebar stages a projection describes, sorted by their first event.
|
||||
* The API serves the same rows through `/runs/{id}/stages`; this derivation
|
||||
* lets a view (and a test) build them from the projection alone.
|
||||
*/
|
||||
export function stagesFromProjection(projection: RunProjection): Stage[] {
|
||||
const stages: Stage[] = [];
|
||||
for (const [id, stage] of Object.entries(projection.stages ?? {})) {
|
||||
const at = id.lastIndexOf("@");
|
||||
const name = at > 0 ? id.slice(0, at) : id;
|
||||
const visit = at > 0 ? Number.parseInt(id.slice(at + 1), 10) || 1 : 1;
|
||||
const branch = stage.parallel_branch_id ?? null;
|
||||
const branchAt = branch ? branch.lastIndexOf(":") : -1;
|
||||
const status = STAGE_STATES.has(stage.state) ? stage.state : StageState.PENDING;
|
||||
stages.push({
|
||||
id,
|
||||
name,
|
||||
handler: (stage as { handler?: Stage["handler"] }).handler ?? "agent",
|
||||
nodeId: name,
|
||||
visit,
|
||||
graphVisit: (stage as { graph_visit?: number | null }).graph_visit ?? null,
|
||||
resumedFromStageId: null,
|
||||
parallelGroupId: branch && branchAt > 0 ? branch.slice(0, branchAt) : null,
|
||||
parallelBranchIndex:
|
||||
branch && branchAt > 0 ? Number.parseInt(branch.slice(branchAt + 1), 10) : null,
|
||||
status,
|
||||
duration:
|
||||
stage.timing?.wall_time_ms != null ? formatDurationMs(stage.timing.wall_time_ms) : "--",
|
||||
startedAt: stage.started_at ?? null,
|
||||
providerUsed: stage.provider_used ?? null,
|
||||
usage: stage.usage,
|
||||
firstEventSeq: stage.first_event_seq,
|
||||
} as Stage & { firstEventSeq: number });
|
||||
}
|
||||
stages.sort(
|
||||
(a, b) =>
|
||||
((a as Stage & { firstEventSeq: number }).firstEventSeq ?? 0) -
|
||||
((b as Stage & { firstEventSeq: number }).firstEventSeq ?? 0),
|
||||
);
|
||||
return stages.map(({ firstEventSeq: _, ...stage }: Stage & { firstEventSeq?: number }) => stage);
|
||||
}
|
||||
|
|
@ -12,6 +12,7 @@ import type {
|
|||
Environment,
|
||||
EnvironmentListResponse,
|
||||
EventEnvelope,
|
||||
ListRunEvents200Response,
|
||||
ListRunsDirectionEnum,
|
||||
ListRunsSortEnum,
|
||||
McpServer,
|
||||
|
|
@ -27,6 +28,7 @@ import type {
|
|||
RunArtifactListResponse,
|
||||
RunProjection,
|
||||
Run,
|
||||
RunStreamItem,
|
||||
RunUsage,
|
||||
SandboxDetails,
|
||||
SecretListResponse,
|
||||
|
|
@ -387,13 +389,76 @@ export function useRunStageContextWindow(
|
|||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* `GET /runs/{id}/events` answers in the run engine's envelope: a legacy run
|
||||
* pages `EventEnvelope`s by `since_seq`, a Petri run pages `RunStreamItem`s
|
||||
* by `after`. The stream page names the Petri event contract version it
|
||||
* follows; the legacy page never does.
|
||||
*/
|
||||
function isRunStreamPage(
|
||||
page: ListRunEvents200Response,
|
||||
): page is Extract<ListRunEvents200Response, { event_contract_version: number }> {
|
||||
return "event_contract_version" in page;
|
||||
}
|
||||
|
||||
export function useRunEventsList(id: string | undefined) {
|
||||
return useSWR<EventEnvelope[]>(
|
||||
id ? queryKeys.runs.events(id, 1000) : null,
|
||||
() =>
|
||||
fetchAllStageEvents(`run ${id} events`, (sinceSeq, limit) =>
|
||||
apiData(() => runInternalsApi.listRunEvents(id!, sinceSeq, limit)),
|
||||
fetchAllStageEvents(`run ${id} events`, async (sinceSeq, limit) => {
|
||||
const page = await apiData(() =>
|
||||
runInternalsApi.listRunEvents(id!, sinceSeq, limit),
|
||||
);
|
||||
// A Petri run's events are a stream, read by `useRunStream`; the
|
||||
// legacy list of such a run is empty.
|
||||
if (isRunStreamPage(page)) {
|
||||
return { data: [], meta: { has_more: false } };
|
||||
}
|
||||
return page;
|
||||
}),
|
||||
);
|
||||
}
|
||||
|
||||
const STREAM_PAGE_LIMIT = 1000;
|
||||
const STREAM_MAX_PAGES = 50;
|
||||
|
||||
/**
|
||||
* Every item of a Petri run's stream, paged by `after` (the last
|
||||
* `stream_seq` seen). A legacy run has no stream: its page comes back in
|
||||
* the legacy envelope and reads as empty here.
|
||||
*/
|
||||
async function fetchRunStream(id: string): Promise<RunStreamItem[]> {
|
||||
const items: RunStreamItem[] = [];
|
||||
let after = 0;
|
||||
for (let pages = 0; pages < STREAM_MAX_PAGES; pages += 1) {
|
||||
const page = await apiData(() =>
|
||||
runInternalsApi.listRunEvents(
|
||||
id,
|
||||
undefined,
|
||||
STREAM_PAGE_LIMIT,
|
||||
undefined,
|
||||
undefined,
|
||||
after,
|
||||
),
|
||||
);
|
||||
if (!isRunStreamPage(page)) return items;
|
||||
if (page.data.length === 0) return items;
|
||||
items.push(...page.data);
|
||||
const last = page.data[page.data.length - 1];
|
||||
if (!page.meta.has_more || last.stream_seq <= after) return items;
|
||||
after = last.stream_seq;
|
||||
}
|
||||
console.warn(
|
||||
`Stopped run stream fetch for ${id} after ${STREAM_MAX_PAGES} pages and ${items.length} items because the safety cap was reached.`,
|
||||
);
|
||||
return items;
|
||||
}
|
||||
|
||||
/** A Petri run's stream: Petri's events and Fabro's platform records, in order. */
|
||||
export function useRunStream(id: string | undefined) {
|
||||
return useSWR<RunStreamItem[]>(
|
||||
id ? queryKeys.runs.stream(id) : null,
|
||||
() => fetchRunStream(id!),
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -61,6 +61,8 @@ export const queryKeys = {
|
|||
questions: (id: string, limit = 1, offset = 0) =>
|
||||
["runs", "questions", id, limit, offset] as const,
|
||||
events: (id: string, limit = 1000) => ["runs", "events", id, limit] as const,
|
||||
/** A Petri run's stream: every `RunStreamItem` in `stream_seq` order. */
|
||||
stream: (id: string) => ["runs", "stream", id] as const,
|
||||
stageEvents: (id: string, stageId: string) =>
|
||||
["runs", "stage-events", id, stageId] as const,
|
||||
stageContextWindow: (id: string, stageId: string) =>
|
||||
|
|
|
|||
|
|
@ -1,11 +1,21 @@
|
|||
import { useEffect } from "react";
|
||||
import type { RunStreamItem } from "@qltysh/fabro-api-client";
|
||||
import { useSWRConfig, type Key } from "swr";
|
||||
|
||||
import {
|
||||
subscribeToCrossTabSse,
|
||||
type CrossTabSseCoordinator,
|
||||
} from "./cross-tab-sse";
|
||||
import {
|
||||
isStreamItemPayload,
|
||||
isTerminalLifecycleItem,
|
||||
petriEventName,
|
||||
petriParsed,
|
||||
petriStageLabel,
|
||||
platformRecordKind,
|
||||
} from "./petri-stream";
|
||||
import { queryKeys } from "./query-keys";
|
||||
import { getString } from "./unknown";
|
||||
import {
|
||||
createBrowserEventSource,
|
||||
subscribeToSharedEventSource,
|
||||
|
|
@ -23,6 +33,10 @@ export interface RunEventPayload extends EventPayload {
|
|||
node_id?: string;
|
||||
stage_id?: string;
|
||||
properties?: Record<string, unknown>;
|
||||
/** Set on a Petri run stream item, which is invalidated by its own rules. */
|
||||
stream_seq?: number;
|
||||
kind?: string;
|
||||
item?: unknown;
|
||||
}
|
||||
|
||||
interface RunEventOptions {
|
||||
|
|
@ -322,6 +336,151 @@ export function queryKeysForRunEvent(
|
|||
return [];
|
||||
}
|
||||
|
||||
/**
|
||||
* The SWR keys a Petri run stream item invalidates. Petri's events are
|
||||
* named `<subject>.<verb>`; a platform record by its `kind`. The stage
|
||||
* keys use the subject's `node@visit` label, which is the stage id the
|
||||
* projection keys stages by.
|
||||
*/
|
||||
export function queryKeysForStreamItem(
|
||||
runId: string,
|
||||
item: RunStreamItem,
|
||||
): { keys: Key[]; immediate: boolean } {
|
||||
const stageId = petriStageLabel(item);
|
||||
const stageKeys: Key[] = stageId
|
||||
? [queryKeys.runs.stageEvents(runId, stageId), queryKeys.runs.stageContextWindow(runId, stageId)]
|
||||
: [];
|
||||
const stream = queryKeys.runs.stream(runId);
|
||||
|
||||
if (item.kind === "platform") {
|
||||
const kind = platformRecordKind(item);
|
||||
if (isTerminalLifecycleItem(item)) {
|
||||
return { keys: terminalKeys(runId, stream), immediate: true };
|
||||
}
|
||||
switch (kind) {
|
||||
case "checkpoint":
|
||||
return {
|
||||
keys: [
|
||||
...queryKeys.runs.filesAllScopes(runId),
|
||||
queryKeys.runs.commits(runId),
|
||||
queryKeys.runs.state(runId),
|
||||
stream,
|
||||
],
|
||||
immediate: false,
|
||||
};
|
||||
case "interview.answered":
|
||||
return {
|
||||
keys: [queryKeys.runs.questions(runId, 25, 0), queryKeys.runs.detail(runId), stream],
|
||||
immediate: false,
|
||||
};
|
||||
default:
|
||||
return { keys: [queryKeys.runs.detail(runId), queryKeys.runs.state(runId), stream], immediate: false };
|
||||
}
|
||||
}
|
||||
|
||||
const name = petriEventName(item);
|
||||
switch (name) {
|
||||
case "run.finished":
|
||||
return { keys: terminalKeys(runId, stream), immediate: false };
|
||||
case "run.started":
|
||||
case "run.paused":
|
||||
case "run.unpaused":
|
||||
case "invocation.finished":
|
||||
case "invocation.cancel.requested":
|
||||
case "run.stalled":
|
||||
return { keys: [queryKeys.runs.detail(runId), queryKeys.runs.state(runId), stream], immediate: false };
|
||||
case "visit.started":
|
||||
case "visit.completed":
|
||||
case "retry.scheduled":
|
||||
case "wait.state.changed":
|
||||
case "admission.decided":
|
||||
return {
|
||||
keys: [
|
||||
queryKeys.runs.stages(runId),
|
||||
queryKeys.runs.state(runId),
|
||||
queryKeys.runs.detail(runId),
|
||||
stream,
|
||||
queryKeys.runs.graph(runId, "LR"),
|
||||
queryKeys.runs.graph(runId, "TB"),
|
||||
...stageKeys,
|
||||
],
|
||||
immediate: false,
|
||||
};
|
||||
case "step.progress.recorded": {
|
||||
const parsed = getString(petriParsed(item), "kind");
|
||||
if (parsed === "question" || parsed === "question_expired") {
|
||||
return {
|
||||
keys: [
|
||||
queryKeys.runs.questions(runId, 25, 0),
|
||||
queryKeys.runs.detail(runId),
|
||||
queryKeys.runs.state(runId),
|
||||
stream,
|
||||
...stageKeys,
|
||||
],
|
||||
immediate: false,
|
||||
};
|
||||
}
|
||||
return { keys: [queryKeys.runs.state(runId), stream, ...stageKeys], immediate: false };
|
||||
}
|
||||
case "control.requested":
|
||||
return {
|
||||
keys: [
|
||||
queryKeys.runs.questions(runId, 25, 0),
|
||||
queryKeys.runs.detail(runId),
|
||||
queryKeys.runs.state(runId),
|
||||
stream,
|
||||
...stageKeys,
|
||||
],
|
||||
immediate: false,
|
||||
};
|
||||
case "step.finished":
|
||||
return {
|
||||
keys: [
|
||||
queryKeys.runs.state(runId),
|
||||
queryKeys.runs.usage(runId),
|
||||
queryKeys.runs.stages(runId),
|
||||
queryKeys.runs.detail(runId),
|
||||
stream,
|
||||
...stageKeys,
|
||||
],
|
||||
immediate: false,
|
||||
};
|
||||
case "fork.started":
|
||||
case "branch.completed":
|
||||
case "fork.completed":
|
||||
case "node.expanded":
|
||||
return {
|
||||
keys: [
|
||||
queryKeys.runs.stages(runId),
|
||||
queryKeys.runs.state(runId),
|
||||
stream,
|
||||
queryKeys.runs.graph(runId, "LR"),
|
||||
queryKeys.runs.graph(runId, "TB"),
|
||||
],
|
||||
immediate: false,
|
||||
};
|
||||
case "routing.resolved":
|
||||
case "route.applied":
|
||||
return { keys: [stream, ...stageKeys], immediate: false };
|
||||
default:
|
||||
return { keys: [stream], immediate: false };
|
||||
}
|
||||
}
|
||||
|
||||
function terminalKeys(runId: string, stream: Key): Key[] {
|
||||
return [
|
||||
queryKeys.runs.detail(runId),
|
||||
queryKeys.runs.state(runId),
|
||||
...queryKeys.runs.filesAllScopes(runId),
|
||||
queryKeys.runs.commits(runId),
|
||||
queryKeys.runs.usage(runId),
|
||||
queryKeys.runs.stages(runId),
|
||||
stream,
|
||||
queryKeys.runs.graph(runId, "LR"),
|
||||
queryKeys.runs.graph(runId, "TB"),
|
||||
];
|
||||
}
|
||||
|
||||
export function subscribeToRunEvents(
|
||||
runId: string,
|
||||
mutate: MutateFn,
|
||||
|
|
@ -357,6 +516,9 @@ export function subscribeToRunEvents(
|
|||
}
|
||||
|
||||
function runInvalidation(runId: string, payload: RunEventPayload) {
|
||||
if (isStreamItemPayload(payload)) {
|
||||
return queryKeysForStreamItem(runId, payload);
|
||||
}
|
||||
const event = payload.event;
|
||||
if (!event) return { keys: [], immediate: false };
|
||||
|
||||
|
|
@ -375,6 +537,7 @@ function resyncKeysForRun(runId: string) {
|
|||
queryKeys.runs.usage(runId),
|
||||
queryKeys.runs.stages(runId),
|
||||
queryKeys.runs.events(runId, 1000),
|
||||
queryKeys.runs.stream(runId),
|
||||
queryKeys.runs.graph(runId, "LR"),
|
||||
queryKeys.runs.graph(runId, "TB"),
|
||||
queryKeys.runs.questions(runId, 25, 0),
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
import { useMemo, useState } from "react";
|
||||
import { useParams, useSearchParams } from "react-router";
|
||||
import type { EventEnvelope } from "@qltysh/fabro-api-client";
|
||||
import type { EventEnvelope, RunStreamItem } from "@qltysh/fabro-api-client";
|
||||
|
||||
import {
|
||||
DebugEventDetailsPanel,
|
||||
|
|
@ -13,9 +13,23 @@ import {
|
|||
debugCategoryLabel,
|
||||
} from "../components/event-debug-helpers";
|
||||
import { RunWaterfall } from "../components/run-waterfall";
|
||||
import type { RunPhase } from "../lib/run-phases";
|
||||
import { StageSidebar } from "../components/stage-sidebar";
|
||||
import { EmptyState, ErrorState, LoadingState } from "../components/state";
|
||||
import { useRun, useRunEventsList, useRunStages } from "../lib/queries";
|
||||
import {
|
||||
debugRowSearchText,
|
||||
debugRowsFromStream,
|
||||
deriveRunPhasesFromStream,
|
||||
isPetriRun,
|
||||
type DebugRow,
|
||||
} from "../lib/petri-stream";
|
||||
import {
|
||||
useRun,
|
||||
useRunEventsList,
|
||||
useRunStages,
|
||||
useRunState,
|
||||
useRunStream,
|
||||
} from "../lib/queries";
|
||||
import { mapRunStagesToSidebarStages } from "../lib/stage-sidebar";
|
||||
|
||||
export const handle = { wide: true, fullHeight: true };
|
||||
|
|
@ -23,12 +37,34 @@ export const handle = { wide: true, fullHeight: true };
|
|||
type ViewMode = "waterfall" | "events";
|
||||
|
||||
const EMPTY_EVENTS: EventEnvelope[] = [];
|
||||
const EMPTY_STREAM: RunStreamItem[] = [];
|
||||
const EMPTY_ROWS: DebugRow[] = [];
|
||||
|
||||
export default function RunEvents() {
|
||||
const { id } = useParams();
|
||||
const runQuery = useRun(id);
|
||||
const runStateQuery = useRunState(id);
|
||||
const stagesQuery = useRunStages(id);
|
||||
const eventsQuery = useRunEventsList(id);
|
||||
// A Petri run's events are its stream; a legacy run's the stored events.
|
||||
// Which one is known from the run's spec, so the other query stays idle.
|
||||
// A state that cannot be read leaves the engine unknown; the legacy list
|
||||
// then loads as it did before the engine existed.
|
||||
const engineKnown =
|
||||
runStateQuery.data !== undefined || runStateQuery.error !== undefined;
|
||||
const petri = isPetriRun(runStateQuery.data);
|
||||
const eventsQuery = useRunEventsList(engineKnown && !petri ? id : undefined);
|
||||
const streamQuery = useRunStream(engineKnown && petri ? id : undefined);
|
||||
const streamPhases = useMemo(
|
||||
() =>
|
||||
petri && streamQuery.data && runQuery.data
|
||||
? deriveRunPhasesFromStream(streamQuery.data, runQuery.data.timestamps.created_at)
|
||||
: undefined,
|
||||
[petri, streamQuery.data, runQuery.data],
|
||||
);
|
||||
const streamRows = useMemo(
|
||||
() => (petri && streamQuery.data ? debugRowsFromStream(streamQuery.data) : undefined),
|
||||
[petri, streamQuery.data],
|
||||
);
|
||||
const [searchParams, setSearchParams] = useSearchParams();
|
||||
const view: ViewMode = searchParams.get("view") === "events" ? "events" : "waterfall";
|
||||
const setView = (next: ViewMode) => {
|
||||
|
|
@ -63,19 +99,32 @@ export default function RunEvents() {
|
|||
{view === "waterfall" ? (
|
||||
<WaterfallPane
|
||||
runId={id!}
|
||||
events={eventsQuery.data}
|
||||
eventsError={eventsQuery.error}
|
||||
events={petri ? (streamQuery.data ? EMPTY_EVENTS : undefined) : eventsQuery.data}
|
||||
phases={streamPhases}
|
||||
eventsError={petri ? streamQuery.error : eventsQuery.error}
|
||||
stagesData={stagesQuery.data}
|
||||
stagesError={stagesQuery.error}
|
||||
createdAt={runQuery.data?.timestamps.created_at}
|
||||
completedAt={runQuery.data?.timestamps.completed_at ?? null}
|
||||
onRetry={() => {
|
||||
void eventsQuery.mutate();
|
||||
void (petri ? streamQuery.mutate() : eventsQuery.mutate());
|
||||
void stagesQuery.mutate();
|
||||
}}
|
||||
view={view}
|
||||
onChangeView={setView}
|
||||
/>
|
||||
) : petri ? (
|
||||
<StreamEventsView
|
||||
rows={streamRows}
|
||||
error={streamQuery.error}
|
||||
onRetry={() => void streamQuery.mutate()}
|
||||
runStart={
|
||||
runQuery.data?.timestamps.started_at ??
|
||||
runQuery.data?.timestamps.created_at
|
||||
}
|
||||
view={view}
|
||||
onChangeView={setView}
|
||||
/>
|
||||
) : (
|
||||
<EventsView
|
||||
events={eventsQuery.data}
|
||||
|
|
@ -131,6 +180,7 @@ function ViewToggle({
|
|||
function WaterfallPane({
|
||||
runId,
|
||||
events,
|
||||
phases,
|
||||
eventsError,
|
||||
stagesData,
|
||||
stagesError,
|
||||
|
|
@ -142,6 +192,7 @@ function WaterfallPane({
|
|||
}: {
|
||||
runId: string;
|
||||
events: EventEnvelope[] | undefined;
|
||||
phases?: RunPhase[];
|
||||
eventsError: unknown;
|
||||
stagesData: ReturnType<typeof useRunStages>["data"];
|
||||
stagesError: unknown;
|
||||
|
|
@ -178,6 +229,7 @@ function WaterfallPane({
|
|||
<RunWaterfall
|
||||
runId={runId}
|
||||
events={events!}
|
||||
phases={phases}
|
||||
stages={stagesData!.data ?? []}
|
||||
createdAtIso={createdAt!}
|
||||
completedAtIso={completedAt}
|
||||
|
|
@ -332,6 +384,196 @@ function EventsView({
|
|||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* The events list of a Petri run: one row per stream item, named by the
|
||||
* Petri event (`<subject>.<verb>`) or the platform record kind, with the
|
||||
* raw item in the details panel.
|
||||
*/
|
||||
export function StreamEventsView({
|
||||
rows,
|
||||
error,
|
||||
onRetry,
|
||||
runStart,
|
||||
view,
|
||||
onChangeView,
|
||||
}: {
|
||||
rows: DebugRow[] | undefined;
|
||||
error: unknown;
|
||||
onRetry: () => void;
|
||||
runStart: string | undefined;
|
||||
view: ViewMode;
|
||||
onChangeView: (v: ViewMode) => void;
|
||||
}) {
|
||||
const [openSeq, setOpenSeq] = useState<number | null>(null);
|
||||
const [selectedCategories, setSelectedCategories] = useState<string[]>([]);
|
||||
const [search, setSearch] = useState("");
|
||||
|
||||
const all = rows ?? EMPTY_ROWS;
|
||||
|
||||
const availableCategories = useMemo<string[]>(() => {
|
||||
const set = new Set<string>();
|
||||
for (const row of all) set.add(row.category);
|
||||
return Array.from(set).sort();
|
||||
}, [all]);
|
||||
|
||||
const filtered = useMemo<DebugRow[]>(() => {
|
||||
const useCategoryFilter = selectedCategories.length > 0;
|
||||
const cats = new Set(selectedCategories);
|
||||
const needle = search.toLowerCase();
|
||||
return all.filter((row) => {
|
||||
if (useCategoryFilter && !cats.has(row.category)) return false;
|
||||
if (needle && !debugRowSearchText(row).includes(needle)) return false;
|
||||
return true;
|
||||
});
|
||||
}, [all, selectedCategories, search]);
|
||||
|
||||
const openRow = useMemo<DebugRow | null>(
|
||||
() => (openSeq != null ? all.find((row) => row.seq === openSeq) ?? null : null),
|
||||
[all, openSeq],
|
||||
);
|
||||
const openPayload = useMemo(
|
||||
() =>
|
||||
openRow
|
||||
? {
|
||||
event: openRow.event,
|
||||
stream_seq: openRow.seq,
|
||||
kind: openRow.item.kind,
|
||||
stage: openRow.stageLabel,
|
||||
recorded_at: openRow.ts,
|
||||
item: openRow.item.item,
|
||||
}
|
||||
: null,
|
||||
[openRow],
|
||||
);
|
||||
|
||||
const allCategoriesSelected =
|
||||
selectedCategories.length === 0 ||
|
||||
selectedCategories.length === availableCategories.length;
|
||||
const isFiltering = !allCategoriesSelected || search.length > 0;
|
||||
|
||||
function clearFilters() {
|
||||
setSelectedCategories([]);
|
||||
setSearch("");
|
||||
}
|
||||
|
||||
if (error) {
|
||||
return (
|
||||
<div className="min-w-0 flex-1 pt-3">
|
||||
<ErrorState
|
||||
title="Couldn't load events"
|
||||
description={errorMessage(error)}
|
||||
onRetry={onRetry}
|
||||
/>
|
||||
</div>
|
||||
);
|
||||
}
|
||||
if (rows === undefined) {
|
||||
return (
|
||||
<div className="min-w-0 flex-1 pt-3">
|
||||
<LoadingState label="Loading events…" />
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
return (
|
||||
<>
|
||||
<div className="flex min-h-0 min-w-0 flex-1 flex-col pt-3">
|
||||
<div className="shrink-0 border-b border-line">
|
||||
<div className="pl-3 pr-4 sm:pr-6 lg:pr-8">
|
||||
<div className="flex flex-wrap items-center gap-x-3 gap-y-2 pb-3">
|
||||
<div className="flex flex-1 flex-wrap items-center gap-2">
|
||||
<ViewToggle value={view} onChange={onChangeView} />
|
||||
<MultiSelectFilter<string>
|
||||
selected={selectedCategories}
|
||||
options={availableCategories}
|
||||
labelOf={debugCategoryLabel}
|
||||
onChange={setSelectedCategories}
|
||||
emptyMeansAll
|
||||
/>
|
||||
<EventSearchInput value={search} onChange={setSearch} />
|
||||
{isFiltering && (
|
||||
<button
|
||||
type="button"
|
||||
onClick={clearFilters}
|
||||
className="rounded px-2 py-1 text-xs text-fg-muted transition-colors hover:bg-overlay hover:text-fg-2 focus-visible:outline-2 focus-visible:outline-offset-1 focus-visible:outline-teal-500"
|
||||
>
|
||||
Clear
|
||||
</button>
|
||||
)}
|
||||
</div>
|
||||
{all.length > 0 && (
|
||||
<span className="text-xs tabular-nums text-fg-muted">
|
||||
{isFiltering
|
||||
? `${filtered.length.toLocaleString()} of ${all.length.toLocaleString()} items`
|
||||
: `${all.length.toLocaleString()} items`}
|
||||
</span>
|
||||
)}
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
<div className="min-h-0 flex-1 overflow-y-auto pt-2 pb-[calc(1.5rem+var(--fabro-interview-dock-clearance,0px))]">
|
||||
{all.length === 0 ? (
|
||||
<div className="px-2 py-12">
|
||||
<EmptyState
|
||||
title="No events yet"
|
||||
description="Events will appear here as the run executes."
|
||||
/>
|
||||
</div>
|
||||
) : filtered.length === 0 ? (
|
||||
<div className="px-2 py-6 text-sm text-fg-muted">
|
||||
No events match these filters.
|
||||
</div>
|
||||
) : (
|
||||
filtered.map((row) => (
|
||||
<StreamEventRow
|
||||
key={`stream-${row.seq}`}
|
||||
row={row}
|
||||
runStart={runStart}
|
||||
selected={openSeq === row.seq}
|
||||
onSelect={() => setOpenSeq(row.seq)}
|
||||
/>
|
||||
))
|
||||
)}
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<DebugEventDetailsPanel event={openPayload} onClose={() => 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 (
|
||||
<div className="grid grid-cols-[1fr_auto] items-center">
|
||||
<DebugEventRow
|
||||
event={row}
|
||||
runStart={runStart}
|
||||
selected={selected}
|
||||
onSelect={onSelect}
|
||||
/>
|
||||
{row.stageLabel && (
|
||||
<span
|
||||
data-stage={row.stageLabel}
|
||||
className="pr-5 font-mono text-[11px] text-fg-muted"
|
||||
>
|
||||
{row.stageLabel}
|
||||
</span>
|
||||
)}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
function errorMessage(error: unknown): string | undefined {
|
||||
return error instanceof Error ? error.message : undefined;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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", () => ({
|
||||
|
|
|
|||
|
|
@ -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() {
|
|||
</div>
|
||||
|
||||
<div className="flex min-h-0 min-w-0 flex-1 flex-col gap-4 pb-[var(--fabro-interview-dock-clearance,0px)]">
|
||||
<div className="shrink-0">
|
||||
<div className="shrink-0 space-y-4">
|
||||
<RunSummaryPanel runId={id!} />
|
||||
<PlatformRecordsPanel runId={id!} />
|
||||
</div>
|
||||
{graphSvg === undefined && graphQuery.isLoading ? (
|
||||
<div className="flex-1" />
|
||||
|
|
|
|||
|
|
@ -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<string, PendingTool>();
|
||||
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<TurnType, { kind: "tool" }>;
|
||||
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<typeof extractStageContext>;
|
||||
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 (
|
||||
<div className="min-h-0 flex-1 overflow-y-auto pt-6 pb-[calc(1.5rem+var(--fabro-interview-dock-clearance,0px))]">
|
||||
{effectiveTab === "chat" ? (
|
||||
|
|
@ -2169,23 +2363,29 @@ function StageActivityBody({
|
|||
turn={commandTurn}
|
||||
/>
|
||||
) : renderer === "human" ? (
|
||||
<HumanQA stage={selectedStage} events={debugEvents} />
|
||||
<HumanQA stage={selectedStage} events={legacyEvents} pairs={petri?.pairs} />
|
||||
) : renderer === "conditional" ? (
|
||||
<ConditionalDecision
|
||||
stage={selectedStage}
|
||||
runEvents={runEvents}
|
||||
edge={petri?.edge}
|
||||
allStages={stages}
|
||||
runId={runId}
|
||||
/>
|
||||
) : renderer === "parallel" ? (
|
||||
<ParallelChildren
|
||||
stage={selectedStage}
|
||||
events={debugEvents}
|
||||
events={legacyEvents}
|
||||
overview={petri?.overview}
|
||||
runId={runId}
|
||||
allStages={stages}
|
||||
/>
|
||||
) : renderer === "fan_in" ? (
|
||||
<FanInResults stage={selectedStage} events={debugEvents} />
|
||||
<FanInResults
|
||||
stage={selectedStage}
|
||||
events={legacyEvents}
|
||||
reducer={petri?.reducer}
|
||||
/>
|
||||
) : renderer === "wait" ? (
|
||||
<WaitStatus stage={selectedStage} />
|
||||
) : (
|
||||
|
|
@ -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<RunStreamItem[]>(
|
||||
() => (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<EventEnvelope[]>(() => {
|
||||
const debugEvents = useMemo<DebugListItem[]>(() => {
|
||||
if (petri) return debugRowsFromStream(stageItems);
|
||||
return (stageEventsQuery.data ?? []).filter(
|
||||
(event) => activityEventStageId(event) === selectedStageId,
|
||||
);
|
||||
}, [stageEventsQuery.data, selectedStageId]);
|
||||
}, [petri, stageItems, stageEventsQuery.data, selectedStageId]);
|
||||
const petriViews = useMemo<PetriStageViews | undefined>(
|
||||
() =>
|
||||
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<EventsTab[]>(
|
||||
() =>
|
||||
|
|
@ -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<CommandTurn | null>(() => {
|
||||
if (effectiveTab !== "primary" || renderer !== "command") return null;
|
||||
|
|
@ -2347,34 +2577,27 @@ function RunStageActivityStage({
|
|||
}
|
||||
return null;
|
||||
}, [isPrimaryAgent, displayItems, panelSelection]);
|
||||
const openDebugEvent = useMemo<EventEnvelope | null>(
|
||||
() =>
|
||||
isDebug && openDebugSeq != null
|
||||
? (debugEvents.find((e) => e.seq === openDebugSeq) ?? null)
|
||||
: null,
|
||||
[isDebug, debugEvents, openDebugSeq],
|
||||
);
|
||||
const openDebugEvent = useMemo<EventDisplayPayload | null>(() => {
|
||||
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<DebugCategory[]>(() => {
|
||||
if (!isDebug) return [];
|
||||
const set = new Set<DebugCategory>();
|
||||
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<EventEnvelope[]>(() => {
|
||||
const filteredDebugEvents = useMemo<DebugListItem[]>(() => {
|
||||
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}
|
||||
/>
|
||||
</div>
|
||||
|
||||
|
|
@ -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<PetriRunData | undefined>(
|
||||
() =>
|
||||
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}
|
||||
/>
|
||||
</div>
|
||||
);
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue