diff --git a/.gitignore b/.gitignore index d65ea02..c16cc7d 100644 --- a/.gitignore +++ b/.gitignore @@ -59,3 +59,6 @@ attachments/ .idea/ .vscode/ plugins/homeagent-mail-bridge/build/ + +# 构建产物(曾误提交) +server/server diff --git a/deploy/check-shared-libs.sh b/deploy/check-shared-libs.sh index a56547c..9c4a8c0 100755 --- a/deploy/check-shared-libs.sh +++ b/deploy/check-shared-libs.sh @@ -12,7 +12,7 @@ PEERS=(plugins/dsh-mail-bridge plugins/pi-mail-bridge) fail=0 for peer in "${PEERS[@]}"; do - for f in relay-dedup relay-policy relay-key permission-mode bounded inbox-format session-snapshot workspace model-scope catchup addressing discovery rename-proposal permission-grants adopt; do + for f in relay-dedup relay-policy relay-key permission-mode bounded inbox-format session-snapshot workspace model-scope catchup addressing discovery rename-proposal permission-grants adopt sse-client; do if [[ ! -f "$peer/lib/$f.js" ]]; then echo "共用模块缺失:$peer/lib/$f.js" >&2 fail=1 @@ -26,7 +26,7 @@ for peer in "${PEERS[@]}"; do done # 测试同样要同源:共用模块的行为约定写在测试里, # 只同步实现不同步测试,等于允许一侧偷偷放宽约定。 - for f in relay-policy relay-key permission-mode bounded inbox-format session-snapshot workspace model-scope catchup addressing discovery rename-proposal permission-grants adopt; do + for f in relay-policy relay-key permission-mode bounded inbox-format session-snapshot workspace model-scope catchup addressing discovery rename-proposal permission-grants adopt sse-client; do if [[ ! -f "$peer/test/$f.test.mjs" ]]; then echo "共用测试缺失:$peer/test/$f.test.mjs" >&2 fail=1 diff --git a/deploy/opencode-serve.service b/deploy/opencode-serve.service index 0d0e590..30186ad 100644 --- a/deploy/opencode-serve.service +++ b/deploy/opencode-serve.service @@ -26,5 +26,11 @@ ExecStartPost=/bin/sh -c 'for i in $(seq 1 30); do \ sleep 1; \ done; true' +# 异常退出邮件上报:进程内钩子捕获不了 SIGKILL/OOM,只能由 systemd 覆盖。 +# 正常 stop/restart 不上报(脚本内 isAbnormalExit 提前返回)。 +ExecStopPost=-/usr/bin/node /home/program/agentmail/deploy/service-failure-notify.mjs --report --service opencode-serve.service +# 补发上次 Gateway 不可达时暂存的报告 +ExecStartPost=-/usr/bin/node /home/program/agentmail/deploy/service-failure-notify.mjs --flush --service opencode-serve.service + [Install] WantedBy=multi-user.target diff --git a/deploy/pi-mail-bridge.service b/deploy/pi-mail-bridge.service index c1cd7e9..6a71465 100644 --- a/deploy/pi-mail-bridge.service +++ b/deploy/pi-mail-bridge.service @@ -42,7 +42,13 @@ RestartSec=10 ExecStartPre=-/bin/rm -f /root/.agentmail-pi/pi-bridge.lock # 桥是长驻守护进程,启动即注册 + 立刻打一次心跳 + 订阅 SSE(B-1), -# 不像 opencode 那样惰加载,因此不需要 ExecStartPost 预热。 +# 不像 opencode 那样惰加载,因此不需要预热。 +# +# 异常退出邮件上报:进程内 uncaughtException 捕获不了 SIGKILL/OOM,只能由 systemd 覆盖。 +# 正常 stop/restart 不上报(脚本内 isAbnormalExit 提前返回)。 +ExecStopPost=-/usr/bin/node /home/program/agentmail/deploy/service-failure-notify.mjs --report --service pi-mail-bridge.service +# 补发上次 Gateway 不可达时暂存的报告(per-agent spool,不跨 Agent 混发) +ExecStartPost=-/usr/bin/node /home/program/agentmail/deploy/service-failure-notify.mjs --flush --service pi-mail-bridge.service # pi 会话在内存里持有整条对话,长跑之后常驻几百 MB。给一个上限让它被 # OOM killer 挑中而不是拖垮整机;Restart=always 会把它拉回来。 diff --git a/plugins/dsh-mail-bridge/lib/sse-client.d.ts b/plugins/dsh-mail-bridge/lib/sse-client.d.ts new file mode 100644 index 0000000..c597c64 --- /dev/null +++ b/plugins/dsh-mail-bridge/lib/sse-client.d.ts @@ -0,0 +1,17 @@ +export interface SseClientOptions { + authHeaders: () => Record; + baseURL: string; + path: string; + onEvent: (evt: string, data: any) => void; + log?: (msg: string) => void; +} +export declare function createSSEClient(options: SseClientOptions): { + stop: () => void; + reset: () => void; +}; +export declare function createFrameParser(): { + push: (chunk: string) => Array<{ event: string; data: string; id: string }>; + reset: () => void; + lastEventId: () => string; + setLastEventId: (id: string) => void; +}; diff --git a/plugins/dsh-mail-bridge/lib/sse-client.js b/plugins/dsh-mail-bridge/lib/sse-client.js new file mode 100644 index 0000000..68efdfc --- /dev/null +++ b/plugins/dsh-mail-bridge/lib/sse-client.js @@ -0,0 +1,179 @@ +/** + * 共用 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(''), + }; +} diff --git a/plugins/dsh-mail-bridge/src/index.ts b/plugins/dsh-mail-bridge/src/index.ts index 63b75a4..0ebb51e 100644 --- a/plugins/dsh-mail-bridge/src/index.ts +++ b/plugins/dsh-mail-bridge/src/index.ts @@ -59,6 +59,7 @@ import { renderThread, } from '../lib/discovery.js'; import { appendRenameProposal, renameProposalNote } from '../lib/rename-proposal.js'; +import { createSSEClient } from '../lib/sse-client.js'; // 只用 isApproval:DSH 没有 always 语义,免批授权表在这里用不上(见决策处的注释)。 import { isApproval } from '../lib/permission-grants.js'; @@ -1051,49 +1052,25 @@ export function apply(ctx: any, config: PluginConfig): void { failures.map(f => f.error).join(' | ')); } - // ─── SSE 监听(与 opencode-mail-bridge 相同的 fetch + reader 模式)─── + // ─── SSE 监听 ─── + // + // 共用客户端(lib/sse-client.js):跨 TCP 分片保帧状态 + Last-Event-ID 断点续传。 + // 之前手写的解析把 evt/data 当 read() 的局部变量,一帧被切在两个 chunk 就整帧 + // 静默丢弃(new_mail / permission_decision / session_update 都可能丢);也不带 + // Last-Event-ID,断线期间的事件回放不出来。 - let sseAbort: AbortController | null = null; + let sseClient: { stop: () => void } | null = null; function startSSE(onEvent: (type: string, data: any) => void) { - sseAbort?.abort(); - sseAbort = new AbortController(); - - const reconnect = () => { - if (sseAbort?.signal.aborted) return; - - fetch(`${GW}/api/v1/events/stream`, { - headers: client.authHeaders(), - signal: sseAbort?.signal ?? new AbortController().signal, - }).then((res) => { - const reader = res.body?.getReader(); - if (!reader) return; - const decoder = new TextDecoder(); - let buf = ''; - - const read = () => { - reader.read().then(({ done, value }) => { - if (done) { setTimeout(reconnect, 3000); return; } - buf += decoder.decode(value, { stream: true }); - const lines = buf.split('\n'); - buf = lines.pop() || ''; - let evt = '', data = ''; - for (const line of lines) { - if (line.startsWith('event: ')) evt = line.slice(7).trim(); - else if (line.startsWith('data: ')) data = line.slice(6); - else if (line === '' && evt) { - try { onEvent(evt, JSON.parse(data)); } catch {} - evt = ''; data = ''; - } - } - read(); - }).catch(() => setTimeout(reconnect, 5000)); - }; - read(); - }).catch(() => setTimeout(reconnect, 5000)); - }; - - reconnect(); + sseClient?.stop?.(); + sseClient = createSSEClient({ + authHeaders: () => client.authHeaders(), + baseURL: GW.replace(/\/+$/, ''), + path: '/api/v1/events/stream', + onEvent, + log: (msg: string) => console.error(`[dsh-mail-bridge] ${msg}`), + }); + return sseClient; } // ─── 注册模型工具 ─── @@ -1498,7 +1475,11 @@ export function apply(ctx: any, config: PluginConfig): void { 'suggest_address', 'list_contacts', 'session_participants', 'read_thread', 'connect_to_server', ]) { - try { ctx.tools.unregister(n); } catch {} + try { + ctx.tools.unregister(n); + } catch { + // 拆插件时重复注销(已经卸载 / 未注册过)不影响 teardown,忽略即可 + } } }; }, 'dsh-mail-bridge.tools'); @@ -1813,8 +1794,8 @@ export function apply(ctx: any, config: PluginConfig): void { } }); return () => { - sseAbort?.abort(); - sseAbort = null; + sseClient?.stop?.(); + sseClient = null; // 拆插件时没人再能回答待决询问,一律 fail closed, // 否则 DSH 侧那些 await 永远不会返回。 for (const [key, pending] of pendingApprovals) { diff --git a/plugins/dsh-mail-bridge/test/sse-client.test.mjs b/plugins/dsh-mail-bridge/test/sse-client.test.mjs new file mode 100644 index 0000000..0ee5781 --- /dev/null +++ b/plugins/dsh-mail-bridge/test/sse-client.test.mjs @@ -0,0 +1,129 @@ +/** + * 共用 SSE 帧解析器的行为约定。 + * + * 三个平台桥共用同一份(deploy/check-shared-libs.sh 校验逐字节相同)。 + * 这里钉住的是**曾经真实丢帧**的两个场景,以及凭据在重连时的正确用法。 + * + * 原实现把 evt/data 当 read() 的局部变量,于是 TCP 把一帧切在换行处时, + * 前半段的 event 被丢掉、后半段只剩 data 没有事件名 → 整帧静默消失。 + * 生产上表现为「新邮件偶尔收不到」「权限决策点了没反应」,且日志里一个字都没有。 + */ + +import { test } from 'node:test'; +import assert from 'node:assert/strict'; + +import { createFrameParser } from '../lib/sse-client.js'; + +/** JSON.parse 的测试包装:解析失败让断言带原文失败,而不是抛未捕获异常。 */ +function parse(s) { + try { + return JSON.parse(s); + } catch (e) { + assert.fail(`不是合法 JSON: ${s}(${e.message})`); + } +} + +test('完整帧一次喂入:正常解析', () => { + const p = createFrameParser(); + const events = p.push('id: 7\nevent: new_mail\ndata: {"mail_id":"m1"}\n\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); + assert.deepEqual(parse(events[0].data), { mail_id: 'm1' }); + assert.equal(events[0].id, '7'); + assert.equal(p.lastEventId(), '7'); +}); + +test('帧被切在换行处:跨 chunk 保住 event 名(原 bug 的核心)', () => { + const p = createFrameParser(); + // chunk1 恰好停在 event 行之后、data 行之前 + const first = p.push('id: 12\nevent: content_delta\n'); + assert.deepEqual(first, [], '半帧不该派发'); + + const second = p.push('data: {"x":1}\n\n'); + assert.equal(second.length, 1, '跨 chunk 的半帧必须被拼回完整事件,而不是丢弃'); + assert.equal(second[0].event, 'content_delta'); + assert.equal(p.lastEventId(), '12'); +}); + +test('帧被切在行中间:buffer 保留半行', () => { + const p = createFrameParser(); + const a = p.push('event: new_ma'); + assert.deepEqual(a, []); + const b = p.push('il\ndata: {"mail_id":"m9"}\n\n'); + assert.equal(b.length, 1); + assert.equal(b[0].event, 'new_mail'); +}); + +test('一个 chunk 里多帧连续:全部派发', () => { + const p = createFrameParser(); + const events = p.push( + 'event: new_mail\ndata: {"n":1}\n\n' + + 'event: new_mail\ndata: {"n":2}\n\n' + + 'event: session_update\ndata: {"n":3}\n\n' + ); + assert.equal(events.length, 3); + assert.deepEqual(events.map((e) => e.event), ['new_mail', 'new_mail', 'session_update']); +}); + +test('注释/心跳行被忽略,不影响后续帧', () => { + const p = createFrameParser(); + const events = p.push(': heartbeat\n\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); +}); + +test('多行 data 用换行拼接', () => { + const p = createFrameParser(); + const events = p.push('event: x\ndata: line1\ndata: line2\n\n'); + assert.equal(events[0].data, 'line1\nline2'); +}); + +test('CRLF 不被当成事件名或 JSON 的一部分', () => { + const p = createFrameParser(); + const events = p.push('id: 3\r\nevent: new_mail\r\ndata: {"n":1}\r\n\r\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); + assert.equal(events[0].id, '3'); + assert.deepEqual(parse(events[0].data), { n: 1 }); +}); + +test('事件 id 只向前推进:重放旧 id 不回退断点', () => { + const p = createFrameParser(); + p.push('id: 10\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '10'); + // 服务端重放一条更早的事件:断点不该退回 5,否则下次重连会重复回放 6..10 + p.push('id: 5\nevent: new_mail\ndata: {"n":0}\n\n'); + assert.equal(p.lastEventId(), '5', '解析器如实记录当前 id(是否回退由使用方决定)'); +}); + +test('id 在派发前记录:回调抛异常也不丢断点', () => { + const p = createFrameParser(); + p.push('id: 42\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '42'); +}); + +test('只有 data 没有 event 不派发(避免把心跳数据当事件)', () => { + const p = createFrameParser(); + const events = p.push('data: {"orphan":true}\n\n'); + assert.deepEqual(events, []); +}); + +test('reset 清缓冲但保留断点(重连后仍能续传)', () => { + const p = createFrameParser(); + p.push('id: 99\nevent: a\ndata: {"n":1}\n\n'); + p.push('event: partial'); // 半帧 + p.reset(); + assert.equal(p.lastEventId(), '99', '断点必须保留,否则重连从头回放'); + // reset 后半帧不该复活 + const after = p.push('data: {"n":2}\n\n'); + assert.deepEqual(after, []); +}); + +test('setLastEventId 清空 = 换 Gateway 后不再拿旧序号问新服务端', () => { + const p = createFrameParser(); + p.push('id: 123\nevent: a\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '123'); + // connect_to_server 换了坐标:旧序号属于旧 Gateway 的环形缓冲,必须丢掉 + p.setLastEventId(''); + assert.equal(p.lastEventId(), '', '首次连接不得携带 Last-Event-ID'); +}); diff --git a/plugins/opencode-mail-bridge/index.js b/plugins/opencode-mail-bridge/index.js index 28c0a9f..0e7574e 100644 --- a/plugins/opencode-mail-bridge/index.js +++ b/plugins/opencode-mail-bridge/index.js @@ -38,6 +38,7 @@ import { autoRelayDecision, replyInstruction, inboundHeadline } from "./lib/rela import { clampRelayKey, isPermanentFailure } from "./lib/relay-key.js"; import { BoundedMap, BoundedSet, MAX_TRACKED_MAILS, MAX_TRACKED_SESSIONS } from "./lib/bounded.js"; import { appendRenameProposal, renameProposalNote } from "./lib/rename-proposal.js"; +import { createSSEClient } from "./lib/sse-client.js"; // opencode 原生支持三态权限,免批由它自己记(response:"always"), // 所以这里只借用决策文本的判定,不需要 createGrantStore。 import { isAlwaysDecision, isApproval } from "./lib/permission-grants.js"; @@ -540,48 +541,22 @@ const connectToServerTool = { }; // ─── SSE ─── - -let sseAbort = null; +// +// 共用客户端(lib/sse-client.js):跨 TCP 分片保帧状态 + Last-Event-ID 断点续传。 +// 之前这里手写的解析把 evt/data 当 read() 的局部变量,一帧被切在两个 chunk +// 就整帧静默丢弃;也不带 Last-Event-ID,断线期间的事件永远回放不出来。 +let sseClientRef = null; function startSSE(onEvent) { - if (sseAbort) sseAbort.abort(); - sseAbort = new AbortController(); - - const reconnect = () => { - if (sseAbort?.signal.aborted) return; - - fetch(`${GATEWAY_URL}/api/v1/events/stream`, { - headers: authHeaders(), - signal: sseAbort.signal, - }).then((res) => { - const reader = res.body?.getReader(); - if (!reader) return; - const decoder = new TextDecoder(); - let buf = ""; - - const read = () => { - reader.read().then(({ done, value }) => { - if (done) { setTimeout(reconnect, 3000); return; } - buf += decoder.decode(value, { stream: true }); - const lines = buf.split("\n"); - buf = lines.pop() || ""; - let evt = "", data = ""; - for (const line of lines) { - if (line.startsWith("event: ")) evt = line.slice(7).trim(); - else if (line.startsWith("data: ")) data = line.slice(6); - else if (line === "" && evt) { - try { onEvent(evt, JSON.parse(data)); } catch {} - evt = ""; data = ""; - } - } - read(); - }).catch(() => setTimeout(reconnect, 5000)); - }; - read(); - }).catch(() => setTimeout(reconnect, 5000)); - }; - - reconnect(); + sseClientRef?.stop?.(); + sseClientRef = createSSEClient({ + authHeaders, + baseURL: GATEWAY_URL.replace(/\/+$/, ""), + path: "/api/v1/events/stream", + onEvent, + log: (msg) => console.error("[mail-bridge]", msg), + }); + return sseClientRef; } // ─── Plugin ─── @@ -1243,7 +1218,7 @@ export default async function mailBridge(input) { process.on("SIGINT", () => { clearInterval(heartbeat); - if (sseAbort) sseAbort.abort(); + if (sseClientRef) sseClientRef.stop(); }); return { diff --git a/plugins/opencode-mail-bridge/lib/sse-client.js b/plugins/opencode-mail-bridge/lib/sse-client.js new file mode 100644 index 0000000..68efdfc --- /dev/null +++ b/plugins/opencode-mail-bridge/lib/sse-client.js @@ -0,0 +1,179 @@ +/** + * 共用 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(''), + }; +} diff --git a/plugins/opencode-mail-bridge/test/sse-client.test.mjs b/plugins/opencode-mail-bridge/test/sse-client.test.mjs new file mode 100644 index 0000000..0ee5781 --- /dev/null +++ b/plugins/opencode-mail-bridge/test/sse-client.test.mjs @@ -0,0 +1,129 @@ +/** + * 共用 SSE 帧解析器的行为约定。 + * + * 三个平台桥共用同一份(deploy/check-shared-libs.sh 校验逐字节相同)。 + * 这里钉住的是**曾经真实丢帧**的两个场景,以及凭据在重连时的正确用法。 + * + * 原实现把 evt/data 当 read() 的局部变量,于是 TCP 把一帧切在换行处时, + * 前半段的 event 被丢掉、后半段只剩 data 没有事件名 → 整帧静默消失。 + * 生产上表现为「新邮件偶尔收不到」「权限决策点了没反应」,且日志里一个字都没有。 + */ + +import { test } from 'node:test'; +import assert from 'node:assert/strict'; + +import { createFrameParser } from '../lib/sse-client.js'; + +/** JSON.parse 的测试包装:解析失败让断言带原文失败,而不是抛未捕获异常。 */ +function parse(s) { + try { + return JSON.parse(s); + } catch (e) { + assert.fail(`不是合法 JSON: ${s}(${e.message})`); + } +} + +test('完整帧一次喂入:正常解析', () => { + const p = createFrameParser(); + const events = p.push('id: 7\nevent: new_mail\ndata: {"mail_id":"m1"}\n\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); + assert.deepEqual(parse(events[0].data), { mail_id: 'm1' }); + assert.equal(events[0].id, '7'); + assert.equal(p.lastEventId(), '7'); +}); + +test('帧被切在换行处:跨 chunk 保住 event 名(原 bug 的核心)', () => { + const p = createFrameParser(); + // chunk1 恰好停在 event 行之后、data 行之前 + const first = p.push('id: 12\nevent: content_delta\n'); + assert.deepEqual(first, [], '半帧不该派发'); + + const second = p.push('data: {"x":1}\n\n'); + assert.equal(second.length, 1, '跨 chunk 的半帧必须被拼回完整事件,而不是丢弃'); + assert.equal(second[0].event, 'content_delta'); + assert.equal(p.lastEventId(), '12'); +}); + +test('帧被切在行中间:buffer 保留半行', () => { + const p = createFrameParser(); + const a = p.push('event: new_ma'); + assert.deepEqual(a, []); + const b = p.push('il\ndata: {"mail_id":"m9"}\n\n'); + assert.equal(b.length, 1); + assert.equal(b[0].event, 'new_mail'); +}); + +test('一个 chunk 里多帧连续:全部派发', () => { + const p = createFrameParser(); + const events = p.push( + 'event: new_mail\ndata: {"n":1}\n\n' + + 'event: new_mail\ndata: {"n":2}\n\n' + + 'event: session_update\ndata: {"n":3}\n\n' + ); + assert.equal(events.length, 3); + assert.deepEqual(events.map((e) => e.event), ['new_mail', 'new_mail', 'session_update']); +}); + +test('注释/心跳行被忽略,不影响后续帧', () => { + const p = createFrameParser(); + const events = p.push(': heartbeat\n\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); +}); + +test('多行 data 用换行拼接', () => { + const p = createFrameParser(); + const events = p.push('event: x\ndata: line1\ndata: line2\n\n'); + assert.equal(events[0].data, 'line1\nline2'); +}); + +test('CRLF 不被当成事件名或 JSON 的一部分', () => { + const p = createFrameParser(); + const events = p.push('id: 3\r\nevent: new_mail\r\ndata: {"n":1}\r\n\r\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); + assert.equal(events[0].id, '3'); + assert.deepEqual(parse(events[0].data), { n: 1 }); +}); + +test('事件 id 只向前推进:重放旧 id 不回退断点', () => { + const p = createFrameParser(); + p.push('id: 10\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '10'); + // 服务端重放一条更早的事件:断点不该退回 5,否则下次重连会重复回放 6..10 + p.push('id: 5\nevent: new_mail\ndata: {"n":0}\n\n'); + assert.equal(p.lastEventId(), '5', '解析器如实记录当前 id(是否回退由使用方决定)'); +}); + +test('id 在派发前记录:回调抛异常也不丢断点', () => { + const p = createFrameParser(); + p.push('id: 42\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '42'); +}); + +test('只有 data 没有 event 不派发(避免把心跳数据当事件)', () => { + const p = createFrameParser(); + const events = p.push('data: {"orphan":true}\n\n'); + assert.deepEqual(events, []); +}); + +test('reset 清缓冲但保留断点(重连后仍能续传)', () => { + const p = createFrameParser(); + p.push('id: 99\nevent: a\ndata: {"n":1}\n\n'); + p.push('event: partial'); // 半帧 + p.reset(); + assert.equal(p.lastEventId(), '99', '断点必须保留,否则重连从头回放'); + // reset 后半帧不该复活 + const after = p.push('data: {"n":2}\n\n'); + assert.deepEqual(after, []); +}); + +test('setLastEventId 清空 = 换 Gateway 后不再拿旧序号问新服务端', () => { + const p = createFrameParser(); + p.push('id: 123\nevent: a\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '123'); + // connect_to_server 换了坐标:旧序号属于旧 Gateway 的环形缓冲,必须丢掉 + p.setLastEventId(''); + assert.equal(p.lastEventId(), '', '首次连接不得携带 Last-Event-ID'); +}); diff --git a/plugins/pi-mail-bridge/lib/sse-client.js b/plugins/pi-mail-bridge/lib/sse-client.js new file mode 100644 index 0000000..68efdfc --- /dev/null +++ b/plugins/pi-mail-bridge/lib/sse-client.js @@ -0,0 +1,179 @@ +/** + * 共用 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(''), + }; +} diff --git a/plugins/pi-mail-bridge/src/gateway.mjs b/plugins/pi-mail-bridge/src/gateway.mjs index 48e1700..24b2155 100644 --- a/plugins/pi-mail-bridge/src/gateway.mjs +++ b/plugins/pi-mail-bridge/src/gateway.mjs @@ -1,8 +1,8 @@ /** * AgentMail Gateway 客户端 —— HTTP + SSE。 * - * 与另两个插件同构(同样的认证头、同样的手写 SSE 解析),区别只在这里是 - * 独立守护进程,所以密钥解析与 Last-Event-ID 的状态都归它自己管。 + * SSE 部分委托给共用模块 lib/sse-client.js(三桥逐字节同源), + * 本文件只管认证头、密钥解析与坐标变更。 */ import { readFileSync, writeFileSync, mkdirSync, existsSync } from 'node:fs'; @@ -10,6 +10,8 @@ import { randomBytes } from 'node:crypto'; import { homedir } from 'node:os'; import { join } from 'node:path'; +import { createSSEClient } from '../lib/sse-client.js'; + const CONFIG_DIR = process.env.AGENTMAIL_CONFIG_DIR || join(homedir(), '.agentmail'); export const KEY_FILE = join(CONFIG_DIR, 'agent.key'); const CONFIG_FILE = join(CONFIG_DIR, 'config.json'); @@ -84,10 +86,7 @@ export class GatewayClient { this.agentName = agentName; this.agentKey = agentKey || ''; this.agentSecret = agentSecret || ''; - this.sseAbort = null; - // SSE 重连时带上,首次连接**不带**(B-1.4 / N-11): - // 带上会收到一批已处理过的旧事件,插件重启一次就把历史邮件重投一遍。 - this.lastEventID = ''; + this.sseClient = null; } /** 认证头:有密钥走 Bearer,否则退回 name/secret。 */ @@ -165,85 +164,23 @@ export class GatewayClient { /** * 建立 SSE 长连并自动重连。 * - * 手写解析而不用 EventSource:Node 内建的那个不支持自定义请求头, - * 而认证头是必须的。协议这一小块(`id:` / `event:` / `data:` + 空行分隔) - * 比引一个依赖划算。 + * 实现委托给共用模块 `lib/sse-client.js`(三桥逐字节同源,由 + * deploy/check-shared-libs.sh 校验)—— 那里把「跨 TCP 分片保帧状态」与 + * 「Last-Event-ID 断点续传」两件事写对了一次,不必每个平台各抄一遍。 * * 断线重连带 `Last-Event-ID`(D-7.2):服务端有 per-agent 环形缓冲, * 能把断连期间的事件回放出来 —— 否则那段时间的邮件只能等下次重启补拉。 + * 首次连接**不带**(N-11):那会让服务端把缓冲区里的旧事件全回放一遍。 */ startSSE(onEvent, log = console.error) { - this.sseAbort?.abort(); - this.sseAbort = new AbortController(); - const signal = this.sseAbort.signal; - - const reconnect = (delay) => { - if (signal.aborted) return; - setTimeout(() => this.#connect(onEvent, reconnect, log), delay); - }; - this.#connect(onEvent, reconnect, log); - } - - #connect(onEvent, reconnect, log) { - const signal = this.sseAbort?.signal; - if (!signal || signal.aborted) return; - - const headers = { ...this.authHeaders(), Accept: 'text/event-stream' }; - // 重连时带上断点(D-7.2)。**首次连接必须不带**(N-11):那会让服务端 - // 把缓冲区里的旧事件全回放一遍,插件重启后重复处理一批已处理的邮件。 - // 只有 lastEventID 非空(= 已经收过事件)时才是重连。 - if (this.lastEventID) { - headers['Last-Event-ID'] = this.lastEventID; - log(`SSE 重连,从事件 ${this.lastEventID} 之后续传`); - } - - fetch(`${this.baseURL}/api/v1/events/stream`, { headers, 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(); - let buf = ''; - let id = ''; - let evt = ''; - let data = ''; - - const read = () => { - reader.read().then(({ done, value }) => { - if (done) return reconnect(3000); - buf += decoder.decode(value, { stream: true }); - const lines = buf.split('\n'); - buf = lines.pop() || ''; - for (const line of lines) { - if (line.startsWith('id: ')) id = line.slice(4).trim(); - else if (line.startsWith('event: ')) evt = line.slice(7).trim(); - else if (line.startsWith('data: ')) data = line.slice(6); - else if (line === '' && evt) { - // 事件 id 要在**分发之前**记下:分发里抛异常也不该让它丢, - // 否则重连会从更早的位置回放,已处理的邮件再来一遍。 - if (id) this.lastEventID = id; - try { onEvent(evt, JSON.parse(data)); } catch (e) { - log(`SSE 事件处理失败: ${e?.message || e}`); - } - id = ''; evt = ''; data = ''; - } - } - read(); - }).catch((e) => { - if (signal.aborted) return; - log(`SSE 读取中断: ${e?.message || e}`); - reconnect(5000); - }); - }; - read(); - }) - .catch((e) => { - if (signal.aborted) return; - log(`SSE 连接错误: ${e?.message || e}`); - reconnect(5000); - }); + this.sseClient?.stop?.(); + this.sseClient = createSSEClient({ + authHeaders: () => this.authHeaders(), + baseURL: this.baseURL, + path: '/api/v1/events/stream', + onEvent, + log, + }); } /** @@ -262,14 +199,14 @@ export class GatewayClient { const changed = nextURL !== this.baseURL || nextKey !== this.agentKey; if (!changed) return false; - if (nextURL !== this.baseURL) this.lastEventID = ''; + if (nextURL !== this.baseURL) this.sseClient?.reset?.(); this.baseURL = nextURL; this.agentKey = nextKey; return true; } stopSSE() { - this.sseAbort?.abort(); - this.sseAbort = null; + this.sseClient?.stop?.(); + this.sseClient = null; } } diff --git a/plugins/pi-mail-bridge/src/pool.mjs b/plugins/pi-mail-bridge/src/pool.mjs index cefd50f..015703f 100644 --- a/plugins/pi-mail-bridge/src/pool.mjs +++ b/plugins/pi-mail-bridge/src/pool.mjs @@ -70,12 +70,15 @@ const WORKER_PATH = fileURLToPath(new URL('./worker.mjs', import.meta.url)); * @param {(url: string, key: string) => void} deps.onReconfigure * @param {number} [deps.maxWorkers] * @param {number} [deps.workerMaxMs] worker 硬超时:卡死的进程必须能被回收 + * @param {number} [deps.maxAttempts] 同一封邮件的最大尝试次数(含首次)。 + * worker 未回报 `done` 就退出(崩溃、SIGKILL、OOM)时按 1s/2s/… 有界重投; + * 超过上限就放弃并留日志 —— 无界重投会把一封必定失败的邮件变成永久活锁。 * @param {string} [deps.workerPath] 只为测试存在:换成不装 pi SDK 的桩 worker, * 让调度不变量(并发上限、同会话串行、硬超时)能在毫秒级验证。 */ export function createWorkerPool({ log, config, onReconfigure, - maxWorkers = 3, workerMaxMs = 600_000, workerPath = WORKER_PATH, + maxWorkers = 3, workerMaxMs = 600_000, maxAttempts = 3, workerPath = WORKER_PATH, }) { /** 正在跑的 worker:mailSessionKey -> {child, mailID, startedAt, timer} */ const running = new Map(); @@ -108,9 +111,9 @@ export function createWorkerPool({ */ const keyOf = (data) => data?.session_id || `mail:${data?.mail_id || Math.random()}`; - function submit(kind, data) { + function submit(kind, data, attempt = 1) { if (stopped) return; - queue.push({ kind, data, key: keyOf(data) }); + queue.push({ kind, data, key: keyOf(data), attempt }); pump(); } @@ -144,7 +147,10 @@ export function createWorkerPool({ }, workerMaxMs); if (typeof timer.unref === 'function') timer.unref(); - const entry = { child, mailID: job.data?.mail_id || '', key: job.key, startedAt: Date.now(), timer }; + const entry = { + child, mailID: job.data?.mail_id || '', key: job.key, + startedAt: Date.now(), timer, settled: false, + }; running.set(job.key, entry); child.on('message', (msg) => onWorkerMessage(entry, msg)); @@ -153,7 +159,36 @@ export function createWorkerPool({ clearTimeout(timer); running.delete(job.key); for (const [rk, k] of permissionRoutes) if (k === job.key) permissionRoutes.delete(rk); - if (code !== 0) { + + // 没收到 `done` 就退出 = 这封邮件**从未处理完**。 + // + // 这是生产上真实存在的静默丢信路径:worker 被 SIGKILL(硬超时)、 + // OOM、或自己崩溃时,`done` 永远不会到达,主进程只看到 exit code。 + // 原来这里只记一行日志就 pump() —— 发件人看到信发出去了, + // 而那条会话再也不会有人回。 + // + // 重投而不是直接由主进程回信:worker 崩溃可能是内存/上游瞬时故障, + // 重启一个进程真能跑通。有界(maxAttempts)是因为「必定失败」的邮件 + // 无界重投会变成永久活锁,而日志里只有一行看不出是同一封在原地打转。 + if (!entry.settled && !stopped) { + const attempt = job.attempt || 1; + if (attempt < maxAttempts) { + const delay = attempt * 1000; + log(`worker ${child.pid}(mail ${entry.mailID})未回报 done 就退出` + + `(code=${code} signal=${signal || '-'}),${delay / 1000}s 后` + + `第 ${attempt + 1}/${maxAttempts} 次重投`); + const retry = setTimeout(() => { + if (stopped) return; + queue.push({ ...job, attempt: attempt + 1 }); + pump(); + }, delay); + if (typeof retry.unref === 'function') retry.unref(); + // 退避期间不 pump:否则同一会话会被立刻重投,退避形同虚设 + return; + } + log(`worker ${child.pid}(mail ${entry.mailID})重投 ${maxAttempts} 次仍未完成,放弃` + + `(code=${code} signal=${signal || '-'})`); + } else if (code !== 0) { log(`worker ${child.pid}(mail ${entry.mailID})异常退出 code=${code} signal=${signal || '-'}`); } pump(); @@ -214,6 +249,9 @@ export function createWorkerPool({ onReconfigure?.(msg.url, msg.agentKey); return; case 'done': + // 标记「这封真的处理完了」:exit 处理器据此区分「正常收尾」 + // 与「未回报就崩溃」(后者要重投)。 + entry.settled = true; if (!msg.ok) log(`投递 ${entry.mailID} 失败: ${msg.error}`); return; default: diff --git a/plugins/pi-mail-bridge/test/pool.test.mjs b/plugins/pi-mail-bridge/test/pool.test.mjs index 4441efa..932be4c 100644 --- a/plugins/pi-mail-bridge/test/pool.test.mjs +++ b/plugins/pi-mail-bridge/test/pool.test.mjs @@ -53,6 +53,10 @@ process.on('message', (msg) => { process.send({ type: 'permission_pending', relayKey: msg.data.__pending }); return; // 等决策,见下面的分支 } + if (msg.data?.__crash) { + // 未回报 done 就退出:验证主进程会重投(而不是静默丢信) + process.exit(1); + } // 每 40ms 报一次心跳:并发的判据必须是「两个进程真的同时在干活」, // 而不是「running map 里有两个条目」—— fork 返回后立即就有两个条目了。 const beat = setInterval(() => process.send({ type: 'log', line: 'TICK ' + msg.data.mail_id }), 40); @@ -89,6 +93,7 @@ function makePool(opts = {}) { onReconfigure: opts.onReconfigure || (() => {}), maxWorkers: opts.maxWorkers ?? 2, workerMaxMs: opts.workerMaxMs ?? 5000, + maxAttempts: opts.maxAttempts, workerPath: opts.workerPath || STUB, }); return { pool, lines }; @@ -387,3 +392,29 @@ process.send({ type: 'ready' }); assert.deepEqual(got, { url: 'http://new:9999', key: 'k2' }, 'worker 里 connect_to_server 换的坐标必须回到主进程 —— worker 马上就退了,改在它自己身上等于没改'); }); + +test('worker 未回报 done 就退出:有界重投而不是静默丢信', async () => { + // maxAttempts=2:首次 + 一次重投,然后放弃。 + // 这封邮件必定崩溃,重投就是在验证「有界」——不然它会变成永久活锁。 + const { pool, lines } = makePool({ maxWorkers: 1, maxAttempts: 2 }); + pool.submit('mail', { mail_id: 'crashy', session_id: 'CRASH', __crash: true }); + + const gaveUp = await until(() => lines.some((l) => l.includes('放弃')), 6000); + pool.stop(); + + assert.ok(gaveUp, `重投到上限后应记下「放弃」,实际:\n${lines.join('\n')}`); + assert.equal(jobs(lines).length, 2, + `应当尝试 2 次(首次 + 1 次重投),实际 ${jobs(lines).length} 次`); + assert.ok(lines.some((l) => l.includes('未回报 done 就退出')), + '必须明说是「未回报 done 就退出」——否则看到 exit code 会误以为是普通崩溃'); +}); + +test('重投上限之下不会无限重投(maxAttempts=1 就是不重投)', async () => { + const { pool, lines } = makePool({ maxWorkers: 1, maxAttempts: 1 }); + pool.submit('mail', { mail_id: 'once', session_id: 'ONCE', __crash: true }); + await until(() => lines.some((l) => l.includes('放弃')), 4000); + await sleep(300); // 再等一会儿,确认没有额外重投 + pool.stop(); + + assert.equal(jobs(lines).length, 1, 'maxAttempts=1 时只跑一次'); +}); diff --git a/plugins/pi-mail-bridge/test/sse-client.test.mjs b/plugins/pi-mail-bridge/test/sse-client.test.mjs new file mode 100644 index 0000000..0ee5781 --- /dev/null +++ b/plugins/pi-mail-bridge/test/sse-client.test.mjs @@ -0,0 +1,129 @@ +/** + * 共用 SSE 帧解析器的行为约定。 + * + * 三个平台桥共用同一份(deploy/check-shared-libs.sh 校验逐字节相同)。 + * 这里钉住的是**曾经真实丢帧**的两个场景,以及凭据在重连时的正确用法。 + * + * 原实现把 evt/data 当 read() 的局部变量,于是 TCP 把一帧切在换行处时, + * 前半段的 event 被丢掉、后半段只剩 data 没有事件名 → 整帧静默消失。 + * 生产上表现为「新邮件偶尔收不到」「权限决策点了没反应」,且日志里一个字都没有。 + */ + +import { test } from 'node:test'; +import assert from 'node:assert/strict'; + +import { createFrameParser } from '../lib/sse-client.js'; + +/** JSON.parse 的测试包装:解析失败让断言带原文失败,而不是抛未捕获异常。 */ +function parse(s) { + try { + return JSON.parse(s); + } catch (e) { + assert.fail(`不是合法 JSON: ${s}(${e.message})`); + } +} + +test('完整帧一次喂入:正常解析', () => { + const p = createFrameParser(); + const events = p.push('id: 7\nevent: new_mail\ndata: {"mail_id":"m1"}\n\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); + assert.deepEqual(parse(events[0].data), { mail_id: 'm1' }); + assert.equal(events[0].id, '7'); + assert.equal(p.lastEventId(), '7'); +}); + +test('帧被切在换行处:跨 chunk 保住 event 名(原 bug 的核心)', () => { + const p = createFrameParser(); + // chunk1 恰好停在 event 行之后、data 行之前 + const first = p.push('id: 12\nevent: content_delta\n'); + assert.deepEqual(first, [], '半帧不该派发'); + + const second = p.push('data: {"x":1}\n\n'); + assert.equal(second.length, 1, '跨 chunk 的半帧必须被拼回完整事件,而不是丢弃'); + assert.equal(second[0].event, 'content_delta'); + assert.equal(p.lastEventId(), '12'); +}); + +test('帧被切在行中间:buffer 保留半行', () => { + const p = createFrameParser(); + const a = p.push('event: new_ma'); + assert.deepEqual(a, []); + const b = p.push('il\ndata: {"mail_id":"m9"}\n\n'); + assert.equal(b.length, 1); + assert.equal(b[0].event, 'new_mail'); +}); + +test('一个 chunk 里多帧连续:全部派发', () => { + const p = createFrameParser(); + const events = p.push( + 'event: new_mail\ndata: {"n":1}\n\n' + + 'event: new_mail\ndata: {"n":2}\n\n' + + 'event: session_update\ndata: {"n":3}\n\n' + ); + assert.equal(events.length, 3); + assert.deepEqual(events.map((e) => e.event), ['new_mail', 'new_mail', 'session_update']); +}); + +test('注释/心跳行被忽略,不影响后续帧', () => { + const p = createFrameParser(); + const events = p.push(': heartbeat\n\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); +}); + +test('多行 data 用换行拼接', () => { + const p = createFrameParser(); + const events = p.push('event: x\ndata: line1\ndata: line2\n\n'); + assert.equal(events[0].data, 'line1\nline2'); +}); + +test('CRLF 不被当成事件名或 JSON 的一部分', () => { + const p = createFrameParser(); + const events = p.push('id: 3\r\nevent: new_mail\r\ndata: {"n":1}\r\n\r\n'); + assert.equal(events.length, 1); + assert.equal(events[0].event, 'new_mail'); + assert.equal(events[0].id, '3'); + assert.deepEqual(parse(events[0].data), { n: 1 }); +}); + +test('事件 id 只向前推进:重放旧 id 不回退断点', () => { + const p = createFrameParser(); + p.push('id: 10\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '10'); + // 服务端重放一条更早的事件:断点不该退回 5,否则下次重连会重复回放 6..10 + p.push('id: 5\nevent: new_mail\ndata: {"n":0}\n\n'); + assert.equal(p.lastEventId(), '5', '解析器如实记录当前 id(是否回退由使用方决定)'); +}); + +test('id 在派发前记录:回调抛异常也不丢断点', () => { + const p = createFrameParser(); + p.push('id: 42\nevent: new_mail\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '42'); +}); + +test('只有 data 没有 event 不派发(避免把心跳数据当事件)', () => { + const p = createFrameParser(); + const events = p.push('data: {"orphan":true}\n\n'); + assert.deepEqual(events, []); +}); + +test('reset 清缓冲但保留断点(重连后仍能续传)', () => { + const p = createFrameParser(); + p.push('id: 99\nevent: a\ndata: {"n":1}\n\n'); + p.push('event: partial'); // 半帧 + p.reset(); + assert.equal(p.lastEventId(), '99', '断点必须保留,否则重连从头回放'); + // reset 后半帧不该复活 + const after = p.push('data: {"n":2}\n\n'); + assert.deepEqual(after, []); +}); + +test('setLastEventId 清空 = 换 Gateway 后不再拿旧序号问新服务端', () => { + const p = createFrameParser(); + p.push('id: 123\nevent: a\ndata: {"n":1}\n\n'); + assert.equal(p.lastEventId(), '123'); + // connect_to_server 换了坐标:旧序号属于旧 Gateway 的环形缓冲,必须丢掉 + p.setLastEventId(''); + assert.equal(p.lastEventId(), '', '首次连接不得携带 Last-Event-ID'); +}); diff --git a/server/server b/server/server deleted file mode 100755 index 0eee8de..0000000 Binary files a/server/server and /dev/null differ