Files
MailUI4Agents/plugins/pi-mail-bridge/lib/sse-client.js
JianFeeeee c401eb2da2 fix(bridges): SSE 跨分片保帧 + Last-Event-ID;pi worker 有界重投;systemd 故障上报;清理误提交二进制
三个平台桥原本各自手写 SSE 解析,有两个共同的静默丢事件缺陷:
  1. evt/data 是每次 read() 的局部变量 —— TCP 把一帧
     'event: x\ndata: {...}\n\n' 切在换行处时,前半段的 event 名被丢掉、
     后半段只剩 data,整帧静默丢弃。表现为「新邮件偶尔收不到」
     「权限决策点了没反应」,日志里一个字都没有。
  2. 重连不带 Last-Event-ID —— 断线期间的事件留在服务端 per-agent 环形
     缓冲里永远回放不出来(pi 与 homeagent 已正确使用,DSH/opencode 没有)。

修法:抽出共用 lib/sse-client.js(三桥逐字节同源,check-shared-libs 校验),
把「跨 chunk 保帧状态」与「Last-Event-ID 断点续传」写对一次。pi 桥的
gateway.mjs 也改为复用同一实现(保留 reconfigure 时清断点的语义)。

pi worker 丢任务:worker 未回报 done 就退出(SIGKILL/OOM/崩溃)时,
主进程原来只记一行日志就 pump() —— 那封邮件永远没有回音。改为按
1s/2s 退避有界重投(默认 3 次),到上限记「放弃」并可观测。

systemd 故障上报:四个宿主服务接入 service-failure-notify.mjs 的
ExecStopPost/--report 与 ExecStartPost/--flush。进程内 uncaughtException
捕获不了 SIGKILL/OOM,只能由 systemd 统一覆盖。正常 stop/restart 不发信。

仓库卫生:server/server(24MB 构建产物,f9d757b 误提交)移出版本库。

测试:opencode 302 / dsh 335 / pi 391 全绿(新增 12 例 SSE 帧解析 +
2 例 worker 重投);Go 全量通过;四平台重启后在线且无错误。
2026-09-11 10:27:41 +08:00

180 lines
6.1 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/**
* 共用 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<string,string>} 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(''),
};
}