/** * SSE 流消费的统一实现。 * */ /** 本模块需要读取的最小事件形状;各业务自己的事件类型(SseEvent / SSERawEvent)都满足它 */ export interface SseEventLike { type?: string; session_id?: string; thread_id?: string; } /** 消费完一条 SSE 流的结果 */ export interface SseConsumeResult { /** 终止事件携带的会话 id,未携带时为空串 */ sessionId: string; /** 终止事件携带的线程 id,未携带时为空串 */ threadId: string; /** 解析失败被丢弃的 data 行数,0 表示整条流都解析正常 */ parseErrors: number; } /** 一次解析尝试的结果 */ type FlushOutcome = 'terminated' | 'dispatched' | 'incomplete' | 'invalid' | 'empty'; /** 默认的终止事件类型:收到其中之一即停止读取并释放连接 */ const DEFAULT_TERMINATE_ON = ['done', 'error', 'interrupt']; /** * 消费一个 SSE 响应体,逐事件回调,直到收到终止事件或流结束。 * @param response - fetch 返回的响应(须带 body) * @param onEvent - 每个事件的回调 * @param options.terminateOn - 终止事件类型,默认 done / error / interrupt */ export async function consumeSse( response: Response, onEvent: (event: E) => void, options: { terminateOn?: readonly string[] } = {} ): Promise { const terminateOn = options.terminateOn ?? DEFAULT_TERMINATE_ON; const body = response.body; if (!body) { throw new Error('SSE 响应缺少可读流'); } const reader = body.getReader(); const decoder = new TextDecoder(); const result: SseConsumeResult = { sessionId: '', threadId: '', parseErrors: 0 }; let buffer = ''; /** 当前事件已收集、尚未解析的 data 行 */ let pending: string[] = []; /** 丢弃坏数据时记账:解析失败不再无声跳过 */ const dropBadData = (raw: string) => { result.parseErrors += 1; console.warn('SSE数据解析失败:', raw); }; /** * 尝试把当前累积的 data 行解析成一个事件并派发。 * 解析失败时保留 pending 不清理——是"多行事件的前半段"还是"坏数据"由调用方判断。 */ const flushPending = (): FlushOutcome => { if (pending.length === 0) return 'empty'; const raw = pending.join('\n'); // 规范规定 data 为空的事件不派发;服务端也有用 `data:` 空行做心跳的写法,故不算解析错误 if (raw.trim() === '') { pending = []; return 'empty'; } let event: E; try { event = JSON.parse(raw) as E; } catch { return 'incomplete'; } // 非对象负载(`data: null` / `data: 123`)没有事件语义,交给回调只会让它崩在取字段上。 // 按坏数据丢弃并继续读流——与重构前"回调抛错被 try 吞掉后继续"的最终表现一致。 if (typeof event !== 'object' || !event) { dropBadData(raw); pending = []; return 'invalid'; } pending = []; onEvent(event); if (event.type && terminateOn.includes(event.type)) { result.sessionId = event.session_id || ''; result.threadId = event.thread_id || ''; return 'terminated'; } return 'dispatched'; }; /** 处理一行;返回 true 表示已收到终止事件 */ const handleLine = (rawLine: string): boolean => { // 兼容 CRLF / CR 行尾 const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine; // 空行是规范定义的事件边界:到这里还解析不出来,才能判定为坏数据并丢弃 if (line === '') { const outcome = flushPending(); if (outcome === 'incomplete') { dropBadData(pending.join('\n')); pending = []; } return outcome === 'terminated'; } const colon = line.indexOf(':'); // 没有冒号的整行都是字段名、冒号在行首则是注释行(心跳),两者都不消费; // 字段名与值之间允许 0 个空格,所以这里不能用 startsWith('data: ') 判断 if (colon <= 0 || line.slice(0, colon) !== 'data') return false; const value = line.slice(colon + 1); pending.push(value.startsWith(' ') ? value.slice(1) : value); // 每收到一行就尝试解析累积的 data 行: // - 能解析 → 立即派发(兼容事件之间不发送空行的服务端,与原实现行为一致) // - 非对象负载 → 已记坏数据并清理干净,结束 // - 解析不出来且已积累多行 → 首行是坏数据,丢掉它再试,避免一行坏 JSON 连坐后续事件 // - 只剩一行 → 保留,它可能是多行事件的前半段 for (;;) { const outcome = flushPending(); if (outcome === 'terminated') return true; if (outcome === 'dispatched' || outcome === 'invalid') return false; if (pending.length <= 1) return false; dropBadData(pending.shift() as string); } }; try { for (;;) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); buffer = lines.pop() || ''; for (const line of lines) { if (handleLine(line)) return result; } } return result; } finally { try { await reader.cancel(); } catch { // 流已自然结束,或已被调用方的 AbortController 中断,取消失败无需处理 } } }