| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161 |
- /**
- * 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<E extends SseEventLike>(
- response: Response,
- onEvent: (event: E) => void,
- options: { terminateOn?: readonly string[] } = {}
- ): Promise<SseConsumeResult> {
- 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 中断,取消失败无需处理
- }
- }
- }
|