// 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 narrowSseEvent(event, data); } function isRecord(v: unknown): v is Record { return typeof v === "object" && v !== null && !Array.isArray(v); } // 逐事件类型守卫:按 event 名收窄 data 形状;形状不符则返回 null(安全跳过)。 function narrowSseEvent(event: string, data: unknown): SseEvent | null { if (!isRecord(data)) return null; switch (event) { case "token": return typeof data.text === "string" ? { event: "token", data: { text: data.text } } : null; case "done": return typeof data.length === "number" ? { event: "done", data: { length: data.length } } : null; case "error": return typeof data.code === "string" && typeof data.message === "string" ? { event: "error", data: { code: data.code, message: data.message, request_id: typeof data.request_id === "string" ? data.request_id : null, }, } : null; default: return null; } } // 增量缓冲:吃进一段文本,吐出已完成的事件块(以空行分隔),保留未完成尾部。 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); 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; } }