// SSE 帧解析 + 草稿流归一(对齐 C3 / ARCH §7.3)。 // 帧格式:`event:\ndata:\n\n`。 // 纯逻辑,便于单测(不依赖浏览器/DOM)。 export interface TokenEvent { event: "token"; data: { text: string }; } export interface DoneEvent { event: "done"; data: { length: number }; } export interface ErrorEvent { event: "error"; data: { code: string; message: string; request_id?: string | null }; } export type SseEvent = TokenEvent | DoneEvent | ErrorEvent; const KNOWN_EVENTS = new Set(["token", "done", "error"]); // 把一个完整 SSE 块(多行)解析成事件;无法解析则返回 null(跳过)。 export function parseSseBlock(block: string): SseEvent | null { let event = ""; const dataLines: string[] = []; for (const rawLine of block.split("\n")) { const line = rawLine.replace(/\r$/, ""); if (line.startsWith(":")) continue; // 注释/心跳 const sep = line.indexOf(":"); if (sep === -1) continue; const field = line.slice(0, sep); const value = line.slice(sep + 1).replace(/^ /, ""); if (field === "event") event = value; else if (field === "data") dataLines.push(value); } if (!KNOWN_EVENTS.has(event) || dataLines.length === 0) return null; let data: unknown; try { data = JSON.parse(dataLines.join("\n")); } catch { return null; } return { event, data } as SseEvent; } // 增量缓冲:吃进一段文本,吐出已完成的事件块(以空行分隔),保留未完成尾部。 export class SseFrameBuffer { private buf = ""; push(chunk: string): SseEvent[] { this.buf += chunk; const events: SseEvent[] = []; let idx: number; // 块之间以空行(\n\n,兼容 \r\n\r\n)分隔。 while ((idx = this.findBoundary(this.buf)) !== -1) { const block = this.buf.slice(0, idx.valueOf()); this.buf = this.buf.slice(this.boundaryEnd(this.buf, idx)); const evt = parseSseBlock(block); if (evt) events.push(evt); } return events; } private findBoundary(s: string): number { const a = s.indexOf("\n\n"); const b = s.indexOf("\r\n\r\n"); if (a === -1) return b; if (b === -1) return a; return Math.min(a, b); } private boundaryEnd(s: string, idx: number): number { return s.startsWith("\r\n\r\n", idx) ? idx + 4 : idx + 2; } } export type StreamPhase = "idle" | "streaming" | "done" | "error" | "aborted"; export interface StreamState { phase: StreamPhase; text: string; error: { code: string; message: string; request_id?: string | null } | null; } export const initialStreamState: StreamState = { phase: "idle", text: "", error: null, }; // 纯 reducer:把单个事件折叠进状态(打字机文本累积)。 export function reduceStream(state: StreamState, event: SseEvent): StreamState { switch (event.event) { case "token": return { ...state, phase: "streaming", text: state.text + event.data.text, }; case "done": return { ...state, phase: "done" }; case "error": return { ...state, phase: "error", error: event.data }; default: return state; } }