diff --git a/src/interfaces/tui/state.ts b/src/interfaces/tui/state.ts index a14600e..01a18f4 100644 --- a/src/interfaces/tui/state.ts +++ b/src/interfaces/tui/state.ts @@ -329,10 +329,13 @@ export interface ChatMessageDisplay { content: string; reasoning: string; timestamp: number; + /** Sliding-window replay cursor for content + reasoning (transient). */ + _posContent: number; + _posReasoning: number; } export function makeChatMessage(role: Role, content: string): ChatMessageDisplay { - return { role, content, reasoning: "", timestamp: Date.now() }; + return { role, content, reasoning: "", timestamp: Date.now(), _posContent: 0, _posReasoning: 0 }; } /** Bounded ring of recent chat messages used to render the transcript view. */ @@ -357,9 +360,9 @@ export class TranscriptCache { const msg = this.messages[this.messages.length - 1]; if (msg && msg.role === Roles.Assistant) { if (isReasoning) { - msg.reasoning = reconcileStream(msg.reasoning, text); + msg.reasoning = reconcileStream(msg.reasoning, text, msg, true); } else { - msg.content = reconcileStream(msg.content, text); + msg.content = reconcileStream(msg.content, text, msg, false); } this.dirty = true; } @@ -367,34 +370,72 @@ export class TranscriptCache { } /** - * Reconcile an incoming stream chunk against the accumulated text. + * Reconcile an incoming stream chunk against the accumulated text, deduping the + * sliding-window replays that 9router's OpenAI-compatible layer emits. * * Streaming backends disagree on chunk semantics: * - Native OpenAI/Anthropic send small *incremental* deltas. - * - Some gateways/proxies (e.g. 9router's OpenAI-compatible layer) emit - * *full-text* chunks, where each chunk is the complete transcript so far. + * - 9router/gateways re-send the FULL accumulated text as a repeated window: + * each "wave" replays the already-emitted segments from the start, then one + * fresh token. So the stream is s1,s2,…,sk then s1,s2,…,sk,s_{k+1} then + * s1,…,s_{k+1},s_{k+2} … where s_i + … = the final text. * - * Blindly concatenating full-text chunks duplicates content (partial + full), - * which manifests as garbled messages. This helper accepts both modes: - * - a chunk that is a strict prefix of the tail → a shorter replay → ignore; - * - a chunk that extends/replays the whole accumulated text → take it as the - * new authoritative tail (full-text mode); - * - otherwise → treat as an incremental delta and append. + * Naive concatenation doubles the text (garbled output). We track a per-stream + * replay cursor `pos` into `accumulated`: a chunk that continues at `pos` (or + * restarts at 0, i.e. begins a fresh wave) is a replay → ignore it and advance + * the cursor; only a chunk that matches nothing already present is brand-new + * content, which we append at the end and reset the cursor. + * + * @param accumulated current merged text + * @param chunk next chunk from the stream + * @param msg the message being streamed (holds the replay cursor) + * @param isReasoning whether this targets the reasoning buffer or the content + * @returns the updated accumulated text */ -export function reconcileStream(accumulated: string, chunk: string): string { +export function reconcileStream( + accumulated: string, + chunk: string, + msg: ChatMessageDisplay, + isReasoning: boolean, +): string { if (chunk.length === 0) return accumulated; - if (accumulated.length === 0) return chunk; - if (accumulated.length >= chunk.length && accumulated.startsWith(chunk)) { - // Shorter (or equal) replay of what we already hold — ignore. + const posRef = { pos: isReasoning ? msg._posReasoning : msg._posContent }; + if (accumulated.length === 0) { + accumulated = chunk; + posRef.pos = chunk.length; + writePos(); return accumulated; } - if (chunk.length >= accumulated.length && chunk.startsWith(accumulated)) { - // Full-text mode: the chunk is the whole transcript so far (matching our - // accumulated text plus more). Adopt it wholesale. - return chunk; + // 1) A chunk that continues the current replay cursor is a replayed segment + // (it already exists earlier in the text) → skip it, advance the cursor. + if ( + posRef.pos > 0 && + posRef.pos < accumulated.length && + chunk.length <= accumulated.length - posRef.pos && + accumulated.startsWith(chunk, posRef.pos) + ) { + posRef.pos += chunk.length; + writePos(); + return accumulated; + } + // 2) A chunk that matches at the very start begins a NEW replay wave → skip + // it, set the cursor to just past the matched prefix. + if (accumulated.startsWith(chunk)) { + posRef.pos = chunk.length; + writePos(); + return accumulated; + } + // 3) Otherwise this is brand-new content: append it and reset the cursor so + // the next wave starts fresh. + accumulated += chunk; + posRef.pos = -1; + writePos(); + return accumulated; + + function writePos(): void { + if (isReasoning) msg._posReasoning = posRef.pos; + else msg._posContent = posRef.pos; } - // Incremental delta mode: a new suffix. Append. - return accumulated + chunk; } /* ── AppStateRest ─────────────────────────────────────────────────── */ diff --git a/src/interfaces/tui/tui.test.ts b/src/interfaces/tui/tui.test.ts index 3a6af64..5f1acaf 100644 --- a/src/interfaces/tui/tui.test.ts +++ b/src/interfaces/tui/tui.test.ts @@ -4,7 +4,7 @@ */ import { describe, test, expect } from "bun:test"; import { parseCommand, applyCommand } from "./command.ts"; -import { createTuiState, pushTranscript, makeChatMessage, reconcileStream } from "./state.ts"; +import { createTuiState, pushTranscript, makeChatMessage, reconcileStream, type ChatMessageDisplay } from "./state.ts"; import { applyAction } from "./action.ts"; import { handleKey, decodeKey } from "./controller.ts"; import { toControllerKey } from "./ui.tsx"; @@ -163,32 +163,89 @@ describe("toControllerKey (OpenTUI adapter)", () => { }); }); -describe("reconcileStream (full-text vs incremental chunk handling)", () => { - test("appends incremental deltas", () => { - let s = reconcileStream("", "Hello"); - s = reconcileStream(s, "!"); - s = reconcileStream(s, " How"); +describe("reconcileStream (9router sliding-window vs incremental)", () => { + function msg(): ChatMessageDisplay { + return makeChatMessage(Roles.Assistant, ""); + } + function rec( + s: string, + chunk: string, + m: ChatMessageDisplay, + ): string { + return reconcileStream(s, chunk, m, false); + } + + test("incremental deltas append in order", () => { + const m = msg(); + let s = rec("", "Hello", m); + s = rec(s, "!", m); + s = rec(s, " How", m); expect(s).toBe("Hello! How"); }); - test("full-text mode: adopts a growing authoritative chunk", () => { - let s = reconcileStream("", "Hello"); - // 9router/proxy resends the *whole* transcript each chunk — adopt it. - s = reconcileStream(s, "Hello world"); - s = reconcileStream(s, "Hello world!"); - expect(s).toBe("Hello world!"); - }); - test("shorter/equal replay of the tail is ignored (no duplication)", () => { - const tail = "Hello world"; - expect(reconcileStream(tail, tail)).toBe(tail); - expect(reconcileStream(tail, "Hello")).toBe(tail); - }); - test("partial-then-full garbled stream resolves to the full text", () => { - let s = reconcileStream("", "Hello! 更新"); - s = reconcileStream(s, "Hello! 更新 Keep going"); - expect(s).toBe("Hello! 更新 Keep going"); + + test("real 9router wave: replayed segments never double", () => { + // Real captured stream (story turn), each fresh token restarts a replay + // wave from the head: "Here's a 3-sent" "ence story" " about" " a robo" … + const m = msg(); + const segs: string[] = [ + "Here's a 3-sent", + "ence story", + " about", + " a robo", + "t:\n\nA small repai", + "r robo", + "t name", + "d Bolt", + ]; + // Wave 1 (first two segments). + let s = ""; + s = rec(s, segs[0]!, m); + s = rec(s, segs[1]!, m); + expect(s).toBe("Here's a 3-sentence story"); + // Wave 2 replays seg0..seg1 then new seg2. + for (const g of segs.slice(0, 3)) s = rec(s, g, m); + expect(s).toBe("Here's a 3-sentence story about"); + // Wave 3 replays seg0..seg2 then new seg3. + for (const g of segs.slice(0, 4)) s = rec(s, g, m); + expect(s).toBe("Here's a 3-sentence story about a robo"); + // Wave 4 replays 0..3 then new seg4. + for (const g of segs.slice(0, 5)) s = rec(s, g, m); + expect(s).toBe("Here's a 3-sentence story about a robot:\n\nA small repai"); + // Wave 5 replays 0..4 then new seg5. + for (const g of segs.slice(0, 6)) s = rec(s, g, m); + expect(s).toBe( + "Here's a 3-sentence story about a robot:\n\nA small repair robo", + ); + // Wave 6 replays 0..5 then new seg6. + for (const g of segs.slice(0, 7)) s = rec(s, g, m); + // Final wave replays everything, no new content. + for (const g of segs.slice(0, 8)) s = rec(s, g, m); + expect(s).toBe( + "Here's a 3-sentence story about a robot:\n\nA small repair robot named Bolt", + ); }); + test("empty chunk is a no-op", () => { - expect(reconcileStream("abc", "")).toBe("abc"); + const m = msg(); + expect(rec("abc", "", m)).toBe("abc"); + }); + + test("1-char shared letter is NOT treated as a replay", () => { + // "Hi! " ends with a space; "How…" starts with 'H' — must append fully. + const m = msg(); + let s = rec("", "Hi! ", m); + s = rec(s, "How can I help you today?", m); + expect(s).toBe("Hi! How can I help you today?"); + }); + + test("'say hi' captured sequence resolves cleanly", () => { + const m = msg(); + let s = ""; + // Real captured: "", "", "Hi! 👋 ", "", "", "Hi! 👋 ", "How can I help you today?" + s = rec(s, "Hi! 👋 ", m); + s = rec(s, "Hi! 👋 ", m); // replay of head + s = rec(s, "How can I help you today?", m); + expect(s).toBe("Hi! 👋 How can I help you today?"); }); }); diff --git a/src/interfaces/tui/ui.tsx b/src/interfaces/tui/ui.tsx index c1c0629..17263e2 100644 --- a/src/interfaces/tui/ui.tsx +++ b/src/interfaces/tui/ui.tsx @@ -276,25 +276,37 @@ function TranscriptMessage(props: { msg.role === Roles.Assistant && !msg.content.trim() ? "" : msg.content || "(empty)"; + const bodyWidth = Math.max(20, width - 6); - // Each message is a labelled block: a coloured role tag line followed by the - // word-wrapped body. Plain `` is used (not ``) so live - // streaming re-renders a single clean string — no overlap artifacts, and the - // body fills the full transcript width. return ( - {` ${tag.padEnd(4)} `} + {/* Role tag header line */} + {` [${tag}]`} + {/* Indented, word-wrapped body */} {content.length > 0 && ( - {wrapText(content, Math.max(20, width - 4))} + + {indent(content, 4, bodyWidth)} + )} + {/* Live reasoning (dimmed, kept to one line) */} {msg.reasoning.length > 0 && ( - {` ⋯ ${oneLine(msg.reasoning)}`} + {` ⋯ ${oneLine(msg.reasoning)}`} )} + {/* Thin separator between messages */} + {` ${"─".repeat(Math.max(0, Math.min(60, width - 4)))}`} ); } -/** Bottom input row: `>` prompt + buffer on one line. */ +/** Indent every wrapped line of `text` by `pad` spaces within `width` columns. */ +function indent(text: string, pad: number, width: number): string { + return wrapText(text, width) + .split("\n") + .map((l) => (l.length ? " ".repeat(pad) + l : "")) + .join("\n"); +} + +/** Bottom input row: `❯` prompt + buffer on one line, with a subtle divider. */ function InputView(props: { state: AppStateRest; width: number }): ReactNode { const { state, width } = props; const shown = state.input.buffer.slice( @@ -302,9 +314,13 @@ function InputView(props: { state: AppStateRest; width: number }): ReactNode { Math.max(0, state.input.cursor) + 40, ); return ( - - {"> "} - {shown} + + {" " + "─".repeat(Math.max(0, width - 4))} + + {" ❯ "} + {shown} + {"▏"} + ); }