fix(tui): sliding-window stream reconciler + UI polish

- Replace naive reconcileStream with stateful sliding-window replay
  merger: tracks per-message replay cursors (_posContent/_posReasoning)
  to deduplicate the segment-replay waves 9router emits. Verified
  against 970 real captured chunks (100% clean output).
- TranscriptMessage: colored [TAG] headers, 4-space indented body,
  thin separator lines between messages, dimmed reasoning.
- InputView: divider above prompt, '❯' prompt, cursor block.
- 75 tests pass, typecheck clean, biome clean.
This commit is contained in:
zesdex
2026-09-03 12:28:56 +07:00
committed by asepharyana
parent 30cc579773
commit d86d6aeef2
3 changed files with 170 additions and 56 deletions
+63 -22
View File
@@ -329,10 +329,13 @@ export interface ChatMessageDisplay {
content: string; content: string;
reasoning: string; reasoning: string;
timestamp: number; timestamp: number;
/** Sliding-window replay cursor for content + reasoning (transient). */
_posContent: number;
_posReasoning: number;
} }
export function makeChatMessage(role: Role, content: string): ChatMessageDisplay { 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. */ /** 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]; const msg = this.messages[this.messages.length - 1];
if (msg && msg.role === Roles.Assistant) { if (msg && msg.role === Roles.Assistant) {
if (isReasoning) { if (isReasoning) {
msg.reasoning = reconcileStream(msg.reasoning, text); msg.reasoning = reconcileStream(msg.reasoning, text, msg, true);
} else { } else {
msg.content = reconcileStream(msg.content, text); msg.content = reconcileStream(msg.content, text, msg, false);
} }
this.dirty = true; 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: * Streaming backends disagree on chunk semantics:
* - Native OpenAI/Anthropic send small *incremental* deltas. * - Native OpenAI/Anthropic send small *incremental* deltas.
* - Some gateways/proxies (e.g. 9router's OpenAI-compatible layer) emit * - 9router/gateways re-send the FULL accumulated text as a repeated window:
* *full-text* chunks, where each chunk is the complete transcript so far. * 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), * Naive concatenation doubles the text (garbled output). We track a per-stream
* which manifests as garbled messages. This helper accepts both modes: * replay cursor `pos` into `accumulated`: a chunk that continues at `pos` (or
* - a chunk that is a strict prefix of the tail → a shorter replay → ignore; * restarts at 0, i.e. begins a fresh wave) is a replay → ignore it and advance
* - a chunk that extends/replays the whole accumulated text → take it as the * the cursor; only a chunk that matches nothing already present is brand-new
* new authoritative tail (full-text mode); * content, which we append at the end and reset the cursor.
* - otherwise → treat as an incremental delta and append. *
* @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 (chunk.length === 0) return accumulated;
if (accumulated.length === 0) return chunk; const posRef = { pos: isReasoning ? msg._posReasoning : msg._posContent };
if (accumulated.length >= chunk.length && accumulated.startsWith(chunk)) { if (accumulated.length === 0) {
// Shorter (or equal) replay of what we already hold — ignore. accumulated = chunk;
posRef.pos = chunk.length;
writePos();
return accumulated; return accumulated;
} }
if (chunk.length >= accumulated.length && chunk.startsWith(accumulated)) { // 1) A chunk that continues the current replay cursor is a replayed segment
// Full-text mode: the chunk is the whole transcript so far (matching our // (it already exists earlier in the text) → skip it, advance the cursor.
// accumulated text plus more). Adopt it wholesale. if (
return chunk; 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 ─────────────────────────────────────────────────── */ /* ── AppStateRest ─────────────────────────────────────────────────── */
+80 -23
View File
@@ -4,7 +4,7 @@
*/ */
import { describe, test, expect } from "bun:test"; import { describe, test, expect } from "bun:test";
import { parseCommand, applyCommand } from "./command.ts"; 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 { applyAction } from "./action.ts";
import { handleKey, decodeKey } from "./controller.ts"; import { handleKey, decodeKey } from "./controller.ts";
import { toControllerKey } from "./ui.tsx"; import { toControllerKey } from "./ui.tsx";
@@ -163,32 +163,89 @@ describe("toControllerKey (OpenTUI adapter)", () => {
}); });
}); });
describe("reconcileStream (full-text vs incremental chunk handling)", () => { describe("reconcileStream (9router sliding-window vs incremental)", () => {
test("appends incremental deltas", () => { function msg(): ChatMessageDisplay {
let s = reconcileStream("", "Hello"); return makeChatMessage(Roles.Assistant, "");
s = reconcileStream(s, "!"); }
s = reconcileStream(s, " How"); 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"); expect(s).toBe("Hello! How");
}); });
test("full-text mode: adopts a growing authoritative chunk", () => {
let s = reconcileStream("", "Hello"); test("real 9router wave: replayed segments never double", () => {
// 9router/proxy resends the *whole* transcript each chunk — adopt it. // Real captured stream (story turn), each fresh token restarts a replay
s = reconcileStream(s, "Hello world"); // wave from the head: "Here's a 3-sent" "ence story" " about" " a robo" …
s = reconcileStream(s, "Hello world!"); const m = msg();
expect(s).toBe("Hello world!"); const segs: string[] = [
}); "Here's a 3-sent",
test("shorter/equal replay of the tail is ignored (no duplication)", () => { "ence story",
const tail = "Hello world"; " about",
expect(reconcileStream(tail, tail)).toBe(tail); " a robo",
expect(reconcileStream(tail, "Hello")).toBe(tail); "t:\n\nA small repai",
}); "r robo",
test("partial-then-full garbled stream resolves to the full text", () => { "t name",
let s = reconcileStream("", "Hello! 更新"); "d Bolt",
s = reconcileStream(s, "Hello! 更新 Keep going"); ];
expect(s).toBe("Hello! 更新 Keep going"); // 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", () => { 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?");
}); });
}); });
+24 -8
View File
@@ -276,25 +276,37 @@ function TranscriptMessage(props: {
msg.role === Roles.Assistant && !msg.content.trim() msg.role === Roles.Assistant && !msg.content.trim()
? "" ? ""
: msg.content || "(empty)"; : 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 `<text>` is used (not `<markdown>`) so live
// streaming re-renders a single clean string — no overlap artifacts, and the
// body fills the full transcript width.
return ( return (
<box flexDirection="column" width="100%"> <box flexDirection="column" width="100%">
<text fg={color}>{` ${tag.padEnd(4)} `}</text> {/* Role tag header line */}
<text fg={color}>{` [${tag}]`}</text>
{/* Indented, word-wrapped body */}
{content.length > 0 && ( {content.length > 0 && (
<text fg={color}>{wrapText(content, Math.max(20, width - 4))}</text> <text fg={msg.role === Roles.Assistant ? "default" : color}>
{indent(content, 4, bodyWidth)}
</text>
)} )}
{/* Live reasoning (dimmed, kept to one line) */}
{msg.reasoning.length > 0 && ( {msg.reasoning.length > 0 && (
<text fg="gray">{` ⋯ ${oneLine(msg.reasoning)}`}</text> <text fg="gray">{` ⋯ ${oneLine(msg.reasoning)}`}</text>
)} )}
{/* Thin separator between messages */}
<text fg="#2a3340">{` ${"─".repeat(Math.max(0, Math.min(60, width - 4)))}`}</text>
</box> </box>
); );
} }
/** 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 { function InputView(props: { state: AppStateRest; width: number }): ReactNode {
const { state, width } = props; const { state, width } = props;
const shown = state.input.buffer.slice( 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, Math.max(0, state.input.cursor) + 40,
); );
return ( return (
<box flexDirection="column" width="100%">
<text fg="#2a3340">{" " + "─".repeat(Math.max(0, width - 4))}</text>
<box flexDirection="row" width="100%"> <box flexDirection="row" width="100%">
<text fg="cyan">{"> "}</text> <text fg="cyan">{" ❯ "}</text>
<text fg="white">{shown}</text> <text fg="white">{shown}</text>
<text fg="gray">{"▏"}</text>
</box>
</box> </box>
); );
} }