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 全量通过;四平台重启后在线且无错误。
This commit is contained in:
2026-09-11 10:27:41 +08:00
parent f9d757b5e5
commit c401eb2da2
17 changed files with 1093 additions and 175 deletions

3
.gitignore vendored
View File

@ -59,3 +59,6 @@ attachments/
.idea/
.vscode/
plugins/homeagent-mail-bridge/build/
# 构建产物(曾误提交)
server/server

View File

@ -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

View File

@ -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

View File

@ -42,7 +42,13 @@ RestartSec=10
ExecStartPre=-/bin/rm -f /root/.agentmail-pi/pi-bridge.lock
# 桥是长驻守护进程,启动即注册 + 立刻打一次心跳 + 订阅 SSEB-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 会把它拉回来。

View File

@ -0,0 +1,17 @@
export interface SseClientOptions {
authHeaders: () => Record<string, string>;
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;
};

View File

@ -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<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(''),
};
}

View File

@ -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';
// 只用 isApprovalDSH 没有 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) {

View File

@ -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');
});

View File

@ -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 {

View File

@ -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<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(''),
};
}

View File

@ -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');
});

View File

@ -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<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(''),
};
}

View File

@ -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 长连并自动重连。
*
* 手写解析而不用 EventSourceNode 内建的那个不支持自定义请求头,
* 而认证头是必须的。协议这一小块(`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;
}
}

View File

@ -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,
}) {
/** 正在跑的 workermailSessionKey -> {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:

View File

@ -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 时只跑一次');
});

View File

@ -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');
});

Binary file not shown.