From 3941a24a288c3140c3091037a2bd9ba7182145bb Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 01:03:19 -0400 Subject: [PATCH] Test the web app's Petri views over the captured scenario fixtures The Pebble envelopes a step records are read as `CodingAgentEvent`s (`{seq, stream_id, session_id, timestamp, event: {Variant}}`), so an agent stage's chat shows the `UserInput` prompt and the assistant's answer; a command step's exit status comes from its output; a user's login names who answered a gate. `lib/petri-stream.test.ts` checks the derivations over the hello, command, parallel and gate fixtures: stream density, names, stage labels with the fork's delegates skipped, the gate's question and answer with the principal, the run phases from the lifecycle records, the notice between the branches, the fork's branches from the projection, the edge a stage took, and the envelopes. `run-events.test.tsx` checks the SWR keys a stream item invalidates and the run filter on the coordinated stream. `run-petri.render.test.tsx` renders the stage list, the chat, the parallel children, the fan-in, the events list, the waterfall, the Q&A, the decision and the platform records for each fixture. Co-Authored-By: Claude Fable 5.1 --- apps/fabro-web/app/lib/petri-fixtures.ts | 25 ++ apps/fabro-web/app/lib/petri-stream.test.ts | 204 +++++++++++++ apps/fabro-web/app/lib/petri-stream.ts | 53 +++- apps/fabro-web/app/lib/run-events.test.tsx | 99 +++++++ .../app/routes/run-petri.render.test.tsx | 279 ++++++++++++++++++ apps/fabro-web/app/routes/run-stages.tsx | 12 +- 6 files changed, 659 insertions(+), 13 deletions(-) create mode 100644 apps/fabro-web/app/lib/petri-fixtures.ts create mode 100644 apps/fabro-web/app/lib/petri-stream.test.ts create mode 100644 apps/fabro-web/app/routes/run-petri.render.test.tsx diff --git a/apps/fabro-web/app/lib/petri-fixtures.ts b/apps/fabro-web/app/lib/petri-fixtures.ts new file mode 100644 index 000000000..b795fffe5 --- /dev/null +++ b/apps/fabro-web/app/lib/petri-fixtures.ts @@ -0,0 +1,25 @@ +/** + * The Petri scenario fixtures the server tests capture + * (`lib/apps/fabro-server/tests/it/scenario/petri_stream.rs` under + * `FABRO_CAPTURE_PETRI_FIXTURES`): a settled run's projection and its whole + * stream. Test-only; `tsc` excludes the tests that import this module. + */ +import { readFileSync } from "node:fs"; +import { dirname, join } from "node:path"; +import { fileURLToPath } from "node:url"; +import type { RunProjection, RunStreamItem } from "@qltysh/fabro-api-client"; + +export type PetriFixtureName = "hello" | "command" | "parallel" | "gate"; + +export interface PetriFixture { + run_id: string; + projection: RunProjection; + stream: RunStreamItem[]; +} + +const FIXTURES_DIR = join(dirname(fileURLToPath(import.meta.url)), "..", "test-fixtures", "petri"); + +export function loadPetriFixture(name: PetriFixtureName): PetriFixture { + const text = readFileSync(join(FIXTURES_DIR, `${name}.json`), "utf8"); + return JSON.parse(text) as PetriFixture; +} diff --git a/apps/fabro-web/app/lib/petri-stream.test.ts b/apps/fabro-web/app/lib/petri-stream.test.ts new file mode 100644 index 000000000..65f5ef074 --- /dev/null +++ b/apps/fabro-web/app/lib/petri-stream.test.ts @@ -0,0 +1,204 @@ +import { describe, expect, test } from "bun:test"; + +import { loadPetriFixture } from "./petri-fixtures"; +import { + agentEnvelopesOf, + commandOutcomeOf, + debugRowsFromStream, + deriveRunPhasesFromStream, + extractPetriStageContext, + findPetriEdgeForStage, + isPetriRun, + isStreamItemPayload, + isTerminalLifecycleItem, + itemsForStage, + parallelOverviewFromProjection, + parsePetriInterviewPairs, + petriEventName, + petriStageLabel, + platformRecordKind, + platformRecordsOf, + reducerTranscriptFromProjection, + stagesFromProjection, + streamItemName, +} from "./petri-stream"; + +const hello = loadPetriFixture("hello"); +const command = loadPetriFixture("command"); +const parallel = loadPetriFixture("parallel"); +const gate = loadPetriFixture("gate"); + +describe("stream items", () => { + test("a fixture run executes on Petri and its stream is dense", () => { + for (const fixture of [hello, command, parallel, gate]) { + expect(isPetriRun(fixture.projection)).toBe(true); + const seqs = fixture.stream.map((item) => item.stream_seq); + expect(seqs).toEqual(seqs.map((_, index) => index + 1)); + for (const item of fixture.stream) { + expect(isStreamItemPayload(item)).toBe(true); + expect(item.run_id).toBe(fixture.run_id); + } + } + expect(isPetriRun({ spec: { engine: { kind: "legacy" } } } as never)).toBe(false); + expect(isStreamItemPayload({ event: "run.completed", seq: 3 })).toBe(false); + }); + + test("a Petri item is named by its recorded event and a platform item by its kind", () => { + const names = command.stream.map(streamItemName); + expect(names[0]).toBe("run.created"); + expect(names).toContain("run.started"); + expect(names).toContain("visit.started"); + expect(names).toContain("step.finished"); + expect(names[names.length - 2]).toBe("run.finished"); + expect(names[names.length - 1]).toBe("run.lifecycle"); + const created = command.stream[0]; + expect(platformRecordKind(created)).toBe("run.created"); + expect(petriEventName(created)).toBeUndefined(); + }); + + test("the stage label is the subject's node@visit and skips the fork's delegates", () => { + const labels = new Set( + parallel.stream.map(petriStageLabel).filter((label): label is string => label != null), + ); + expect(labels).toEqual(new Set(["start@1", "fork@1", "a@1", "b@1", "merge@1", "exit@1"])); + // The parent execution holds a `parallel.branch` delegate named after + // each branch; only the child execution's own node is the stage. + const starts = parallel.stream.filter( + (item) => petriEventName(item) === "visit.started" && petriStageLabel(item) === "a@1", + ); + expect(starts).toHaveLength(1); + expect(itemsForStage(command.stream, "say@1").map(petriEventName)).toEqual([ + "visit.started", + "wait.state.changed", + "admission.decided", + "step.started", + "wait.state.changed", + "step.progress.recorded", + "step.finished", + "visit.completed", + "routing.resolved", + "route.applied", + "token.emitted", + ]); + }); +}); + +describe("questions", () => { + test("a gate's question pairs with its delivered answer and the answering principal", () => { + const pairs = parsePetriInterviewPairs(itemsForStage(gate.stream, "gate@1").concat( + gate.stream.filter((item) => platformRecordKind(item) === "interview.answered"), + )); + expect(pairs).toHaveLength(1); + const [pair] = pairs; + expect(pair.question.questionId).toBe("gate#2"); + expect(pair.question.question).toBe("Go?"); + expect(pair.question.questionType).toBe("yes_no"); + expect(pair.question.options.map((option) => option.key)).toEqual(["Y", "N"]); + expect(pair.question.allowFreeform).toBe(false); + expect(pair.resolution).toMatchObject({ kind: "answered", answer: "N", actor: "dev" }); + expect(pair.resolution?.kind === "answered" && pair.resolution.durationMs).toBeGreaterThan(0); + }); + + test("a run without a gate asks nothing", () => { + expect(parsePetriInterviewPairs(command.stream)).toEqual([]); + }); +}); + +describe("run phases", () => { + test("the phases come from the platform lifecycle records", () => { + const createdAt = parallel.projection.status_updated_at; + const created = parallel.stream[0]; + const phases = deriveRunPhasesFromStream( + parallel.stream, + new Date(created.recorded_at).toISOString(), + ); + expect(phases.map((phase) => phase.kind)).toEqual(["submitted", "runnable", "initializing"]); + for (const phase of phases) { + expect(phase.endMs).not.toBeNull(); + expect(phase.startMs).toBeLessThanOrEqual(phase.endMs!); + } + expect(createdAt).toBeDefined(); + const terminal = parallel.stream.filter(isTerminalLifecycleItem); + expect(terminal).toHaveLength(1); + expect(terminal[0]).toBe(parallel.stream[parallel.stream.length - 1]); + }); +}); + +describe("platform records", () => { + test("a notice recorded between the branches lists with its message", () => { + const records = platformRecordsOf(parallel.stream); + expect(records.map((record) => record.kind)).toEqual(["run.notice"]); + expect(records[0].detail).toBe("recorded while both branches ran"); + expect(records[0].stageKey).toBeNull(); + }); + + test("the debug rows name every item and carry its stage", () => { + const rows = debugRowsFromStream(parallel.stream); + expect(rows).toHaveLength(parallel.stream.length); + const notice = rows.find((row) => row.event === "run.notice"); + expect(notice?.category).toBe("platform"); + const started = rows.find((row) => row.event === "visit.started" && row.stageLabel === "b@1"); + expect(started?.category).toBe("petri"); + expect(rows.every((row) => !Number.isNaN(Date.parse(row.ts)))).toBe(true); + }); +}); + +describe("stage renderers", () => { + test("the fork's branches and results come from the projection", () => { + const fork = parallel.projection.stages["fork@1"]; + const overview = parallelOverviewFromProjection(fork); + expect(overview.branchCount).toBe(2); + expect(overview.results.map((result) => [result.id, result.index, result.status])).toEqual([ + ["a", 0, "succeeded"], + ["b", 1, "succeeded"], + ]); + const stages = stagesFromProjection(parallel.projection); + const branches = stages.filter((stage) => stage.parallelGroupId === "fork@1"); + expect(branches.map((stage) => [stage.id, stage.parallelBranchIndex])).toEqual([ + ["a@1", 0], + ["b@1", 1], + ]); + expect(stages.map((stage) => stage.id)).toEqual([ + "start@1", + "fork@1", + "a@1", + "b@1", + "merge@1", + "exit@1", + ]); + }); + + test("the fan-in with no reducer has no transcript", () => { + expect(reducerTranscriptFromProjection(parallel.projection.stages["merge@1"])).toBeNull(); + const greet = reducerTranscriptFromProjection(hello.projection.stages["greet@1"]); + expect(greet?.response).toBe("A haiku, added."); + }); + + test("the edge a stage took is its route.applied target", () => { + expect(findPetriEdgeForStage(gate.stream, "gate@1")).toEqual({ + fromNode: "gate", + toNode: "no", + reason: "condition", + condition: null, + isJump: false, + }); + expect(findPetriEdgeForStage(gate.stream, "exit@1")).toBeNull(); + }); + + test("a command stage's outcome is read from its final step.finished", () => { + const say = itemsForStage(command.stream, "say@1"); + expect(commandOutcomeOf(say).exitCode).toBe(0); + expect(extractPetriStageContext(say)).toBeNull(); + }); + + test("an agent stage's Pebble envelopes are read with their variant and session", () => { + const envelopes = agentEnvelopesOf(itemsForStage(hello.stream, "greet@1")); + expect(envelopes.length).toBeGreaterThan(0); + expect(envelopes[0].variant).toBe("SessionStarted"); + expect(envelopes[0].payload).toEqual({ provider: "openai", model: "gpt-5.4" }); + expect(envelopes.every((envelope) => envelope.sessionId?.startsWith("ses_"))).toBe(true); + const message = envelopes.find((envelope) => envelope.variant === "AssistantMessage"); + expect(message?.payload.text).toBe("A haiku, added."); + expect(agentEnvelopesOf(itemsForStage(command.stream, "say@1"))).toEqual([]); + }); +}); diff --git a/apps/fabro-web/app/lib/petri-stream.ts b/apps/fabro-web/app/lib/petri-stream.ts index 1087467b4..390590297 100644 --- a/apps/fabro-web/app/lib/petri-stream.ts +++ b/apps/fabro-web/app/lib/petri-stream.ts @@ -23,6 +23,7 @@ import type { StageContextData, } from "../components/stage-renderers/helpers"; import { principalLabel } from "../components/stage-renderers/helpers"; +import { principalDisplay } from "./principal-display"; import type { Stage } from "./stage-sidebar"; import type { RunPhase, RunPhaseKind } from "./run-phases"; import { formatDurationMs } from "./format"; @@ -180,6 +181,20 @@ function parseOptions(value: unknown): InterviewOption[] { return out; } +/** + * Who answered, from the `interview.answered` record's `Principal`: the + * user's login, or the legacy label for an actor shaped as the old events + * carried it. + */ +function answeringPrincipalLabel(principal: unknown): string | null { + if (!isRecord(principal)) return null; + const kind = getString(principal, "kind"); + if (kind === "user" && getString(principal, "login")) { + return principalDisplay(principal as unknown as Parameters[0]).label; + } + return principalLabel(principal); +} + function answerText(answer: UnknownRecord): string { const choice = getString(answer, "choice"); if (choice) return choice; @@ -214,7 +229,7 @@ export function parsePetriInterviewPairs(stream: PetriStream): HumanInterviewPai const rec = record(item); if (getString(rec, "kind") === "interview.answered") { const question = getString(rec, "question"); - if (question) actors.set(question, principalLabel(rec?.principal)); + if (question) actors.set(question, answeringPrincipalLabel(rec?.principal)); } continue; } @@ -541,7 +556,14 @@ export function reducerTranscriptFromProjection( }; } -const ENGINE_CONTEXT_KEYS = new Set(["last_stage", "last_response", "command.output"]); +// The command step's own bookkeeping (`command.output`, `failure_class`) +// joins the engine keys the legacy Context tab hides. +const ENGINE_CONTEXT_KEYS = new Set([ + "last_stage", + "last_response", + "command.output", + "failure_class", +]); const ENGINE_CONTEXT_PREFIXES = ["response.", "internal.", "current.", "human.gate.", "parallel."]; function isEngineContextKey(key: string): boolean { @@ -588,15 +610,18 @@ export interface PetriAgentEnvelope { /** * 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`). + * object, Pebble's `CodingAgentEvent` as recorded: `{seq, stream_id, + * session_id, parent_session_id?, timestamp, event: {Variant: {...}}}`. */ 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; + const envelope = getObject(custom, "event"); + if (!custom || !envelope || !getString(custom, "kind")) continue; + const event = getObject(envelope, "event"); + if (!event) continue; let variant: string | null = null; let payload: UnknownRecord = {}; for (const [key, value] of Object.entries(event)) { @@ -606,12 +631,12 @@ export function agentEnvelopesOf(items: PetriStream): PetriAgentEnvelope[] { } if (!variant) continue; out.push({ - ts: streamItemTs(item), + ts: getString(envelope, "timestamp") ?? streamItemTs(item), streamSeq: item.stream_seq, variant, payload, - sessionId: getString(custom, "session_id") ?? null, - parentSessionId: getString(custom, "parent_session_id") ?? null, + sessionId: getString(envelope, "session_id") ?? null, + parentSessionId: getString(envelope, "parent_session_id") ?? null, }); } return out; @@ -629,7 +654,10 @@ export function commandScriptOf(items: PetriStream): string | null { return null; } -/** The exit code and duration of the stage's final `step.finished`. */ +/** + * The exit code and duration of the stage's final `step.finished`: the + * command step's output carries `exit_status`, its metrics the duration. + */ export function commandOutcomeOf(items: PetriStream): { exitCode: number | null; durationMs: number; @@ -638,8 +666,11 @@ export function commandOutcomeOf(items: PetriStream): { 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; + const outcome = getObject(petriBody(item), "outcome"); + const output = getObject(outcome, "output"); + const metrics = getObject(outcome, "metrics"); + exitCode = + getNumber(output, "exit_status") ?? getNumber(metrics, "exit_code") ?? exitCode; durationMs = getNumber(metrics, "duration_ms") ?? durationMs; } return { exitCode, durationMs }; diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 794e183c8..356d8cc35 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -1,8 +1,10 @@ import { describe, expect, test } from "bun:test"; import type { Key } from "swr"; +import { loadPetriFixture } from "./petri-fixtures"; import { queryKeysForRunEvent, + queryKeysForStreamItem, subscribeToRunEvents, } from "./run-events"; import { @@ -227,6 +229,103 @@ describe("queryKeysForRunEvent", () => { }); }); +describe("queryKeysForStreamItem", () => { + const parallel = loadPetriFixture("parallel"); + const gate = loadPetriFixture("gate"); + const runId = "run-petri"; + const named = (name: string, stage?: string) => + parallel.stream.find((item) => { + const body = (item.item as { record?: { body?: { event?: string } } }).record?.body; + const derived = (item.item as { derived?: { event?: string } }).derived; + const subject = (item.item as { subject?: { node?: { name?: string } } }).subject; + return ( + (body?.event ?? derived?.event) === name && + (stage === undefined || subject?.node?.name === stage) + ); + })!; + + test("a stage's visit invalidates the stage list, the state, the stream and its stage keys", () => { + const { keys, immediate } = queryKeysForStreamItem(runId, named("visit.started", "merge")); + expect(immediate).toBe(false); + expect(keys).toEqual([ + queryKeys.runs.stages(runId), + queryKeys.runs.state(runId), + queryKeys.runs.detail(runId), + queryKeys.runs.stream(runId), + queryKeys.runs.graph(runId, "LR"), + queryKeys.runs.graph(runId, "TB"), + queryKeys.runs.stageEvents(runId, "merge@1"), + queryKeys.runs.stageContextWindow(runId, "merge@1"), + ]); + }); + + test("a platform notice refreshes the run summary; the terminal lifecycle record is immediate", () => { + const notice = parallel.stream.find( + (item) => item.kind === "platform" && (item.item as { record: { kind: string } }).record.kind === "run.notice", + )!; + expect(queryKeysForStreamItem(runId, notice)).toEqual({ + keys: [queryKeys.runs.detail(runId), queryKeys.runs.state(runId), queryKeys.runs.stream(runId)], + immediate: false, + }); + const terminal = parallel.stream[parallel.stream.length - 1]; + const result = queryKeysForStreamItem(runId, terminal); + expect(result.immediate).toBe(true); + expect(result.keys).toContainEqual(queryKeys.runs.usage(runId)); + expect(result.keys).toContainEqual(queryKeys.runs.stream(runId)); + }); + + test("a question and its answer refresh the questions list", () => { + const question = gate.stream.find( + (item) => (item.item as { derived?: { parsed?: { kind?: string } } }).derived?.parsed?.kind === "question", + )!; + expect(queryKeysForStreamItem(runId, question).keys[0]).toEqual( + queryKeys.runs.questions(runId, 25, 0), + ); + const answer = gate.stream.find( + (item) => (item.item as { record?: { body?: { event?: string } } }).record?.body?.event === "control.requested", + )!; + expect(queryKeysForStreamItem(runId, answer).keys).toContainEqual( + queryKeys.runs.stageEvents(runId, "gate@1"), + ); + }); + + test("a run stream item on the attach stream is invalidated by its own rules", async () => { + const source = new FakeEventSource(); + const keys: Key[] = []; + // The coordinated stream carries every run, so the item's `run_id` is + // what keeps another run's item from invalidating this one. + const coordinator = createCoordinator(() => source); + const cleanup = subscribeToRunEvents( + runId, + (key) => { + keys.push(key); + return Promise.resolve(); + }, + () => { + throw new Error("source should be created by coordinator"); + }, + { debounceMs: 0, coordinator }, + ); + await waitFor(() => source.onmessage !== null); + keys.length = 0; + source.emit({ ...named("step.finished", "a"), run_id: runId }); + expect(keys).toEqual([ + queryKeys.runs.state(runId), + queryKeys.runs.usage(runId), + queryKeys.runs.stages(runId), + queryKeys.runs.detail(runId), + queryKeys.runs.stream(runId), + queryKeys.runs.stageEvents(runId, "a@1"), + queryKeys.runs.stageContextWindow(runId, "a@1"), + ]); + keys.length = 0; + source.emit({ ...named("step.finished", "a"), run_id: "another-run" }); + expect(keys).toEqual([]); + cleanup(); + coordinator.close(); + }); +}); + describe("subscribeToRunEvents", () => { test("coordinated mode uses the global attach stream and filters by run_id", async () => { const source = new FakeEventSource(); diff --git a/apps/fabro-web/app/routes/run-petri.render.test.tsx b/apps/fabro-web/app/routes/run-petri.render.test.tsx new file mode 100644 index 000000000..8f660e546 --- /dev/null +++ b/apps/fabro-web/app/routes/run-petri.render.test.tsx @@ -0,0 +1,279 @@ +/** + * The run detail's views over a Petri run, rendered from the projection and + * the stream the server tests captured (`test-fixtures/petri/*.json`): one + * scenario per fixture, every view `VIEWS.md` lists that the web app draws + * from those two sources. + */ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import type { ReactElement } from "react"; +import type { RunStage } from "@qltysh/fabro-api-client"; +import TestRenderer, { act } from "react-test-renderer"; +import { MemoryRouter } from "react-router"; + +import { PlatformRecordsPanelView } from "../components/platform-records-panel"; +import { RunWaterfall } from "../components/run-waterfall"; +import { StageSidebar } from "../components/stage-sidebar"; +import { FanInResults } from "../components/stage-renderers/fan-in-results"; +import { HumanQA } from "../components/stage-renderers/human-qa"; +import { ParallelChildren } from "../components/stage-renderers/parallel-children"; +import { ConditionalDecision } from "../components/stage-renderers/conditional-decision"; +import { loadPetriFixture, type PetriFixture } from "../lib/petri-fixtures"; +import { + debugRowsFromStream, + deriveRunPhasesFromStream, + findPetriEdgeForStage, + itemsForStage, + parallelOverviewFromProjection, + parsePetriInterviewPairs, + platformRecordsOf, + reducerTranscriptFromProjection, + stagesFromProjection, +} from "../lib/petri-stream"; +import { setupReactTestEnv } from "../lib/test-utils"; +import { StreamEventsView } from "./run-events"; +import { StageChatView, buildPetriStageActivity } from "./run-stages"; + +let teardown: () => void; +beforeEach(() => { + teardown = setupReactTestEnv(); +}); +afterEach(() => teardown()); + +function render(element: ReactElement): string { + let renderer!: TestRenderer.ReactTestRenderer; + act(() => { + renderer = TestRenderer.create( + {element}, + ); + }); + const json = JSON.stringify(renderer.toJSON()); + act(() => renderer.unmount()); + return json; +} + +function runStages(fixture: PetriFixture): RunStage[] { + return stagesFromProjection(fixture.projection).map((stage) => ({ + id: stage.id, + name: stage.name, + handler: stage.handler, + status: stage.status, + node_id: stage.nodeId, + visit: stage.visit, + started_at: stage.startedAt, + wall_time_ms: fixture.projection.stages[stage.id]?.timing?.wall_time_ms, + usage: stage.usage, + parallel_group_id: stage.parallelGroupId ?? undefined, + parallel_branch_index: stage.parallelBranchIndex ?? undefined, + })); +} + +function createdAt(fixture: PetriFixture): string { + return new Date(fixture.stream[0].recorded_at).toISOString(); +} + +describe("a command-only run", () => { + const fixture = loadPetriFixture("command"); + const stages = stagesFromProjection(fixture.projection); + + test("the stage list shows every stage with its state", () => { + expect(stages.map((stage) => [stage.id, stage.status])).toEqual([ + ["start@1", "succeeded"], + ["say@1", "succeeded"], + ["exit@1", "succeeded"], + ]); + const html = render(); + for (const name of ["start", "say", "exit"]) expect(html).toContain(name); + // The sidebar shows a stage's state as its icon's tone: mint is succeeded. + expect((html.match(/text-mint/g) ?? []).length).toBe(3); + expect(html).not.toContain("animate-pulse"); + }); + + test("the command stage is one command turn with its exit status and output size", () => { + const say = fixture.projection.stages["say@1"]; + const activity = buildPetriStageActivity( + itemsForStage(fixture.stream, "say@1"), + say, + "command", + ); + expect(activity.turns).toHaveLength(1); + expect(activity.turns[0]).toMatchObject({ + kind: "command", + running: false, + exitCode: 0, + outputBytes: say.output_bytes, + }); + }); + + test("the waterfall's phases come from the lifecycle records", () => { + const html = render( + , + ); + for (const label of ["Submitted", "Runnable", "Initializing", "say"]) { + expect(html).toContain(label); + } + }); +}); + +describe("the hello run on the twin", () => { + const fixture = loadPetriFixture("hello"); + const stages = stagesFromProjection(fixture.projection); + + test("the agent stage's chat shows the prompt and the agent's response", () => { + const greet = stages.find((stage) => stage.id === "greet@1")!; + expect(greet.handler).toBe("agent"); + const activity = buildPetriStageActivity( + itemsForStage(fixture.stream, "greet@1"), + fixture.projection.stages["greet@1"], + "agent", + ); + expect(activity.turns.map((turn) => turn.kind)).toEqual(["system", "assistant"]); + expect(activity.turns[0]).toMatchObject({ kind: "system" }); + expect((activity.turns[0] as { content: string }).content).toContain("Add a haiku"); + expect(activity.turns[1]).toMatchObject({ + kind: "assistant", + content: "A haiku, added.", + inputTokens: 1, + outputTokens: 5, + }); + const html = render( + , + ); + expect(html).toContain("A haiku, added."); + }); + + test("the stage list shows the agent stage succeeded", () => { + const greet = stages.find((stage) => stage.id === "greet@1")!; + expect(greet.status).toBe("succeeded"); + expect(greet.providerUsed?.model).toBe("gpt-5.4"); + const html = render(); + expect(html).toContain("greet"); + expect(html).toContain("text-mint"); + }); +}); + +describe("a two-branch parallel run", () => { + const fixture = loadPetriFixture("parallel"); + const stages = stagesFromProjection(fixture.projection); + + test("the fork lists both branches under it with their outcomes", () => { + const fork = stages.find((stage) => stage.id === "fork@1")!; + const html = render( + , + ); + expect(html).toContain("Branches"); + expect(html).toContain("/runs/run-1/stages/a@1"); + expect(html).toContain("/runs/run-1/stages/b@1"); + expect((html.match(/Succeeded/g) ?? []).length).toBeGreaterThanOrEqual(2); + }); + + test("the fan-in joined the branches", () => { + const merge = stages.find((stage) => stage.id === "merge@1")!; + const html = render( + , + ); + expect(html).toContain("Joined"); + expect(html).not.toContain("Reducer transcript"); + }); + + test("the events view lists Petri events by name and the platform notice", () => { + const html = render( + {}} + runStart={createdAt(fixture)} + view="events" + onChangeView={() => {}} + />, + ); + for (const name of ["run.started", "visit.started", "fork.completed", "run.finished"]) { + expect(html).toContain(name); + } + expect(html).toContain("run.notice"); + expect(html).toContain('"data-stage":"a@1"'); + expect(html).toContain(`${fixture.stream.length} items`); + }); + + test("the overview lists the platform records", () => { + const html = render( + , + ); + expect(html).toContain("Platform records"); + expect(html).toContain("Notice"); + expect(html).toContain("recorded while both branches ran"); + }); + + test("the sidebar groups the branches under the fork", () => { + const branches = stages.filter((stage) => stage.parallelGroupId === "fork@1"); + expect(branches.map((stage) => stage.id)).toEqual(["a@1", "b@1"]); + const html = render(); + for (const name of ["fork", "a", "b", "merge"]) expect(html).toContain(name); + }); +}); + +describe("a human gate answered through the API", () => { + const fixture = loadPetriFixture("gate"); + const stages = stagesFromProjection(fixture.projection); + + test("the Q&A shows the question, its options and the answer with who gave it", () => { + const gate = stages.find((stage) => stage.id === "gate@1")!; + expect(gate.handler).toBe("human"); + const html = render( + , + ); + expect(html).toContain("Go?"); + expect(html).toContain("[Y] Yes"); + expect(html).toContain("[N] No"); + expect(html).toContain("dev"); + expect(html).not.toContain("pending"); + }); + + test("the gate's decision took the no edge", () => { + const gate = stages.find((stage) => stage.id === "gate@1")!; + const edge = findPetriEdgeForStage(fixture.stream, "gate@1"); + const html = render( + , + ); + expect(html).toContain("/runs/run-1/stages/no@1"); + expect(html).not.toContain("No target"); + }); + + test("no question is left pending in the projection", () => { + expect(Object.keys(fixture.projection.pending_interviews)).toEqual([]); + expect(stages.map((stage) => stage.id)).toEqual(["start@1", "gate@1", "no@1", "exit@1"]); + }); +}); diff --git a/apps/fabro-web/app/routes/run-stages.tsx b/apps/fabro-web/app/routes/run-stages.tsx index 9994760ef..06c802053 100644 --- a/apps/fabro-web/app/routes/run-stages.tsx +++ b/apps/fabro-web/app/routes/run-stages.tsx @@ -580,13 +580,21 @@ export function buildPetriStageActivity( return { turns, pendingTools: [] }; } - if (stage?.prompt) { + const envelopes = agentEnvelopesOf(items); + // The prompt is the session's `UserInput`; the projection's `prompt` stands + // in for a stage whose session recorded none (a prompt node). + if (stage?.prompt && !envelopes.some((envelope) => envelope.variant === "UserInput")) { turns.push({ kind: "system", ts: startTs, content: stage.prompt }); } let sawAssistantMessage = false; - for (const envelope of agentEnvelopesOf(items)) { + for (const envelope of envelopes) { const { payload } = envelope; switch (envelope.variant) { + case "UserInput": { + const text = getString(payload, "text") ?? ""; + if (text) turns.push({ kind: "system", ts: envelope.ts, content: text }); + break; + } case "AssistantMessage": { sawAssistantMessage = true; const tokens = getObject(getObject(payload, "usage"), "tokens") ?? getObject(payload, "usage") ?? {};