diff --git a/packages/chat-ui/src/chat-workspace.tsx b/packages/chat-ui/src/chat-workspace.tsx index 7e6a1f05d..a474de8de 100644 --- a/packages/chat-ui/src/chat-workspace.tsx +++ b/packages/chat-ui/src/chat-workspace.tsx @@ -253,8 +253,16 @@ export function mergeStreamingReply( ): readonly TimelineMessageItem[] { // A pending reply with no tokens yet stays off the timeline — an // empty bubble with no timestamp reads as broken; the typing pulse - // in the incoming-message slot owns that phase until the first delta lands. - if (streamingReply === null || streamingReply.text === "") return items; + // in the incoming-message slot owns that phase until the first delta + // lands. A `"replied"` turn renders nothing: its reply is already a + // persisted message. + if ( + streamingReply === null || + streamingReply.phase === "replied" || + streamingReply.text === "" + ) { + return items; + } const agent = participants.find((participant) => isAgentAddress(participant.address), ); diff --git a/packages/chat-ui/src/streaming-reply.test.ts b/packages/chat-ui/src/streaming-reply.test.ts index b58e52e91..4ad61b539 100644 --- a/packages/chat-ui/src/streaming-reply.test.ts +++ b/packages/chat-ui/src/streaming-reply.test.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test"; import { hydrateStreamingReplyFromTurn, + isPendingReply, nextStreamingReplyState, openPendingReply, typingAgentNames, @@ -22,12 +23,18 @@ function delta(text: string) { }); } +const AWAITING_EMPTY = { phase: "awaiting", text: "" } as const; + +function awaiting(text: string) { + return { phase: "awaiting", text } as const; +} + describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing reply)", () => { test("a non chat.agent event never opens or changes the reply", () => { expect( nextStreamingReplyState(null, { eventType: "chat.typing", data: {} }), ).toBeNull(); - const current = { text: "hi" }; + const current = awaiting("hi"); expect( nextStreamingReplyState(current, { eventType: "chat.pin", @@ -42,11 +49,11 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep null, agentEvent({ type: "inference.start", seq: 0, data: { model: "x" } }), ), - ).toEqual({ text: "" }); + ).toEqual(AWAITING_EMPTY); }); test("inference.start never wipes tokens already streamed", () => { - const state = { text: "Hello" }; + const state = awaiting("Hello"); expect( nextStreamingReplyState( state, @@ -61,14 +68,14 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep agentEvent({ type: "inference.start", seq: 0, data: { model: "x" } }), ); state = nextStreamingReplyState(state, delta("Hel")); - expect(state).toEqual({ text: "Hel" }); + expect(state).toEqual(awaiting("Hel")); state = nextStreamingReplyState(state, delta("Hello")); - expect(state).toEqual({ text: "Hello" }); + expect(state).toEqual(awaiting("Hello")); }); test("inference.done clears the reply — the persisted message takes over", () => { const state = nextStreamingReplyState( - { text: "Hello there" }, + awaiting("Hello there"), agentEvent({ type: "inference.done", seq: 5, @@ -79,7 +86,7 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep }); test("inference.done keeps an empty pending reply — the next inference round is still owed", () => { - const pending = { text: "" }; + const pending = awaiting(""); expect( nextStreamingReplyState( pending, @@ -107,7 +114,7 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep test("inference.error clears the reply rather than leaving a stuck cursor", () => { const state = nextStreamingReplyState( - { text: "Hello" }, + awaiting("Hello"), agentEvent({ type: "inference.error", seq: 5, @@ -118,7 +125,7 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep }); test("an event with no known inner shape (tool calls, usage) leaves the reply untouched", () => { - const state = { text: "Hello" }; + const state = awaiting("Hello"); expect( nextStreamingReplyState( state, @@ -137,11 +144,11 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep null, agentEvent({ type: "reactor.start", seq: 0, data: {} }), ), - ).toEqual({ text: "" }); + ).toEqual(AWAITING_EMPTY); }); test("reactor.start never resets an in-progress reply", () => { - const state = { text: "Hello" }; + const state = awaiting("Hello"); expect( nextStreamingReplyState( state, @@ -154,7 +161,7 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep for (const type of ["reactor.done", "reactor.error"]) { expect( nextStreamingReplyState( - { text: "Hello" }, + awaiting("Hello"), agentEvent({ type, seq: 10, data: {} }), ), ).toBeNull(); @@ -162,7 +169,7 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep }); test("a malformed delta payload (no partial.text) is ignored rather than crashing", () => { - const state = { text: "Hello" }; + const state = awaiting("Hello"); expect( nextStreamingReplyState( state, @@ -176,7 +183,7 @@ describe("nextStreamingReplyState (CL-6115: token deltas fold into a growing rep describe("nextStreamingReplyState (CL-6376: the typing pulse clears on a dispatch failure too)", () => { test("a chat.message carrying a turnFailed part clears a pending reply", () => { - const state = { text: "" }; + const state = awaiting(""); expect( nextStreamingReplyState(state, { eventType: "chat.message", @@ -191,7 +198,7 @@ describe("nextStreamingReplyState (CL-6376: the typing pulse clears on a dispatc }); test("an ordinary chat.message (no turnFailed part) leaves the reply untouched", () => { - const state = { text: "" }; + const state = awaiting(""); expect( nextStreamingReplyState(state, { eventType: "chat.message", @@ -210,76 +217,89 @@ describe("nextStreamingReplyState (CL-6376: the typing pulse clears on a dispatc }); }); -describe("nextStreamingReplyState (CL-6432: folded-run turns end on connector.reply / message.run.ended, not reactor.done)", () => { - test("connector.reply clears the reply — the orchestrator posts the persisted message off this very event", () => { +describe("nextStreamingReplyState (CL-6432 reopened: a folded run parks after the reply — post-reply tool rounds never re-open the pulse)", () => { + test("connector.reply moves the turn to the replied phase — the persisted message takes over the timeline", () => { expect( nextStreamingReplyState( - { text: "Hey! What are you working on today?" }, + awaiting("Hey! What are you working on right now?"), agentEvent({ type: "connector.reply", seq: 9, - data: { content: "Hey! What are you working on today?" }, + data: { content: "Hey! What are you working on right now?" }, }), ), - ).toBeNull(); + ).toEqual({ phase: "replied" }); }); - test("connector.reply clears an empty pending reply too — the reply is finalized even when its tokens never streamed here", () => { + test("connector.reply settles an empty pending reply too — the reply is finalized even when its tokens never streamed here", () => { expect( nextStreamingReplyState( - { text: "" }, + awaiting(""), agentEvent({ type: "connector.reply", seq: 9, data: { content: "Hey!" }, }), ), - ).toBeNull(); + ).toEqual({ phase: "replied" }); }); - test("message.run.ended clears an empty pending reply — the turn bracket closed, nothing more streams", () => { - expect( - nextStreamingReplyState( - { text: "" }, - agentEvent({ - type: "message.run.ended", - seq: 12, - data: { status: "completed" }, - }), - ), - ).toBeNull(); + test("the replied phase renders nothing — no pulse, no bubble", () => { + const replied = nextStreamingReplyState( + awaiting(""), + agentEvent({ type: "connector.reply", seq: 9, data: { content: "x" } }), + ); + expect(isPendingReply(replied)).toBe(false); + expect(typingAgentNames(replied, [HUMAN, MYRA])).toEqual([]); }); - test("connector.reply and message.run.ended while idle stay idle", () => { - expect( - nextStreamingReplyState( - null, - agentEvent({ type: "connector.reply", seq: 9, data: { content: "x" } }), - ), - ).toBeNull(); + test("message.run.started opens the next turn's pulse — the one event that ends the replied phase from the stream", () => { + const replied = nextStreamingReplyState( + awaiting(""), + agentEvent({ type: "connector.reply", seq: 9, data: { content: "x" } }), + ); expect( nextStreamingReplyState( - null, - agentEvent({ - type: "message.run.ended", - seq: 12, - data: { status: "completed" }, - }), + replied, + agentEvent({ type: "message.run.started", seq: 0, data: {} }), ), - ).toBeNull(); + ).toEqual(AWAITING_EMPTY); + }); + + test("message.run.ended returns to idle from any phase — the bracket closed, nothing more streams", () => { + for (const state of [awaiting(""), awaiting("Hello")]) { + expect( + nextStreamingReplyState( + state, + agentEvent({ + type: "message.run.ended", + seq: 12, + data: { status: "completed" }, + }), + ), + ).toBeNull(); + } }); - test("the Myra repro: a post-reply tool round reopens the pulse and the turn's own terminal events shut it", () => { - // Round 1: the visible reply streams and its inference.done hands off - // to the persisted message. + test("the live Myra sequence: reply posts, memory rounds follow, the run parks — the pulse never comes back", () => { + // Captured from a live folded run (scratch stack, real provider): the + // run brackets open per dequeued message, the visible reply streams + // and posts via connector.reply, then post-reply tool-only rounds + // (memory writes) run inference again, and the run PARKS — no + // message.run.ended ever arrives. let state = openPendingReply(null); + state = nextStreamingReplyState( + state, + agentEvent({ type: "message.run.started", seq: 0, data: {} }), + ); state = nextStreamingReplyState( state, agentEvent({ type: "inference.start", seq: 1, data: { model: "x" } }), ); + expect(isPendingReply(state)).toBe(true); state = nextStreamingReplyState( state, - delta("Hey! What are you working on today?"), + delta("Hey! What are you working on right now?"), ); state = nextStreamingReplyState( state, @@ -289,54 +309,82 @@ describe("nextStreamingReplyState (CL-6432: folded-run turns end on connector.re data: { turn: {}, usage: {}, source: "primary" }, }), ); - expect(state).toBeNull(); - - // Round 2: a tool-only follow-up (memory writes) reopens the pulse and - // its textless inference.done deliberately leaves it up mid-turn. - state = nextStreamingReplyState( - state, - agentEvent({ type: "inference.start", seq: 5, data: { model: "x" } }), - ); - expect(state).toEqual({ text: "" }); state = nextStreamingReplyState( state, agentEvent({ - type: "inference.done", - seq: 7, - data: { turn: {}, usage: {}, source: "primary" }, + type: "connector.reply", + seq: 5, + data: { content: "Hey! What are you working on right now?" }, }), ); - expect(state).toEqual({ text: "" }); + expect(state).toEqual({ phase: "replied" }); - // The folded run never emits reactor.done — its turn ends here. - state = nextStreamingReplyState( - state, + // Post-reply memory rounds: inference.start must NOT re-open the + // pulse, and the textless inference.done must not strand one either. + for (const round of [ + agentEvent({ type: "inference.start", seq: 6, data: { model: "x" } }), agentEvent({ - type: "connector.reply", + type: "inference.tool_call.start", + seq: 7, + data: { callId: "c1", name: "memory_write" }, + }), + agentEvent({ + type: "tool.done", seq: 8, - data: { content: "Hey!" }, + data: { result: { callId: "c1", content: [], isError: false } }, }), - ); - expect(state).toBeNull(); - state = nextStreamingReplyState( - state, agentEvent({ - type: "message.run.ended", + type: "inference.done", seq: 9, - data: { status: "completed" }, + data: { turn: {}, usage: {}, source: "primary" }, }), + agentEvent({ type: "inference.start", seq: 10, data: { model: "x" } }), + agentEvent({ + type: "inference.done", + seq: 11, + data: { turn: {}, usage: {}, source: "primary" }, + }), + ]) { + state = nextStreamingReplyState(state, round); + expect(state).toEqual({ phase: "replied" }); + expect(isPendingReply(state)).toBe(false); + } + // The run parks here: no message.run.ended, and the state stays + // invisible until the next turn's message.run.started. + }); + + test("the next user turn brings the pulse back — replied never suppresses a genuinely new turn", () => { + const replied = nextStreamingReplyState( + awaiting(""), + agentEvent({ type: "connector.reply", seq: 5, data: { content: "x" } }), ); - expect(state).toBeNull(); + // Locally, the send itself re-opens the pulse... + expect(openPendingReply(replied)).toEqual(AWAITING_EMPTY); + // ...and on the stream, the dequeued message's own bracket does. + let state = nextStreamingReplyState( + replied, + agentEvent({ type: "message.run.started", seq: 0, data: {} }), + ); + expect(isPendingReply(state)).toBe(true); + state = nextStreamingReplyState( + state, + agentEvent({ type: "inference.start", seq: 1, data: { model: "x" } }), + ); + expect(isPendingReply(state)).toBe(true); }); }); describe("openPendingReply", () => { test("opens an empty pending reply when idle", () => { - expect(openPendingReply(null)).toEqual({ text: "" }); + expect(openPendingReply(null)).toEqual(AWAITING_EMPTY); + }); + + test("opens an empty pending reply from a replied previous turn", () => { + expect(openPendingReply({ phase: "replied" })).toEqual(AWAITING_EMPTY); }); test("never resets a reply already streaming", () => { - const state = { text: "Hel" }; + const state = awaiting("Hel"); expect(openPendingReply(state)).toBe(state); }); }); @@ -347,23 +395,23 @@ describe("typingAgentNames", () => { }); test("a pending reply with no tokens names the workbench's agent participant", () => { - expect(typingAgentNames({ text: "" }, [HUMAN, MYRA])).toEqual(["Myra"]); + expect(typingAgentNames(awaiting(""), [HUMAN, MYRA])).toEqual(["Myra"]); }); test("once tokens stream the bubble takes over — the typing line goes quiet", () => { - expect(typingAgentNames({ text: "Hel" }, [HUMAN, MYRA])).toEqual([]); + expect(typingAgentNames(awaiting("Hel"), [HUMAN, MYRA])).toEqual([]); }); test('a slugified handle is shown as a display name — "myra" reads "Myra"', () => { expect( - typingAgentNames({ text: "" }, [ + typingAgentNames(awaiting(""), [ { address: "myra@agents.example", handle: "myra" }, ]), ).toEqual(["Myra"]); }); test("no agent participant on the workbench means nobody is named", () => { - expect(typingAgentNames({ text: "" }, [HUMAN])).toEqual([]); + expect(typingAgentNames(awaiting(""), [HUMAN])).toEqual([]); }); }); @@ -375,12 +423,12 @@ describe("hydrateStreamingReplyFromTurn (CL-6380: reattach snapshot)", () => { test("a running turn with committed text opens the reply carrying it", () => { expect( hydrateStreamingReplyFromTurn({ textSnapshot: "streamed so far" }), - ).toEqual({ text: "streamed so far" }); + ).toEqual(awaiting("streamed so far")); }); test("a running turn with no text yet opens the same empty pending pulse as openPendingReply", () => { - expect(hydrateStreamingReplyFromTurn({ textSnapshot: null })).toEqual({ - text: "", - }); + expect(hydrateStreamingReplyFromTurn({ textSnapshot: null })).toEqual( + AWAITING_EMPTY, + ); }); }); diff --git a/packages/chat-ui/src/streaming-reply.ts b/packages/chat-ui/src/streaming-reply.ts index f17166d82..730c6e958 100644 --- a/packages/chat-ui/src/streaming-reply.ts +++ b/packages/chat-ui/src/streaming-reply.ts @@ -17,7 +17,36 @@ import { isAgentAddress } from "@corbits/chat/mentions"; import type { ParticipantRecord } from "./api"; import { displayNameFromHandle } from "./timeline"; -export type StreamingReplyState = { readonly text: string } | null; +/** + * The current turn's reply state, phase-tagged (CL-6432 reopened): + * + * - `"awaiting"` — a turn is in flight and its visible reply hasn't posted + * yet: an empty `text` renders the typing pulse, streamed tokens render + * the growing bubble. + * - `"replied"` — this turn's `connector.reply` has posted. A live folded + * run PARKS after the turn (no `message.run.ended`), and its post-reply + * tool-only rounds (memory writes) still emit `inference.start`/ + * `inference.done` — none of which may re-open the pulse. Renders + * nothing; only a genuinely new turn leaves it. + * - `null` — idle, no turn in flight. + */ +export type StreamingReplyState = + | { readonly phase: "awaiting"; readonly text: string } + | { readonly phase: "replied" } + | null; + +const REPLIED: StreamingReplyState = { phase: "replied" }; + +function awaiting(text: string): StreamingReplyState { + return { phase: "awaiting", text }; +} + +/** Whether the state renders the tokenless typing pulse. The `"replied"` + * phase is deliberately not pending — it renders nothing and must never + * arm the pending-timeout backstop. */ +export function isPendingReply(state: StreamingReplyState): boolean { + return state !== null && state.phase === "awaiting" && state.text === ""; +} type InferenceDeltaEvent = { readonly type: "inference.text.delta"; @@ -70,18 +99,24 @@ function hasTurnFailedPart(data: unknown): boolean { } /** - * The streaming reply's whole state machine, pure: an `inference.start` - * opens an empty in-progress reply if nothing is showing yet (it never - * wipes tokens already streamed), each `inference.text.delta` replaces it - * with that delta's cumulative text, and the turn-terminal events clear - * it — `reactor.done`/`reactor.error` for reactor-style runs, - * `connector.reply`/`message.run.ended` for folded runs, which never emit - * a per-turn `reactor.done` (CL-6432). `inference.done` only clears once tokens - * have streamed (the persisted message takes over); an empty pending - * survives so the typing pulse stays up across tool rounds. `inference.error` - * always clears. A `chat.message` carrying `postUndeliveredNotice`'s - * `turnFailed` part also clears it (see `hasTurnFailedPart`) — the one - * failure path with no `chat.agent` events of its own. Every other event + * The streaming reply's whole state machine, pure and turn-phase aware + * (CL-6432 reopened). `message.run.started` — the harness's per-dequeued- + * message turn begin, the same event the chat orchestrator keys new turns + * off (see `chat-orchestrator.ts`'s use of `messageRunStarted`) — opens a + * fresh awaiting turn. While awaiting, `inference.start`/`reactor.start` + * open the empty pulse (never wiping streamed tokens), each + * `inference.text.delta` replaces the text with its cumulative snapshot, + * and a textless `inference.done` keeps the pulse up across pre-reply tool + * rounds. `connector.reply` — the event the orchestrator posts the + * persisted reply off — moves the turn to `"replied"`: a live folded run + * PARKS here (no `message.run.ended`), and its post-reply tool-only rounds + * (memory writes) still emit `inference.start`/`inference.done`, so in + * `"replied"` every inference/reactor event is inert rather than + * re-opening the pulse. The hard-terminal events — + * `reactor.done`/`reactor.error`, `message.run.ended`, `inference.error`, + * and a `chat.message` carrying `postUndeliveredNotice`'s `turnFailed` + * part (see `hasTurnFailedPart`, the one failure path with no `chat.agent` + * events of its own) — return to idle from any phase. Every other event * type (tool calls, thinking, usage) leaves the current state untouched. */ export function nextStreamingReplyState( @@ -94,27 +129,24 @@ export function nextStreamingReplyState( if (event.eventType !== "chat.agent") return current; const innerType = innerEventType(event.data); - if (innerType === "inference.start") return current ?? { text: "" }; + if (innerType === "message.run.started") return awaiting(""); + if ( + innerType === "reactor.done" || + innerType === "reactor.error" || + innerType === "message.run.ended" || + innerType === "inference.error" + ) { + return null; + } + if (innerType === "connector.reply") return REPLIED; + if (current !== null && current.phase === "replied") return current; + + if (innerType === "inference.start") return current ?? awaiting(""); // `reactor.start` is the earliest "the agent is on it" signal — it // fires before any tokens, often seconds before a slow model's // `inference.start` — so it opens the indicator without waiting for // the first inference call. - if (innerType === "reactor.start") return current ?? { text: "" }; - if (innerType === "reactor.done" || innerType === "reactor.error") { - return null; - } - // A folded run (@corbits/folded-runs — every workbench agent) never - // emits a per-turn `reactor.done`: its turn is finalized by - // `connector.reply` (the very event the chat orchestrator posts the - // persisted reply off, see chat-orchestrator.ts) and bracket-closed by - // `message.run.ended` either way. Without these two clears, a tool-only - // follow-up round after the visible reply (memory writes) reopens the - // empty pending pulse via its `inference.start` and nothing ever shuts - // it (CL-6432). - if (innerType === "connector.reply" || innerType === "message.run.ended") { - return null; - } - if (innerType === "inference.error") return null; + if (innerType === "reactor.start") return current ?? awaiting(""); if (innerType === "inference.done") { if (current === null || current.text === "") return current; return null; @@ -122,7 +154,7 @@ export function nextStreamingReplyState( const delta = parseInferenceDeltaEvent(event.data); if (delta === null) return current; - return { text: delta.data.partial.text }; + return awaiting(delta.data.partial.text); } /** @@ -143,7 +175,10 @@ export function nextStreamingReplyState( export function openPendingReply( current: StreamingReplyState, ): StreamingReplyState { - return current ?? { text: "" }; + // A `"replied"` previous turn is over — the send that called this opens + // the next one, so the pulse comes back. + if (current === null || current.phase === "replied") return awaiting(""); + return current; } /** @@ -160,7 +195,7 @@ export function hydrateStreamingReplyFromTurn( runningTurn: { readonly textSnapshot?: string | null } | null, ): StreamingReplyState { if (runningTurn === null) return null; - return { text: runningTurn.textSnapshot ?? "" }; + return awaiting(runningTurn.textSnapshot ?? ""); } /** How long an empty pending reply may sit with no tokens before the @@ -209,7 +244,7 @@ export function useStreamingReply( }, [workbenchId]); useEffect(() => { - if (streamingReply === null || streamingReply.text !== "") return; + if (!isPendingReply(streamingReply)) return; const timer = setTimeout(() => { // The dependency below re-arms this effect (clearing this exact // timer) the instant `streamingReply` changes, so this callback only @@ -227,16 +262,10 @@ export function useStreamingReply( current: StreamingReplyState, ): StreamingReplyState { const now = Date.now(); - const becamePending = - next !== null && - next.text === "" && - (current === null || current.text !== ""); + const becamePending = isPendingReply(next) && !isPendingReply(current); if (becamePending) pendingSinceRef.current = now; - const leavingPending = - current !== null && - current.text === "" && - (next === null || next.text !== ""); + const leavingPending = isPendingReply(current) && !isPendingReply(next); if ( leavingPending && pendingSinceRef.current !== null && @@ -257,7 +286,7 @@ export function useStreamingReply( clearTimeout(holdTimerRef.current); holdTimerRef.current = null; } - if (next === null || next.text !== "") pendingSinceRef.current = null; + if (!isPendingReply(next)) pendingSinceRef.current = null; return next; } @@ -321,10 +350,10 @@ export function typingAgentNames( streamingReply: StreamingReplyState, participants: readonly ParticipantRecord[], ): readonly string[] { - // Only the tokenless pending phase shows the typing line; once text - // streams, the growing timeline bubble is the signal — showing both - // would double-indicate the same turn. - if (streamingReply === null || streamingReply.text !== "") return []; + // Only the tokenless awaiting phase shows the typing line; once text + // streams, the growing timeline bubble is the signal, and a `"replied"` + // turn shows nothing at all. + if (!isPendingReply(streamingReply)) return []; const agent = participants.find((participant) => isAgentAddress(participant.address), ); diff --git a/packages/chat-ui/test/refresh-and-send.test.ts b/packages/chat-ui/test/refresh-and-send.test.ts index 015cbc3d9..74d8e9c3d 100644 --- a/packages/chat-ui/test/refresh-and-send.test.ts +++ b/packages/chat-ui/test/refresh-and-send.test.ts @@ -575,9 +575,11 @@ describe("mergeStreamingReply (CL-6115: the in-progress agent reply folds into t }); test("a growing reply appends a streaming item attributed to the workbench's agent", () => { - const merged = mergeStreamingReply(serverItems, { text: "Working on it" }, [ - agent, - ]); + const merged = mergeStreamingReply( + serverItems, + { phase: "awaiting", text: "Working on it" }, + [agent], + ); expect(merged).toHaveLength(2); expect(merged[1]).toMatchObject({ streaming: true, @@ -587,12 +589,20 @@ describe("mergeStreamingReply (CL-6115: the in-progress agent reply folds into t }); test("no agent participant to attribute the reply to means no synthetic item", () => { - const merged = mergeStreamingReply(serverItems, { text: "hi" }, []); + const merged = mergeStreamingReply( + serverItems, + { phase: "awaiting", text: "hi" }, + [], + ); expect(merged).toBe(serverItems); }); test("a pending reply with no tokens yet renders no ghost bubble — the typing line owns that phase", () => { - const merged = mergeStreamingReply(serverItems, { text: "" }, [agent]); + const merged = mergeStreamingReply( + serverItems, + { phase: "awaiting", text: "" }, + [agent], + ); expect(merged).toBe(serverItems); }); }); diff --git a/packages/chat-ui/test/use-streaming-reply.test.tsx b/packages/chat-ui/test/use-streaming-reply.test.tsx index aed893a01..7e4a90878 100644 --- a/packages/chat-ui/test/use-streaming-reply.test.tsx +++ b/packages/chat-ui/test/use-streaming-reply.test.tsx @@ -94,20 +94,20 @@ describe("useStreamingReply (CL-6115: live wiring)", () => { seq: 0, data: { model: "x" }, }); - expect(harness.get()).toEqual({ text: "" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "" }); harness.send("chat.agent", delta("Hel")); - expect(harness.get()).toEqual({ text: "Hel" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hel" }); harness.send("chat.agent", delta("Hello")); - expect(harness.get()).toEqual({ text: "Hello" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hello" }); harness.unmount(); }); test("inference.done clears the reply once the turn ends", () => { const harness = mount("chan_a"); harness.send("chat.agent", delta("Hello")); - expect(harness.get()).toEqual({ text: "Hello" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hello" }); harness.send("chat.agent", { type: "inference.done", @@ -121,7 +121,7 @@ describe("useStreamingReply (CL-6115: live wiring)", () => { test("switching workbenches clears whatever was streaming in the one just left", () => { const harness = mount("chan_a"); harness.send("chat.agent", delta("Hello")); - expect(harness.get()).toEqual({ text: "Hello" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hello" }); harness.switchWorkbench("chan_b"); expect(harness.get()).toBeNull(); @@ -131,13 +131,13 @@ describe("useStreamingReply (CL-6115: live wiring)", () => { test("noteAwaitingReply opens a pending reply after a send, and deltas take over", () => { const harness = mount("chan_a"); harness.awaitReply(); - expect(harness.get()).toEqual({ text: "" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "" }); harness.send("chat.agent", delta("Hel")); - expect(harness.get()).toEqual({ text: "Hel" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hel" }); harness.awaitReply(); - expect(harness.get()).toEqual({ text: "Hel" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hel" }); harness.unmount(); }); @@ -145,7 +145,7 @@ describe("useStreamingReply (CL-6115: live wiring)", () => { const harness = mount("chan_a"); harness.send("chat.agent", delta("Hello")); harness.send("chat.typing", { principalId: "prn_other" }); - expect(harness.get()).toEqual({ text: "Hello" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hello" }); harness.unmount(); }); }); @@ -154,7 +154,7 @@ describe("useStreamingReply's reply-timeout backstop (CL-6252 #6)", () => { test("a pending reply with no tokens for the whole clearMs marks replyTimedOut", async () => { const harness = mount("chan_a", 30); harness.awaitReply(); - expect(harness.get()).toEqual({ text: "" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "" }); expect(harness.timedOut()).toBe(false); await harness.settle(60); @@ -169,7 +169,7 @@ describe("useStreamingReply's reply-timeout backstop (CL-6252 #6)", () => { harness.send("chat.agent", delta("Hi")); await harness.settle(60); - expect(harness.get()).toEqual({ text: "Hi" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hi" }); expect(harness.timedOut()).toBe(false); harness.unmount(); }); @@ -203,14 +203,17 @@ describe("useStreamingReply.resumeFromTurn (CL-6380: catch-up on remount)", () = expect(harness.get()).toBeNull(); harness.resumeFromTurn({ textSnapshot: "already streamed so far" }); - expect(harness.get()).toEqual({ text: "already streamed so far" }); + expect(harness.get()).toEqual({ + phase: "awaiting", + text: "already streamed so far", + }); harness.unmount(); }); test("a running turn with no text yet opens the same empty pending pulse as noteAwaitingReply", () => { const harness = mount("chan_a"); harness.resumeFromTurn({ textSnapshot: null }); - expect(harness.get()).toEqual({ text: "" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "" }); harness.unmount(); }); @@ -222,20 +225,20 @@ describe("useStreamingReply.resumeFromTurn (CL-6380: catch-up on remount)", () = data: { model: "x" }, }); harness.send("chat.agent", delta("hi")); - expect(harness.get()).toEqual({ text: "hi" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "hi" }); harness.resumeFromTurn(null); - expect(harness.get()).toEqual({ text: "hi" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "hi" }); harness.unmount(); }); test("a live event that already opened the reply wins over a slower snapshot fetch", () => { const harness = mount("chan_a"); harness.send("chat.agent", delta("live wins")); - expect(harness.get()).toEqual({ text: "live wins" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "live wins" }); harness.resumeFromTurn({ textSnapshot: "stale snapshot" }); - expect(harness.get()).toEqual({ text: "live wins" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "live wins" }); harness.unmount(); }); }); @@ -244,13 +247,13 @@ describe("useStreamingReply's typing-pulse floor", () => { test("a token arriving immediately still leaves the empty pulse up until minVisibleMs", async () => { const harness = mount("chan_a", undefined, 40); harness.awaitReply(); - expect(harness.get()).toEqual({ text: "" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "" }); harness.send("chat.agent", delta("Hi")); - expect(harness.get()).toEqual({ text: "" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "" }); await harness.settle(60); - expect(harness.get()).toEqual({ text: "Hi" }); + expect(harness.get()).toEqual({ phase: "awaiting", text: "Hi" }); harness.unmount(); }); });