/** * 共用 SSE 客户端:跨 TCP 分片保帧状态 + Last-Event-ID 断点续传。 * * 三平台桥原本各自手写 SSE 解析,且都有同一个 bug: * - `evt` / `data` 是每次 `read()` 的局部变量,TCP 把一帧 * `event: xxx\ndata: {...}\n\n` 切在换行处时,第一段只剩 * `event:` 而第二段只有 `data:` —— 整帧被静默丢弃。 * - 重连不带 `Last-Event-ID`,断线期间的事件只在服务端环形 * 缓冲里等着,永远回放不出来(Gateway 有 per-agent ring buffer, * pi 与 homeagent 已正确利用,DSH/opencode 没有)。 * * 这个模块把 pi 桥 `src/gateway.mjs` 里那份验证过的实现抽成共用件, * 三桥逐字节同源(deploy/check-shared-libs.sh 校验)。 */ /** * 增量 SSE 帧解析器。 * * `push(chunk)` 可以喂任意切分的文本片段,返回本次完整解析出的事件数组。 * 所有跨帧状态(缓冲、当前 event/data/id)都保存在闭包里,**不随 chunk 重置** —— * 这正是原实现丢帧的根因。 * * 协议细节: * - `:` 开头 = 注释/心跳,忽略 * - `id:` / `event:` / `data:` 各取字段;`data:` 后的单个空格是分隔符 * - 多行 data 用 `\n` 拼接 * - 空行 = 帧结束;只有 event 与 data 都非空才派发(与旧行为一致) * - 兼容 CRLF * - `lastEventId` 在**派发之前**记下:回调抛异常也不该让断点回退。 * * @returns {{push: (chunk: string) => Array<{event: string, data: string, id: string}>, * reset: () => void, * lastEventId: () => string, * setLastEventId: (id: string) => void}} */ export function createFrameParser() { let buffer = ''; let lastEventId = ''; let curEvent = ''; let curData = ''; let curId = ''; function push(chunk) { buffer += chunk; const events = []; const lines = buffer.split('\n'); // 最后一段可能是被切断的半行,留到下一个 chunk buffer = lines.pop() ?? ''; for (let line of lines) { if (line.length > 0 && line.charAt(line.length - 1) === '\r') { line = line.slice(0, -1); } if (line.startsWith(':')) continue; if (line.startsWith('id:')) { curId = line.slice(3).trim(); } else if (line.startsWith('event:')) { curEvent = line.slice(6).trim(); } else if (line.startsWith('data:')) { let value = line.slice(5); if (value.startsWith(' ')) value = value.slice(1); curData = curData.length > 0 ? `${curData}\n${value}` : value; } else if (line === '') { if (curEvent.length > 0 && curData.length > 0) { if (curId.length > 0) lastEventId = curId; events.push({ event: curEvent, data: curData, id: curId }); } curEvent = ''; curData = ''; curId = ''; } } return events; } function reset() { buffer = ''; curEvent = ''; curData = ''; curId = ''; } return { push, reset, lastEventId: () => lastEventId, setLastEventId: (id) => { lastEventId = id || ''; }, }; } /** * @param {object} deps * @param {() => Record} deps.authHeaders 认证头(每次重连重取,密钥可能已换) * @param {string} deps.baseURL Gateway 基地址(不带末尾 /) * @param {string} deps.path SSE 路径(如 /api/v1/events/stream) * @param {(evt: string, data: any) => void} deps.onEvent 事件分发回调 * @param {(msg: string) => void} [deps.log] 日志回调(默认 console.error) * @returns {{stop: () => void}} stop() 终止重连与在途请求 */ export function createSSEClient({ authHeaders, baseURL, path, onEvent, log = console.error }) { const controller = new AbortController(); const parser = createFrameParser(); function stop() { controller.abort(); } function reconnect(delay) { if (controller.signal.aborted) return; setTimeout(() => connect(), delay); } function connect() { if (controller.signal.aborted) return; const headers = { ...authHeaders(), Accept: 'text/event-stream' }; // 只有 lastEventId 非空(= 已经收过事件)时才是重连:首次连接不带, // 否则服务端会把环形缓冲里的旧事件全回放一遍,插件重启后重复处理一批已处理的邮件。 const lastEventID = parser.lastEventId(); if (lastEventID) { headers['Last-Event-ID'] = lastEventID; log(`SSE 重连,从事件 ${lastEventID} 之后续传`); } fetch(`${baseURL}${path}`, { headers, signal: controller.signal }) .then((res) => { if (!res.ok || !res.body) { log(`SSE 建连失败: HTTP ${res.status}`); return reconnect(5000); } const reader = res.body.getReader(); const decoder = new TextDecoder(); function read() { reader.read().then(({ done, value }) => { if (done) { parser.reset(); return reconnect(3000); } for (const ev of parser.push(decoder.decode(value, { stream: true }))) { try { onEvent(ev.event, JSON.parse(ev.data)); } catch (e) { log(`SSE 事件处理失败: ${e?.message || e}`); } } read(); }).catch((e) => { if (controller.signal.aborted) return; log(`SSE 读取中断: ${e?.message || e}`); parser.reset(); reconnect(5000); }); } read(); }) .catch((e) => { if (controller.signal.aborted) return; log(`SSE 连接错误: ${e?.message || e}`); parser.reset(); reconnect(5000); }); } connect(); return { stop, /** * 清掉断点(不终止连接)。 * * 换 Gateway 地址时必须调:lastEventID 是**旧** Gateway 环形缓冲里的序号, * 拿去问新 Gateway 会命中一段完全无关的历史(或直接被拒), * 得到的事件属于别人的会话。 */ reset: () => parser.setLastEventId(''), }; }