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