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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 01:03:19 -04:00
parent f22cf18c86
commit 3941a24a28
No known key found for this signature in database
6 changed files with 659 additions and 13 deletions

View file

@ -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;
}

View file

@ -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([]);
});
});

View file

@ -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<typeof principalDisplay>[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 };

View file

@ -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();

View file

@ -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(
<MemoryRouter initialEntries={["/runs/run-1"]}>{element}</MemoryRouter>,
);
});
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(<StageSidebar stages={stages} runId="run-1" />);
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(
<RunWaterfall
runId="run-1"
events={[]}
phases={deriveRunPhasesFromStream(fixture.stream, createdAt(fixture))}
stages={runStages(fixture)}
createdAtIso={createdAt(fixture)}
completedAtIso={fixture.projection.conclusion?.timestamp ?? null}
/>,
);
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(
<StageChatView
turns={activity.turns}
pendingTools={activity.pendingTools}
stage={greet}
/>,
);
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(<StageSidebar stages={stages} runId="run-1" />);
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(
<ParallelChildren
stage={fork}
events={[]}
overview={parallelOverviewFromProjection(fixture.projection.stages["fork@1"])}
runId="run-1"
allStages={stages}
/>,
);
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(
<FanInResults
stage={merge}
events={[]}
reducer={reducerTranscriptFromProjection(fixture.projection.stages["merge@1"])}
/>,
);
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(
<StreamEventsView
rows={debugRowsFromStream(fixture.stream)}
error={undefined}
onRetry={() => {}}
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(
<PlatformRecordsPanelView
records={platformRecordsOf(fixture.stream)}
projection={fixture.projection}
/>,
);
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(<StageSidebar stages={stages} runId="run-1" />);
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(
<HumanQA
stage={gate}
events={[]}
pairs={parsePetriInterviewPairs(fixture.stream)}
/>,
);
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(
<ConditionalDecision
stage={gate}
runEvents={[]}
edge={edge}
allStages={stages}
runId="run-1"
/>,
);
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"]);
});
});

View file

@ -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") ?? {};