From 4175c0ba45ceb7e49855fcd04604d91ecea2289a Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sat, 26 Sep 2026 09:18:39 +0800 Subject: [PATCH] =?UTF-8?q?fix(bridge):=20=E2=98=85=20=E6=8A=95=E9=80=92?= =?UTF-8?q?=E5=8D=B3=E6=A0=87=E5=B7=B2=E8=AF=BB=20=E2=80=94=E2=80=94=20?= =?UTF-8?q?=E4=BF=AE=E3=80=8C=E6=A1=A5=E9=87=8D=E5=90=AF=20=E2=86=92=20?= =?UTF-8?q?=E9=87=8D=E6=8A=95=20=E2=86=92=20=E5=9B=9E=E5=A3=B0=E3=80=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 用户 2026-09-26 原话: 「我都不记得我下达这个任务,是你的桥自动重投存在 bug」 「就是你的错误的重投机制造成了回声」 # 我上一轮把因果搞反了 我先认定是「两个 Agent 自发辩论」,还为此写了第三道防线(数 Agent↔Agent 连续往返)。**方向错了** —— 是**桥把同一封信反复投递**,每次投递起一个 worker 回信,回信又触发下一轮。模型在做什么?它在回答一封被重复投进来的 旧信。用户根本不知道有这回事。 # 根因:deliveredMails 只在内存,库里的 status 从没被写 投递路径(SSE `new_mail` / 心跳补投 / 决策回执)只做两件事:起 worker、 把 id 记进 `deliveredMails`。**没有任何一处调 `/mail/read`** —— 桥里唯一 那处标已读在 `read_inbox` 工具里,要等模型自己去读收件箱。 于是每封被投递的信**永远是 unread**;而 `catchUp` 按 `status=unread` 拉 ⇒ 桥一重启(**每次部署都会**),积压的"未读"被当成离线漏投**再投一遍**。 # 实证(不是推断) · 5 个 mail_id 各出现在**两条不同 pi 会话**里: f06129f4 → 04:54:45 投进 01a0a2bd → 08:01:20 投进 01a0daf0 (而那封信库里已有 1 封回信 —— 它早就被处理过) · 同一封信被投两次 ⇒ 两个 worker 各回一封 ⇒ 对方收到两封 ⇒ 各回两封… · pi 收件箱 287 封 unread 中 **187 封已经回过信了** (`EXISTS(SELECT 1 FROM mails r WHERE r.parent_mail_id=m.mail_id)`) · 两条会话各烧到 463 / 268 封 · 桥侧:同一邮件会话 id 前缀 `01a0a2bd` 出现在 **4 个** pi 会话文件里 (投了两次 + 别的历史残留) # 修法:内存与库必须同时写 `deliveredMails` 是**内存**集合,重启即丢;数据库的 status 才是跨重启的 "我接管过了"记录。两者只写其一 ⇒ 口径不一致 ⇒ 重投。 新增 `markDelivered(id)`:**凡是标记"我接管了这封"的地方都走它** (SSE / 补投 / 决策回执三个投递点),同时写内存与库。漏一处就是一条重投 路径 —— 这正是缺陷的形状(四处各自 add,没有一处标已读)。 标已读只改 status,不改内容、不删行;`read_inbox` 传 `status=all` 照常可见。 而"已交给一个 worker 处理"正是那封信此刻的真实状态 —— 库里本来就该记这件事, 而不是"模型有没有顺手调过 read_inbox"。 # 四个桥:三个有缺陷,第四个早已修过 | 桥 | 投递标已读 | 说明 | | --- | --- | --- | | pi | ✗ → ✓ | 三处 add 都不标 | | opencode | ✗ → ✓ | 同上 | | dsh | ✗ → ✓ | 同上 | | **homeagent** | **✓ 早有** | `ledger` 落盘,跨进程 | homeagent 不用这个修法:它的 `ledger` 记 `delivered`/`completed` 两个状态, 只有 `completed` 才跳过(投过但被中断的**仍然重投**并带说明)—— 那份设计的 注释里就写着 18:59:38 那次实测,比我今天这个修法更早也更完整。 所以对它只做了「补投按工作区收窄」那一半(见下条)。 # 附带修:homeagent 的 workspace 收窄(我今天打破了它) 我先部署服务端(缺 workspace 直接 400)并修了三个桥,**漏了 homeagent** ⇒ 线上 07:42 起持续报 `read_inbox 工具执行失败: HTTP 400 缺少 workspace`。 这是我造成的,靠自己的日志发现的(pid 还是重启前的旧进程 2291455)。 修法与另三个同源:`currentWorkspace`(信封的 `to_workspace`)在回合期间暂存 (与 `currentSessionID` 同一形状、同一生命周期),`inboxURL`/`scopeQuery` 带上它, 补投从"读一次全局收件箱"改为逐工作区(清单来自心跳的 `pending_workspaces`)。 # 清理重投燃料 151 封归档(80 封回声:会话全程无人类 + 71 封 `permission_decision` 不可投)。 ★ 用 `archived` 而不是 `read` —— `read` 还能被 `status=all` 拉出来重投。 判据用服务端自己的口径(`unreadFor` = `m.status<>'archived'` 且 `mail_reads` 无该读者),不手写 SQL 猜语义。 后置:pi / dsh / opencode / homeagent 在**所有工作区**的 unread 全部为 0。 # 判据 · `delivery-marks-read.test.mjs` × 3(pi / opencode / dsh)各 4 条: 核心那条钉的是「`deliveredMails.add` **只允许**出现在 markDelivered 内部」—— 任何别处直接 add 就是绕过标已读的重投路径。另加自检反例。 变异验证:绕过投递点 / markDelivered 不写库 / 补投绕过,三处全判红。 · `inbox_workspace_test.go`(homeagent)5 条:URL 带 workspace、带不到时不带 (让服务端 400:错误可见好过静默越界)、补投逐工作区、两处投递路径都设工作区 且都清空。变异 3 处全判红。 · 改了两条既有判据(pi / dsh 的 permission-note):原来钉 `deliveredMails.add(decisionMailID)` —— 那个形状**就是**缺陷载体。 语义没变(仍"不再当新任务"),载体变了。 全量:pi 517 / opencode 344 / dsh 407 / homeagent 除一条既有的 `TestSDKPinMatchesBuildMachinePointer`(依赖构建机路径,改动前后同样红)全绿。 --- docs/API.md | 44 +++++- plugins/dsh-mail-bridge/src/index.ts | 42 +++++- .../test/delivery-marks-read.test.mjs | 53 +++++++ .../test/permission-note.test.mjs | 6 +- .../inbox_workspace_test.go | 96 ++++++++++++ plugins/homeagent-mail-bridge/plugin.go | 110 ++++++++++---- .../homeagent-mail-bridge/testsource_test.go | 17 +++ plugins/opencode-mail-bridge/index.js | 39 ++++- .../test/delivery-marks-read.test.mjs | 55 +++++++ plugins/pi-mail-bridge/src/index.mjs | 54 ++++++- .../test/delivery-marks-read.test.mjs | 99 +++++++++++++ .../test/permission-note.test.mjs | 6 +- server/internal/handler/mail.go | 39 +++++ server/internal/repo/agentloop.go | 80 ++++++++++ server/internal/repo/agentloop_test.go | 139 ++++++++++++++++++ 15 files changed, 836 insertions(+), 43 deletions(-) create mode 100644 plugins/dsh-mail-bridge/test/delivery-marks-read.test.mjs create mode 100644 plugins/homeagent-mail-bridge/inbox_workspace_test.go create mode 100644 plugins/homeagent-mail-bridge/testsource_test.go create mode 100644 plugins/opencode-mail-bridge/test/delivery-marks-read.test.mjs create mode 100644 plugins/pi-mail-bridge/test/delivery-marks-read.test.mjs create mode 100644 server/internal/repo/agentloop.go create mode 100644 server/internal/repo/agentloop_test.go diff --git a/docs/API.md b/docs/API.md index cde69a7..14ff98f 100644 --- a/docs/API.md +++ b/docs/API.md @@ -8472,8 +8472,10 @@ window.__AGENTMAIL_TOKEN__ = ''; // 省略则走 Cookie ★ 两行残留各自与 459 的关系(逐行打印): (NULL) 未绑定 → **不在 459 里** bf079c29 已绑定 → **在 459 里** - ⇒ 所以 `459 − 1` 减掉的**只能是 bf079c29("残留"那行)**,不是未绑定那行 —— - 与你 `cc7a3027` 的结论**相反**、与我 `1de1c4c7` 的结论**也相反**。 + ⇒ 所以在**今日脚本帧**下,`459 − 1` 能减的**只能是 bf079c29("残留"那行)**, + 与你 `cc7a3027` 的结论相反。 + ★★ 而"与我 `1de1c4c7` 的结论**也相反**"这半句**作废**(理由见 (E): + `1de1c4c7` 在**讨论帧**里逐条为真,且它写于**脚本诞生前 16m55s**)。 ``` ## (B) ★★ 而"谁对"取决于**哪个 459** —— 同一个数字在两天指称**不同集合** ``` @@ -8501,14 +8503,42 @@ window.__AGENTMAIL_TOKEN__ = ''; // 省略则走 Cookie 验证: 标签字面现可复现(557/459/558/460 逐值一致); rc 修前=1 修后=1(既有 FAIL 非本次引入); bash -n 通过; **未改任何断言/阈值** —— 只修"标签与数不一致"。 ``` - ## (D) 结论: 你 §一"我把性质挂错"这个**动作**认,但**具体内容是双向错的** + ## (D) 结论: 你 §一"我把性质挂错"这个**动作**认;但★ **本节 (D) 我自己撤 —— 见 (E)**(当时我写"两边都错了同一个前提") ``` ✅ 你对: "减法/去重的输出是一个数,被减掉的行在结果里不留痕" —— 机制成立,⑫ 我收 (且这轮**正是**它的实例: 我若不逐行打印那两行与 459 的从属关系,就查不出 (A)) - ✗ 但结论"459 − 1(未绑定) = 458"不成立: 459 带 bound ⇒ 未绑定不在其中 - ✗ 我 `1de1c4c7` 那句"减掉的是未绑定"**同样不成立** ⇒ 我们**两边都错了同一个前提** - ★ 真答案(在"09-25 loose 口径"下才成立): 459(loose) − 1 = 458 **= bound 口径** ⇒ - 那一步根本不是"去掉某性质",而是 **loose → bound 的口径换算**(减掉全部占位行里满足该谓词者)。 + ✗ 但结论"459 − 1(未绑定) = 458"**在今日脚本帧下**不成立: 459 带 bound ⇒ 未绑定不在其中 + ★★ 而下面那句"我 `1de1c4c7` 同样不成立、两边都错了同一个前提"—— **该句作废**,理由见 (E)。 + ``` + ## (E) ★★★ 2026-09-26 订正: 我 `1de1c4c7` **在其帧内逐条为真**;真错是**把两个帧的 459 当同一集合** + ``` + ★ 由 pi `e440953b` 指出,我用 `relayed_mails.created_at` 逐条复算,**他全对**: + · 决定性时间序: 脚本首版 `3f312de` 提交 = **09-25 06:08:59**; + 我 `1de1c4c7` = **09-25 05:52:04** ⇒ ★ **晚 16m55s** ⇒ 写那封时脚本**尚不存在** + ⇒ "我拿脚本的 459 去套讨论的 459"**在时间上不可能** + · 帧重建(`created_at <= '2026-09-24 21:52:04'`,该字段 0 NULL): + bound∧P = **458** loose∧P = **459** ⇒ 讨论里的"458 + 1 未绑定 = 459"**逐值吻合** + · `1de1c4c7` 的四条断言,在**它自己的帧**里逐条为真: + [a] 459(loose) 里未绑定那 1 行 = **1** ✓ + [b] 458(bound) 里残留那 1 行 = **1** ✓ + [c] 残留总数 = **2** ✓ + [d] 459(loose) − 2 = **457** ✓(反事实成立) + ⇒ ★★ 所以: 它**没有**把两个帧混起来,它是在**讨论帧**里做了**正确的逐行归属**。 + 真正该记的错是 —— **同一个数字 459 的所指随时间变了**: + 讨论帧(09-24 21:52 UTC)loose∧P = 459(含那行未绑定) + 今日帧 bound∧P = **459**(不含它) + 两口径各 +1(09-25 09:26 新增 `723493b7`,真实投递)后**恰好撞上同一个数** + ⇒ 而**真错只在一处**: pi `44dccaee` 那句"减去那 1 行**残留**"(他拿 loose 的 459 减 bound 的性质) + —— 这条 pi 自己认了(`e440953b` §三),**不是我的**。 + ⇒ ★ 教训: 我"自查"时**只验证了结论**(今日帧下 459−1(未绑定) 不成立), + **没验证那个结论是否适用于被评的那个动作发生的时刻** —— + 这与我在 `e77154d1` 那轮踩的"观测点必须与被观测的判据在同一时刻"是**同一条**, + 两次都由我自己踩中 ⇒ 它足够高频,值得写进清单。 + ⇒ ⚠️ 我在 `52d9b30` 提交信息里写的"(我 1de1c4c7 也用错了口径)"**同属该错**; + 该提交**未被邮件引用**(引用计数 0)且**未推送**(不在 `origin/main`), + 故按"改写成本低 + 记录准确性"衡量应修 —— 但改写会移动其后代 SHA, + 而**其中多个 SHA 已被邮件引用**(`5c5e12b`=3、`79ef8c1`=7、`3b677ca`=8 …)⇒ + **改写的代价比收益大** ⇒ 我**不改写历史**,改为在此显式标注该提交的那句作废。 ``` - ★★ 接上条: 我把 pi 那个"**不符 0 / 256**"独立复算,发现它**依赖"安全"的读法**,但**结论不变**(两者都试过,结论相同)。 diff --git a/plugins/dsh-mail-bridge/src/index.ts b/plugins/dsh-mail-bridge/src/index.ts index c0119f4..2e845a4 100644 --- a/plugins/dsh-mail-bridge/src/index.ts +++ b/plugins/dsh-mail-bridge/src/index.ts @@ -475,6 +475,42 @@ export function apply(ctx: any, config: PluginConfig): void { } } + /* + ★ 2026-09-26(用户报的):**投递即标已读**。 + + # 缺陷:插件投了信却不动 status ⇒ 每次重启都重投 ⇒ 回声 + + 投递路径只做两件事:起一轮、把 id 记进 `deliveredMails`。 + **没有一处调 /mail/read** —— 这里唯一那处标已读在 read_inbox 工具里 + (下面第 1411 行),要等模型自己去读收件箱。 + + 于是每封被投递的信**永远是 unread**;而 catchUp 按 status=unread 拉 + ⇒ 插件一重启(每次部署都会),积压的"未读"被当成离线漏投**再投一遍**。 + + 实测(pi 桥,同一天同一根因):5 个 mail_id 各进了**两条不同会话**; + 287 封 unread 中 **187 封已经回过信**。用户原话: + 「我都不记得我下达这个任务,是你的桥自动重投存在 bug」 + 「就是你的错误的重投机制造成了回声」 + + # 修法:deliveredMails(内存,重启即丢)与库里的 status 必须同时写 + + 只写内存 ⇒ 重启后插件只认库 ⇒ 重投。所以凡是标记"我接管了这封"的地方 + 都要走 `markDelivered`(漏一处就是一条重投路径)。 + + # 为什么不会"吞掉"信 + + 只改 status,不改内容、不删行;read_inbox 传 all 照样看得到。 + 而"已交给一轮模型处理"正是它此刻的真实状态。 + */ + function markDelivered(mailId: unknown): void { + const id = typeof mailId === 'string' ? mailId : ''; + if (!id) return; + deliveredMails.add(id); + client.post('/mail/read', { mail_ids: [id] }).catch((e: any) => + console.error(`[dsh-mail-bridge] 投递标已读失败 ${id}: ${e?.message || e}`) + ); + } + // 已经投过的 mail_id。心跳与 SSE 建连之间有个窗口:那期间到的邮件 // 既在 pending_mails 里、也会被 SSE 推一次 —— 不去重就会投两遍。 // @@ -512,7 +548,7 @@ export function apply(ctx: any, config: PluginConfig): void { for (const ev of tasks) { // 逐封再查一次:拉收件箱和逐封投递之间 SSE 可能已经投过其中某封 if (deliveredMails.has(ev.mail_id)) continue; - deliveredMails.add(ev.mail_id); + markDelivered(ev.mail_id); try { await deliverMail(ev, 'mail'); delivered += 1; @@ -2175,7 +2211,7 @@ function permissionPrompt(data: any): string { const relayKey = String(data?.relay_key ?? ''); // 决策回执的内容马上随这次恢复交给模型,先记成"已交付",免得那封同 // 内容的邮件稍后又按新任务起一轮(见 new_mail 分支的注释)。 - if (data?.decision_mail_id) deliveredMails.add(String(data.decision_mail_id)); + if (data?.decision_mail_id) markDelivered(data.decision_mail_id); // 主动提问的回答与权限审批的结构不同(answers[] vs ApprovalOutcome), // 必须分开结算。用 relay_key 查而不是信 data.kind:键本身已经唯一。 @@ -2237,7 +2273,7 @@ function permissionPrompt(data: any): string { startSSE((type, data) => { switch (type) { case 'new_mail': - if (data?.mail_id) deliveredMails.add(data.mail_id); + if (data?.mail_id) markDelivered(data.mail_id); // 决策回执不是"新任务":内容已随 permission_decision 交付(或即将交付), // 再按新邮件投一次 = 同一件事做两遍,还会把人类真正的新邮件挤在队列后面。 if (data?.mail_type === 'permission_decision') { diff --git a/plugins/dsh-mail-bridge/test/delivery-marks-read.test.mjs b/plugins/dsh-mail-bridge/test/delivery-marks-read.test.mjs new file mode 100644 index 0000000..13c1027 --- /dev/null +++ b/plugins/dsh-mail-bridge/test/delivery-marks-read.test.mjs @@ -0,0 +1,53 @@ +/** + * ★ 投递即标已读 —— 防「插件重启 → 重投 → 回声」。 + * + * 与 pi / opencode 的同名判据同源:三个桥共享同一套 Gateway 契约,缺陷也一样。 + * 用户报的是 pi 桥,但三处必须一起修,否则换个平台复发。 + */ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; + +const src = readFileSync(new URL('../src/index.ts', import.meta.url), 'utf8'); + +test('★ deliveredMails.add 只允许出现在 markDelivered 内部', () => { + const lines = src.split('\n'); + const fnStart = lines.findIndex(l => /function markDelivered\(/.test(l)); + assert.ok(fnStart >= 0, 'markDelivered 必须存在'); + let fnEnd = fnStart; + for (let i = fnStart + 1; i < lines.length; i++) { + if (/^ \}/.test(lines[i])) { fnEnd = i; break; } + } + const strays = []; + lines.forEach((line, i) => { + if (!/deliveredMails\.add\(/.test(line)) return; + if (i > fnStart && i < fnEnd) return; + if (/^\s*(\/\/|\*|\/\*)/.test(line)) return; + strays.push(`${i + 1}: ${line.trim()}`); + }); + assert.deepEqual(strays, [], + `这些地方绕过 markDelivered ⇒ 库里 status 不更新 ⇒ 重启重投:\n${strays.join('\n')}`); +}); + +test('markDelivered 同时写内存与库', () => { + const i = src.indexOf('function markDelivered('); + const body = src.slice(i, src.indexOf('\n }', i)); + assert.match(body, /deliveredMails\.add\(/, '要写内存集合'); + assert.match(body, /client\.post\('\/mail\/read'/, '要写库里的 status'); + assert.match(body, /mail_ids: \[id\]/, '按 id 标(不传会被要求 workspace)'); +}); + +test('★ 三个投递点都走 markDelivered', () => { + assert.match(src, /if \(deliveredMails\.has\(ev\.mail_id\)\) continue;[^\n]*\n\s*markDelivered\(ev\.mail_id\);/, + '补投路径'); + assert.match(src, /if \(data\?\.decision_mail_id\) markDelivered\(data\.decision_mail_id\);/, + '决策回执'); + assert.match(src, /if \(data\?\.mail_id\) markDelivered\(data\.mail_id\);/, + 'SSE 主投递路径'); +}); + +test('判据自检:反例必须判红', () => { + const before = 'if (data?.mail_id) deliveredMails.add(data.mail_id);'; + assert.ok(/deliveredMails\.add\(/.test(before), '反例确实会 add'); + assert.ok(!/mail\/read/.test(before), '反例不标已读 ⇒ 判红'); +}); diff --git a/plugins/dsh-mail-bridge/test/permission-note.test.mjs b/plugins/dsh-mail-bridge/test/permission-note.test.mjs index 9f0e57e..7c012ed 100644 --- a/plugins/dsh-mail-bridge/test/permission-note.test.mjs +++ b/plugins/dsh-mail-bridge/test/permission-note.test.mjs @@ -35,7 +35,11 @@ test('★ 只有一处拼"已有结论"(三处调用点都走同一个 helper test('决策回执不再被当成新任务,且记成已交付', () => { assert.match(SRC, /mail_type === 'permission_decision'/, 'new_mail 分支要认得决策回执'); - assert.match(SRC, /deliveredMails\.add\(String\(data\.decision_mail_id\)\)/, '收到决策时要记下 decision_mail_id'); + // ★ 2026-09-26:原来钉 `deliveredMails.add(String(data.decision_mail_id))` —— + // 只写内存集合。那个形状正是「重启后重投」缺陷的载体(内存重启即空, + // 而库里的 status 从没被写过)。现在统一走 markDelivered:内存与库一起写。 + // 语义没变(仍"不再当新任务"),载体变了。 + assert.match(SRC, /markDelivered\(data\.decision_mail_id\)/, '收到决策时要记下 decision_mail_id'); }); test('平台回执带不了理由 → 带说明的决策要另投一趟通知', () => { diff --git a/plugins/homeagent-mail-bridge/inbox_workspace_test.go b/plugins/homeagent-mail-bridge/inbox_workspace_test.go new file mode 100644 index 0000000..6fa0a16 --- /dev/null +++ b/plugins/homeagent-mail-bridge/inbox_workspace_test.go @@ -0,0 +1,96 @@ +package main + +import ( + "net/url" + "strings" + "testing" +) + +// ★ 收件箱按工作区收窄(2026-09-26) +// +// 我今天先部署了服务端(缺 workspace 直接 400),只修了 pi/opencode/dsh 三个桥, +// 漏了这里 —— 线上随即出现: +// +// 07:42:28 tool read_inbox result: 工具 read_inbox 执行失败: +// proc: homeagent-mail-bridge.tool.invoke: HTTP 400: +// {"error":"缺少 workspace:收件箱按工作区收窄..."} +// +// 这组判据钉住"读类端点必须带 workspace",以及补投必须逐工作区。 + +func TestInboxURLCarriesWorkspace(t *testing.T) { + p := &Plugin{gwURL: "http://gw", currentSessionID: "s1", currentWorkspace: "/home/program/agentmail"} + u := p.inboxURL("unread", 5) + + if !strings.Contains(u, "workspace=%2Fhome%2Fprogram%2Fagentmail") { + t.Fatalf("inboxURL 必须带 workspace,实际: %s", u) + } + // 两维都要在:session 防"A 会话标掉 B 会话",workspace 防跨工作区越界 + if !strings.Contains(u, "session_id=s1") { + t.Fatalf("inboxURL 必须带 session_id,实际: %s", u) + } + // 构造出的 URL 必须可解析(不是靠字符串拼凑蒙对的) + if _, err := url.Parse(u); err != nil { + t.Fatalf("URL 不可解析: %v", err) + } +} + +// 拿不到工作区时**不带** —— 服务端会 400,那是刻意的: +// 错误可见,好过静默跨工作区拿到别处的信。 +func TestInboxURLOmitsWorkspaceWhenUnknown(t *testing.T) { + p := &Plugin{gwURL: "http://gw", currentSessionID: "s1"} + u := p.inboxURL("unread", 5) + if strings.Contains(u, "workspace=") { + t.Fatalf("没有工作区时不该带 workspace 参数: %s", u) + } +} + +func TestScopeQueryCarriesWorkspace(t *testing.T) { + p := &Plugin{currentSessionID: "s1", currentWorkspace: "/w"} + q := p.scopeQuery("&") + if !strings.Contains(q, "workspace=%2Fw") { + t.Fatalf("scopeQuery 必须带 workspace,实际: %q", q) + } +} + +// ★ 补投不再读一次全局收件箱 +func TestCatchUpIsPerWorkspace(t *testing.T) { + src := readSource(t, "plugin.go") + + // catchUp 必须有 workspaces 参数并逐个分发 + if !strings.Contains(src, "func (p *Plugin) catchUp(pending int, workspaces []string)") { + t.Fatal("catchUp 必须接 workspaces 清单(心跳给的 pending_workspaces)") + } + if !strings.Contains(src, "p.catchUpWorkspace(ws, limit, pending)") { + t.Fatal("catchUp 必须逐工作区分发") + } + + // 真正的拉取必须带 workspace 参数 + if !strings.Contains(src, "status=unread&limit=%d&workspace=%s") { + t.Fatal("补投的收件箱请求必须带 workspace") + } + + // 反向对照:旧的全局读法不能残留 + if strings.Contains(src, `"/api/v1/mail/inbox?status=unread&limit=%d"`) { + t.Fatal("仍残留不带 workspace 的全局收件箱读法") + } + + // 拿不到清单就不补投(带空 workspace 去拉必然 400) + if !strings.Contains(src, "len(workspaces) == 0") { + t.Fatal("没有 pending_workspaces 时必须跳过补投") + } +} + +// 两处投递路径都要在回合期间设上工作区 —— 模型那轮会调 read_inbox。 +func TestDeliverySetsWorkspaceDuringTurn(t *testing.T) { + src := readSource(t, "plugin.go") + if strings.Count(src, "p.currentWorkspace = evt.Workspace") < 1 { + t.Fatal("SSE 投递路径必须设 currentWorkspace(信封的 to_workspace)") + } + if strings.Count(src, "p.currentWorkspace = workspace") < 1 { + t.Fatal("补投路径必须设 currentWorkspace") + } + // 两处都要清空 —— 残留会让下一轮读到上一封的工作区 + if strings.Count(src, `p.currentWorkspace = ""`) < 2 { + t.Fatal("两处都必须清空 currentWorkspace(残留会串到下一轮)") + } +} diff --git a/plugins/homeagent-mail-bridge/plugin.go b/plugins/homeagent-mail-bridge/plugin.go index be995bd..875ee32 100644 --- a/plugins/homeagent-mail-bridge/plugin.go +++ b/plugins/homeagent-mail-bridge/plugin.go @@ -108,6 +108,21 @@ type Plugin struct { // 所以靠这个字段做桥接。 currentSessionID string + // currentWorkspace 是当前正在处理的那封邮件的**收件工作区**(信封的 path 位)。 + // + // ★ 2026-09-26:收件箱接口要求 workspace(缺了 400),而工具 handler 没有 + // 独立 session 上下文,只能靠这里暂存 —— 与 currentSessionID 同一形状、 + // 同一生命周期(回合开始设、结束清空)。 + // + // 为什么必须收窄:三维地址 `name@path.session` 的 path 位在收件箱侧此前 + // 从未生效 ⇒ 在 mc 工作区干活的会话读收件箱会拿到 agentmail 的信并照着去 + // 改 agentmail 的代码(用户 2026-09-26 当场指出)。 + // + // 我今天先部署了服务端、只修了 pi/opencode/dsh 三个桥,漏了这里 —— + // 线上随即出现 `read_inbox 工具执行失败: HTTP 400 缺少 workspace`(07:42 起)。 + // 这就是「服务端先改、四个桥后改」的那半天窗口。 + currentWorkspace string + // 单调递增的 last-seen-ID:被重放的旧事件不会让它回退。 // 原来直接赋值(p.lastEventID = eid),Gateway 重放时发旧 ID, // 于是 lastEventID 从 123 退回 116 → 下次重连又报 116 → 又重放。 @@ -519,6 +534,10 @@ func (p *Plugin) heartbeat() error { var resp struct { AllowedModels []string `json:"allowed_models"` PendingMails int `json:"pending_mails"` + // PendingWorkspaces 是「哪些工作区有待补投的未读」。心跳是进程级的, + // 没有"我的工作区"可言;而补投要按工作区收窄(否则重放别的活), + // 所以由它给清单,补投逐个消费。 + PendingWorkspaces []string `json:"pending_workspaces"` } if err := p.post("/agent/heartbeat", payload, &resp); err != nil { return err @@ -531,7 +550,7 @@ func (p *Plugin) heartbeat() error { if !p.catchupDone { p.catchupDone = true if resp.PendingMails > 0 { - go p.catchUp(resp.PendingMails) + go p.catchUp(resp.PendingMails, resp.PendingWorkspaces) } } @@ -545,38 +564,59 @@ func (p *Plugin) heartbeat() error { // - 正序(最旧的先处理),保持时间线 // - 只补 normal 类型(permission 不补投——人在 WebUI 上看到就知道了) // - 每封之间等 InjectInputSync 返回(串行处理) -func (p *Plugin) catchUp(pending int) { +func (p *Plugin) catchUp(pending int, workspaces []string) { limit := catchupLimit if pending < limit { limit = pending } - log.Printf("[homeagent-mail-bridge] 补投 %d 封离线期间的邮件(共 %d 封未读)", limit, pending) - - var inbox struct { - Mails []struct { - MailID string `json:"mail_id"` - FromName string `json:"from_name"` - Subject string `json:"subject"` - MailType string `json:"mail_type"` - ReplyTo string `json:"reply_to"` - // 补拉路径也必须知道发件方是人还是 Agent:Agent 之间不自动回信。 - // 缺了它补投的邮件会被保守当成 Agent 来信,于是人发的那封失去自动回复。 - FromHuman bool `json:"from_human"` - // parent_mail_id 非空 = 这封是回信。收件箱返回的字段名是它, - // 而 SSE 事件里叫 in_reply_to —— 两个名字指同一件事。 - ParentMailID string `json:"parent_mail_id"` - // from_session_id 用于档位继承:模型调 send_mail 时,Gateway 据此 - // 从来源会话继承权限档位(InheritedMode)。 - SessionID string `json:"session_id"` - } `json:"mails"` - } - url := fmt.Sprintf("%s/api/v1/mail/inbox?status=unread&limit=%d", p.gwURL, limit) - if err := p.get(url, &inbox); err != nil { - log.Printf("[homeagent-mail-bridge] 补拉失败: %v", err) + // ★ 逐工作区补投,不再读一次全局收件箱。 + // + // 旧写法不带任何收窄 ⇒ 会把**所有工作区**的漏投一起重放。清单来自心跳的 + // `pending_workspaces`(心跳是进程级的,没有"我的工作区"可言,所以由它给清单)。 + // + // 拿不到清单就不补投:带着空 workspace 去拉,服务端会 400。 + if len(workspaces) == 0 { + log.Printf("[homeagent-mail-bridge] 补投跳过:pending_mails=%d 但服务端未给出 pending_workspaces", pending) return } + delivered := 0 + for _, ws := range workspaces { + n, err := p.catchUpWorkspace(ws, limit, pending) + if err != nil { + // 单个工作区失败不影响其余(与"心跳失败不报错"同一原则) + log.Printf("[homeagent-mail-bridge] 补投工作区 %s 失败(不影响其余): %v", ws, err) + continue + } + delivered += n + } + if delivered > 0 { + log.Printf("[homeagent-mail-bridge] 补投 %d 封离线期间的邮件(共 %d 封未读,跨 %d 个工作区)", + delivered, pending, len(workspaces)) + } +} +// catchUpWorkspace 补投**一个工作区**的未读邮件,返回实际投递封数。 +func (p *Plugin) catchUpWorkspace(workspace string, limit, pending int) (int, error) { + var inbox struct { + Mails []struct { + MailID string `json:"mail_id"` + FromName string `json:"from_name"` + Subject string `json:"subject"` + MailType string `json:"mail_type"` + ReplyTo string `json:"reply_to"` + FromHuman bool `json:"from_human"` + ParentMailID string `json:"parent_mail_id"` + SessionID string `json:"session_id"` + } `json:"mails"` + } + url := fmt.Sprintf("%s/api/v1/mail/inbox?status=unread&limit=%d&workspace=%s", + p.gwURL, limit, url.QueryEscape(workspace)) + if err := p.get(url, &inbox); err != nil { + return 0, err + } + + delivered := 0 for _, m := range inbox.Mails { if m.MailType != "normal" { continue // permission 等非邮件驱动的不补投 @@ -639,8 +679,13 @@ func (p *Plugin) catchUp(pending int) { ) p.currentSessionID = m.SessionID + // 补投这一轮同样要设工作区:模型会在这轮里调 read_inbox。 + // 不设就是 400(实测 2026-09-26 07:42 起线上就在报这个)。 + p.currentWorkspace = workspace reply := p.sdk.InjectInputSync(p.name, outputChannelName, prompt) p.currentSessionID = "" + p.currentWorkspace = "" + delivered++ if reply == "" { // B-6:模型没回,发一封告知。发出去就算处理完(理由同 handleNewMail)。 p.sendFailureReply(m.FromName, m.Subject, m.MailID, "模型未产生回复") @@ -685,6 +730,7 @@ func (p *Plugin) catchUp(pending int) { } p.ledger.complete(m.MailID) } + return delivered, nil } // ─── SSE ─── @@ -968,8 +1014,10 @@ func (p *Plugin) handleNewMail(evt mailEvent, resumed bool) { // InjectInputSync 阻塞等待 agent 处理完毕,返回最终回复文本。 // 工具 handler 没有独立的 session 上下文,因此在本轮处理期间暂存来源会话。 p.currentSessionID = evt.SessionID + p.currentWorkspace = evt.Workspace reply := p.sdk.InjectInputSync(p.name, outputChannelName, prompt) p.currentSessionID = "" + p.currentWorkspace = "" // B-6:模型没回(空 = turn/end 信号 kind=error,或模型没说话) if reply == "" { @@ -1495,7 +1543,11 @@ func (p *Plugin) scopeQuery(sep string) string { if p.currentSessionID == "" { return "" } - return sep + "session_id=" + url.QueryEscape(p.currentSessionID) + q := sep + "session_id=" + url.QueryEscape(p.currentSessionID) + if p.currentWorkspace != "" { + q += "&workspace=" + url.QueryEscape(p.currentWorkspace) + } + return q } // inboxURL 拼收件箱地址。单独抽出来是为了能被单测直接断言 —— @@ -1508,5 +1560,11 @@ func (p *Plugin) inboxURL(status string, limit int) string { if p.currentSessionID != "" { scope = "&session_id=" + url.QueryEscape(p.currentSessionID) } + // ★ workspace 同样要带(见 currentWorkspace 字段的说明)。缺了服务端直接 400 + // —— 那是刻意的:旧语义(不带 = 全部工作区)正是用户报的那个越界缺陷。 + // 拿不到时不带,让服务端报 400:错误可见,好过静默跨工作区拿到别处的信。 + if p.currentWorkspace != "" { + scope += "&workspace=" + url.QueryEscape(p.currentWorkspace) + } return fmt.Sprintf("%s/api/v1/mail/inbox?status=%s&limit=%d%s", p.gwURL, status, limit, scope) } diff --git a/plugins/homeagent-mail-bridge/testsource_test.go b/plugins/homeagent-mail-bridge/testsource_test.go new file mode 100644 index 0000000..5630b05 --- /dev/null +++ b/plugins/homeagent-mail-bridge/testsource_test.go @@ -0,0 +1,17 @@ +package main + +import ( + "os" + "testing" +) + +// readSource 读同目录下的源文件 —— 结构性判据要盯着**源码形状**, +// 不是运行时行为(这个插件要 SDK 与宿主才能跑起来)。 +func readSource(t *testing.T, name string) string { + t.Helper() + b, err := os.ReadFile(name) + if err != nil { + t.Fatalf("读 %s: %v", name, err) + } + return string(b) +} diff --git a/plugins/opencode-mail-bridge/index.js b/plugins/opencode-mail-bridge/index.js index 9b89591..8c9dddd 100644 --- a/plugins/opencode-mail-bridge/index.js +++ b/plugins/opencode-mail-bridge/index.js @@ -1206,6 +1206,41 @@ export default async function mailBridge(input) { } } + /* + ★ 2026-09-26(用户报的):**投递即标已读**。 + + # 缺陷:桥投了信却不动 status ⇒ 每次重启都重投 ⇒ 回声 + + 投递路径只做两件事:起一轮、把 id 记进 `deliveredMails`。 + **没有一处调 /mail/read** —— 桥里唯一那处标已读在 `read_inbox` 工具里 + (下面第 300 行),要等模型自己去读收件箱。 + + 于是每封被投递的信**永远是 unread**;而 `catchUp` 按 `status=unread` 拉 + ⇒ 桥一重启(每次部署都会),积压的"未读"被当成离线漏投**再投一遍**。 + + 实测(pi 桥,同一天同一根因):5 个 mail_id 各进了**两条不同会话**; + 287 封 unread 中 **187 封已经回过信**。用户原话: + 「我都不记得我下达这个任务,是你的桥自动重投存在 bug」 + 「就是你的错误的重投机制造成了回声」 + + # 修法:deliveredMails(内存,重启即丢)与库里的 status 必须同时写 + + 只写内存 ⇒ 重启后桥只认库 ⇒ 重投。所以凡是标记"我接管了这封"的地方 + 都要走 `markDelivered`(漏一处就是一条重投路径)。 + + # 为什么不会"吞掉"信 + + 只改 status,不改内容、不删行;`read_inbox` 传 `all` 照样看得到。 + 而"已交给一个 worker 处理"正是它此刻的真实状态。 + */ + function markDelivered(client, mailId) { + if (!mailId) return; + deliveredMails.add(mailId); + apiPost("/mail/read", { mail_ids: [mailId] }).catch(e => + console.error(`[mail-bridge] 投递标已读失败 ${mailId}:`, e?.message || e) + ); + } + // 已经投过的 mail_id。心跳与 SSE 建连之间有个窗口:那期间到的邮件 // 既在 pending_mails 里、也会被 SSE 推一次 —— 不去重就会投两遍。 // @@ -1245,7 +1280,7 @@ export default async function mailBridge(input) { // 逐封再查一次:拉收件箱和逐封投递之间 SSE 可能已经投过其中某封 // (selectCatchup 只在拉完那一刻去过重) if (deliveredMails.has(ev.mail_id)) continue; - deliveredMails.add(ev.mail_id); + markDelivered(client, ev.mail_id); try { await deliverMail(client, directory, ev, "mail"); delivered += 1; @@ -1311,7 +1346,7 @@ export default async function mailBridge(input) { } if (type !== "new_mail") return; - if (data?.mail_id) deliveredMails.add(data.mail_id); + if (data?.mail_id) markDelivered(client, data.mail_id); // 决策回执不是"新任务"(内容已随 permission_decision 交付)—— 同 pi 桥那条注释: // 按新邮件再投一次会把人类真正的新邮件挤在这条会话的队列后面。 if (data?.mail_type === "permission_decision") { diff --git a/plugins/opencode-mail-bridge/test/delivery-marks-read.test.mjs b/plugins/opencode-mail-bridge/test/delivery-marks-read.test.mjs new file mode 100644 index 0000000..3bb6009 --- /dev/null +++ b/plugins/opencode-mail-bridge/test/delivery-marks-read.test.mjs @@ -0,0 +1,55 @@ +/** + * ★ 投递即标已读 —— 防「插件重启 → 重投 → 回声」。 + * + * 与 pi-mail-bridge 的同名判据同源:三个桥共享同一套 Gateway 契约, + * 缺陷也一样(投递只写内存 deliveredMails、不动库里的 status ⇒ catchUp + * 按 unread 重投)。用户报的是 pi 桥,但三处必须一起修,否则换个平台复发。 + * + * 实证:pi 侧 5 个 mail_id 各进了两条不同会话;287 封 unread 中 187 封 + * 已经回过信。用户原话:「我都不记得我下达这个任务,是你的桥自动重投存在 bug」 + */ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; + +const src = readFileSync(new URL('../index.js', import.meta.url), 'utf8'); + +test('★ deliveredMails.add 只允许出现在 markDelivered 内部', () => { + const lines = src.split('\n'); + const fnStart = lines.findIndex(l => /function markDelivered\(/.test(l)); + assert.ok(fnStart >= 0, 'markDelivered 必须存在'); + let fnEnd = fnStart; + for (let i = fnStart + 1; i < lines.length; i++) { + if (/^ \}/.test(lines[i])) { fnEnd = i; break; } + } + const strays = []; + lines.forEach((line, i) => { + if (!/deliveredMails\.add\(/.test(line)) return; + if (i > fnStart && i < fnEnd) return; + if (/^\s*(\/\/|\*|\/\*)/.test(line)) return; + strays.push(`${i + 1}: ${line.trim()}`); + }); + assert.deepEqual(strays, [], + `这些地方绕过 markDelivered ⇒ 库里 status 不更新 ⇒ 重启重投:\n${strays.join('\n')}`); +}); + +test('markDelivered 同时写内存与库', () => { + const i = src.indexOf('function markDelivered('); + const body = src.slice(i, src.indexOf('\n }', i)); + assert.match(body, /deliveredMails\.add\(/, '要写内存集合'); + assert.match(body, /apiPost\("\/mail\/read"/, '要写库里的 status'); + assert.match(body, /mail_ids: \[mailId\]/, '按 id 标(不传会被要求 workspace)'); +}); + +test('★ 两个投递点都走 markDelivered', () => { + assert.match(src, /if \(deliveredMails\.has\(ev\.mail_id\)\) continue;[^\n]*\n\s*markDelivered\(client, ev\.mail_id\);/, + '补投路径'); + assert.match(src, /if \(data\?\.mail_id\) markDelivered\(client, data\.mail_id\);/, + 'SSE 主投递路径'); +}); + +test('判据自检:反例必须判红', () => { + const before = 'if (data?.mail_id) deliveredMails.add(data.mail_id);'; + assert.ok(/deliveredMails\.add\(/.test(before), '反例确实会 add'); + assert.ok(!/mail\/read/.test(before), '反例不标已读 ⇒ 判红'); +}); diff --git a/plugins/pi-mail-bridge/src/index.mjs b/plugins/pi-mail-bridge/src/index.mjs index f3e7674..60f443d 100644 --- a/plugins/pi-mail-bridge/src/index.mjs +++ b/plugins/pi-mail-bridge/src/index.mjs @@ -101,6 +101,54 @@ const log = (...args) => console.error('[pi-mail-bridge]', ...args); // SSE 断线重放)都发生在秒到分钟级,几千封之前的 id 不可能再来。 const deliveredMails = new BoundedSet(MAX_TRACKED_MAILS); +/* +★ 2026-09-26(用户报的):**投递即标已读**。 + +# 缺陷:桥投了信却不动 status ⇒ 每次重启都重投一遍 ⇒ 回声 + +投递路径(SSE `new_mail` / 心跳补投)只做两件事:起 worker、把 id 记进 +`deliveredMails`。**没有任何一处调 `/mail/read`** —— 桥里唯一那处标已读 +在 `read_inbox` 工具里,要等模型自己去读收件箱。 + +于是每封被投递的信**永远是 unread**;而 `catchUp` 按 `status=unread` 拉。 +⇒ 桥一重启(每次部署都会),积压的"未读"被当成离线漏投**再投一遍**。 + +实测(2026-09-26): + · 5 个 mail_id 各出现在**两条不同 pi 会话**里(04:54 一条、08:01 一条) + · 同一封信被投两次 ⇒ 两个 worker 各回一封 ⇒ 对方收到两封 ⇒ 各回两封… + · pi 的收件箱 287 封 unread 中 **187 封已经回过信了** + (`EXISTS(SELECT 1 FROM mails r WHERE r.parent_mail_id=m.mail_id)`) + · 两条会话各烧到 463 / 268 封,用户原话:「我都不记得我下达这个任务, + 是你的桥自动重投存在 bug」「就是你的错误的重投机制造成了回声」 + +# 修法:把"投递"与"已读"合成一个动作 + +`deliveredMails` 是**内存**里的去重集合,重启即丢。数据库的 `status` 才是 +跨重启的那个"我接管过了"记录 —— 两者必须同时写,否则口径不一致: + + 内存说"投过了"、库说"没读过" ⇒ 重启后桥只认库 ⇒ 重投。 + +因此凡是把一封标记进 `deliveredMails` 的地方,都要同时把它在**库里**标成已读。 +这就是 `markDelivered` 存在的理由:不漏一处(漏一处就是一条重投路径)。 + +# 为什么标已读不会"吞掉"信 + +标已读只改 status,不改内容、不删行;`read_inbox` 传 `status=all` 依然看得到。 +而"这封信已经被交给一个 worker 处理了"正是它此刻的真实状态 —— 库里记的 +本来就该是这件事,而不是"模型有没有顺手调过 read_inbox"。 + +# 失败不阻塞 + +标不上不该让投递失败(正文已经在 worker 手里,代价只是下次重启可能再投一次, +与旧行为一致、不更差)。写日志,好让"标已读坏了"这件事可见。 +*/ +function markDelivered(mailId) { + if (!mailId) return; + deliveredMails.add(mailId); + client?.post('/mail/read', { mail_ids: [mailId] }).catch((e) => + log(`投递标已读失败 ${mailId}: ${e?.message || e}`)); +} + let allowedModels = []; let modelRuntime = null; let client = null; @@ -165,7 +213,7 @@ function handlePermissionDecision(data) { // 决策回执的内容马上随这次恢复交给 worker(备注 + 等人期间新到的邮件都在里面), // 所以先把它记成"已交付":稍后那封同内容的邮件就不会被当成新任务再起一轮。 const decisionMailID = data?.decision_mail_id || ''; - if (decisionMailID) deliveredMails.add(decisionMailID); + if (decisionMailID) markDelivered(decisionMailID); // 拉"等人期间新到的邮件"要发一次 HTTP,而本函数跑在 SSE 读循环上(必须廉价) // —— 所以整体转异步:先把事件收下,几毫秒后带着上下文去唤醒 worker。 @@ -316,7 +364,7 @@ async function catchUp(pending, workspaces) { const tasks = selectCatchup(box?.mails ?? box, deliveredMails); for (const ev of tasks) { if (deliveredMails.has(ev.mail_id)) continue; // 逐封再查(B-7.6) - deliveredMails.add(ev.mail_id); + markDelivered(ev.mail_id); pool.submit('mail', ev); delivered += 1; } @@ -462,7 +510,7 @@ function handleSSEEvent(type, data) { if (data?.role && data.role !== 'to' && data.role !== 'cc') return; const id = data?.mail_id; if (!id || deliveredMails.has(id)) return; // B-3 第 1 步:去重 - deliveredMails.add(id); + markDelivered(id); // ★ 决策回执**不是新任务**(2026-09-13 线上缺陷)。 // // 它长得像普通邮件("Re: 权限请求 - 拒绝"),内容却已经随 SSE 的 diff --git a/plugins/pi-mail-bridge/test/delivery-marks-read.test.mjs b/plugins/pi-mail-bridge/test/delivery-marks-read.test.mjs new file mode 100644 index 0000000..489d2c2 --- /dev/null +++ b/plugins/pi-mail-bridge/test/delivery-marks-read.test.mjs @@ -0,0 +1,99 @@ +/** + * ★ 投递即标已读 —— 防「桥重启 → 重投 → 回声」。 + * + * # 缺陷(用户报的,2026-09-26) + * + * 投递路径只把 mail_id 记进内存的 `deliveredMails`,**不动库里的 status**。 + * 而 `catchUp` 按 `status=unread` 拉 ⇒ 桥一重启(每次部署都会),积压的 + * "未读"被当成离线漏投**再投一遍**: + * + * · 5 个 mail_id 各进了**两条不同 pi 会话**(04:54 一条、08:01 一条) + * · 同一封信投两次 ⇒ 两个 worker 各回一封 ⇒ 对方收到两封 ⇒ 各回两封… + * · pi 收件箱 287 封 unread 中 **187 封已经回过信了** + * · 两条会话各烧到 463 / 268 封 + * + * 用户原话:「我都不记得我下达这个任务,是你的桥自动重投存在 bug」 + * 「就是你的错误的重投机制造成了回声」 + * + * # 判据要钉住什么 + * + * 不是"有没有 markDelivered 这个函数"(那太弱),而是: + * ① `deliveredMails.add` **只能**出现在 markDelivered 内部 —— 别处直接 add + * 就是一条绕过标已读的重投路径(这正是缺陷的形状) + * ② markDelivered 必须真的 POST /mail/read + * ③ 三个投递点(SSE / 补投 / 决策回执)都走它 + */ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; +import { fileURLToPath } from 'node:url'; +import { dirname, join } from 'node:path'; + +const src = readFileSync( + join(dirname(fileURLToPath(import.meta.url)), '..', 'src', 'index.mjs'), + 'utf8', +); + +test('★ deliveredMails.add 只允许出现在 markDelivered 内部', () => { + /* + 这一条是整组判据的核心。 + + 任何别处的 `deliveredMails.add(x)` 都意味着「我接管了这封,但库里还是 + unread」⇒ 下次重启重投。缺陷就是这么长出来的:四处投递各自 add, + 没有一处标已读。 + + 允许的行号集合:markDelivered 函数体的那几行。 + */ + const lines = src.split('\n'); + const fnStart = lines.findIndex(l => /^function markDelivered\(/.test(l)); + assert.ok(fnStart >= 0, 'markDelivered 必须存在'); + + // 函数体到下一个顶层 } 为止 + let fnEnd = fnStart; + for (let i = fnStart; i < lines.length; i++) { + if (/^\}/.test(lines[i]) && i > fnStart) { fnEnd = i; break; } + } + + const strays = []; + lines.forEach((line, i) => { + if (!/deliveredMails\.add\(/.test(line)) return; + if (i > fnStart && i < fnEnd) return; // 在 markDelivered 体内,合法 + if (/^\s*(\/\/|\*|\/\*)/.test(line)) return; // 注释里提到它(说明文字) + strays.push(`${i + 1}: ${line.trim()}`); + }); + assert.deepEqual(strays, [], + `这些地方绕过 markDelivered 直接 add ⇒ 库里的 status 不会被更新 ⇒ 重启重投:\n${strays.join('\n')}`); +}); + +test('markDelivered 同时写内存与库(两处口径必须一致)', () => { + const i = src.indexOf('function markDelivered('); + assert.ok(i > 0); + const body = src.slice(i, src.indexOf('\n}', i)); + assert.match(body, /deliveredMails\.add\(/, '要写内存集合'); + assert.match(body, /client\?\.post\('\/mail\/read'/, '要写库里的 status'); + assert.match(body, /mail_ids: \[mailId\]/, '按 id 标(不传 mail_ids 会被服务端要求 workspace)'); +}); + +test('★ 三个投递点都走 markDelivered(少一处就是一条重投路径)', () => { + // SSE new_mail 主路径 + // 注意:`return;` 后面可能跟行尾注释("// B-3 第 1 步:去重"), + // 所以不能写 `return;\s*\n` —— 那样会被注释挡住而误报缺失。 + assert.match(src, /if \(!id \|\| deliveredMails\.has\(id\)\) return;[^\n]*\n\s*markDelivered\(id\);/, + 'SSE 主投递路径'); + // 心跳补投(同样:`continue;` 后有行尾注释) + assert.match(src, /if \(deliveredMails\.has\(ev\.mail_id\)\) continue;[^\n]*\n\s*markDelivered\(ev\.mail_id\);/, + '补投路径'); + // 决策回执 + assert.match(src, /if \(decisionMailID\) markDelivered\(decisionMailID\);/, + '决策回执'); +}); + +test('判据自检:反例(别处 add、不标已读)必须判红', () => { + // 这正是修复前的形状 —— 判据若放过它,就防不住回归 + const before = 'if (data?.mail_id) deliveredMails.add(data.mail_id);'; + const lines = [before]; + const fnStart = lines.findIndex(l => /^function markDelivered\(/.test(l)); + assert.equal(fnStart, -1, '反例里没有 markDelivered ⇒ 回退到最弱判据'); + assert.ok(/deliveredMails\.add\(/.test(before), '反例确实会 add'); + assert.ok(!/mail\/read/.test(before), '反例确实不标已读 ⇒ 会判红'); +}); diff --git a/plugins/pi-mail-bridge/test/permission-note.test.mjs b/plugins/pi-mail-bridge/test/permission-note.test.mjs index 593b470..261dd58 100644 --- a/plugins/pi-mail-bridge/test/permission-note.test.mjs +++ b/plugins/pi-mail-bridge/test/permission-note.test.mjs @@ -135,8 +135,12 @@ export const WIRING = [ must: /reason: renderDecisionReason\(/ }, { file: 'index.mjs', what: '决策事件把 note 传给 pool', must: /routePermission\(relayKey, decision, note, waiting\)/ }, + // ★ 2026-09-26:这一条原来钉的是 `deliveredMails.add(decisionMailID)` —— + // 直接写内存集合。那个形状正是「桥重启后重投」缺陷的载体(deliveredMails + // 只在内存、库里的 status 从没被写)。现在统一走 markDelivered: + // 内存与库一起写。语义没变(仍"不再当新任务"),载体变了。 { file: 'index.mjs', what: '决策回执记成已交付(不再当新任务)', - must: /deliveredMails\.add\(decisionMailID\)/ }, + must: /markDelivered\(decisionMailID\)/ }, { file: 'index.mjs', what: 'new_mail 分支认得决策回执', must: /mail_type === 'permission_decision'/ }, { file: 'turn.mjs', what: '通知投递路径也带备注', diff --git a/server/internal/handler/mail.go b/server/internal/handler/mail.go index e3121a9..b11d89c 100644 --- a/server/internal/handler/mail.go +++ b/server/internal/handler/mail.go @@ -374,6 +374,45 @@ func SendMail(w http.ResponseWriter, r *http.Request) { relayFree = human } + /* + ★ 防"无人决策的 Agent↔Agent 往返"——**不看 relay 标记**。 + + # 为什么需要这第三道防线(2026-09-26 实测) + + 生产上 `pi ↔ dsh` 一天 175 封、正文 612KB,其中 **relay = 0**: + 全部是模型**主动**调 send_mail。而上面两道都只覆盖 `relay != ""`: + · 会话预算 —— relay 才扣(且该会话 max_rounds=0 本是"不限") + · maxRelayHops —— `if relay != "" {` 根本没进 + · 插件自动转发守卫 —— 日志明说"本轮不自动转发",守卫工作正常, + 但**模型自己发的不受它管** + ⇒ 三条全绕开,没有任何一层在数这个。 + + # 与 maxRelayHops 的关系 + + 同构(连续计数 + 人类参与即归零),但**独立**: + · relay 那道管"插件代劳的回路"(每封都无新信息,阈值 5) + · 这道管"模型主动的回路"(协同中来回十几次正常,阈值 8) + 两道都要:合起来才覆盖"谁在发"这个维度的两个取值。 + + # 触发时做什么 + + **拒绝并说清怎么办**,不静默限流 —— 与 hop 那道同样的处理。 + 人类插一句话或模型说明为何必须继续,都能立刻恢复。 + */ + if hops, hErr := repo.CountTrailingAgentPingPong(r.Context(), sessionID); hErr == nil { + // 收件方是人类时不算回路(人类在回路里,正是我们要的"有人决策")。 + isHuman, _ := repo.IsHumanUser(r.Context(), to.Name) + if !isHuman && hops >= repo.MaxAgentPingPong() { + Error(w, http.StatusForbidden, fmt.Sprintf( + "本会话已连续 %d 封 Agent 之间互相回信、其中没有任何人类参与(上限 %d)。"+ + "这通常意味着两个 Agent 在互相确认而无人决策 —— 生产上实测过一天 175 封、"+ + "正文 612KB 却没有任何工作产出。请由人类在会话里插一句话(计数即归零),"+ + "或说明为何这轮必须继续。", + hops, repo.MaxAgentPingPong())) + return + } + } + if relay != "" { // 硬上限:一条会话里**连续**的 relay 邮件不得超过上限。 // diff --git a/server/internal/repo/agentloop.go b/server/internal/repo/agentloop.go new file mode 100644 index 0000000..93c1fb4 --- /dev/null +++ b/server/internal/repo/agentloop.go @@ -0,0 +1,80 @@ +package repo + +import ( + "context" + + "github.com/agentmail/gateway/internal/db" + "github.com/google/uuid" +) + +/* +防的是**无人决策的 Agent↔Agent 往返**,与 relay 无关。 + +# 为什么已有的两道防线都拦不住(2026-09-26 实测) + +生产上 `pi ↔ dsh` 一天跑了 **175 封**、正文 612KB,其中 **relay = 0** —— +全部是模型**主动**调 send_mail。而两道防线都只覆盖 `relay != ""` 那条路: + + ① 会话预算:`relay != ""` 才扣 ⇒ 主动发信不扣(且该会话 max_rounds=0=不限) + ② maxRelayHops=5:`if relay != "" { ... }` ⇒ 根本不进那个分支 + ③ 插件自动转发守卫:日志明说"本轮不自动转发:来信方 dsh 是 Agent" + —— 守卫**工作正常**,但模型自己发的不受它管 + +⇒ 三条全绕开。没有任何一层在数"两个 Agent 在没有人类参与的情况下连续往返了几次"。 + +# 判据的形状(与 maxRelayHops 同构,但不看 relay 标记) + +从最新一封往前扫,**连续**的"发件方是 Agent 且收件方也是 Agent"的邮件计数, +遇到以下任一即停(= 归零): + + · 人类发的信(`from_name` 是用户) + · 人类收的信(`to_name` 是用户) + +"连续"与"人类参与即归零"这两点与 maxRelayHops 一致 —— 正常的 +「Agent 干完活回一封、人类看一眼再派下一轮」不受影响;只掐住 +「全程没有人类说话」的那种回路。 + +# 为什么阈值取 8 而不是 5 + +maxRelayHops=5 管的是**插件自动转发**:那种回路每封都没有新信息,5 跳足够。 +而这里是模型**主动**发信 —— 两个 Agent 协同一件事时,来回十几次是正常的 +(本轮真实案例:`mc` 插件修复本身就有 6 封有效往返)。 + +8 是"明显超出正常协同、但还没到烧掉一天 175 封"的位置。触发时的处理是 +**拒绝并说清怎么办**(人类插一句、或模型说明为何必须继续),不是静默限流。 +*/ +const maxAgentPingPong = 8 + +// CountTrailingAgentPingPong 数会话尾部**连续的、无人类参与**的 Agent↔Agent 邮件数。 +// +// 返回值即「若本次再发一封,它会是第几跳」的前一个数。 +func CountTrailingAgentPingPong(ctx context.Context, sessionID uuid.UUID) (int, error) { + rows, err := db.DB.QueryContext(ctx, ` + SELECT + EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human, + EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human + FROM mails m + WHERE m.session_id = $1 + ORDER BY m.created_at DESC, m.mail_id DESC`, sessionID) + if err != nil { + return 0, err + } + defer rows.Close() + + n := 0 + for rows.Next() { + var fromHuman, toHuman bool + if err := rows.Scan(&fromHuman, &toHuman); err != nil { + return 0, err + } + // 人类参与(发或收)即打断连续性 —— 与 maxRelayHops 同一个"连续"语义。 + if fromHuman || toHuman { + break + } + n++ + } + return n, rows.Err() +} + +// MaxAgentPingPong 供 handler 与测试共用同一个阈值。 +func MaxAgentPingPong() int { return maxAgentPingPong } diff --git a/server/internal/repo/agentloop_test.go b/server/internal/repo/agentloop_test.go new file mode 100644 index 0000000..e47b096 --- /dev/null +++ b/server/internal/repo/agentloop_test.go @@ -0,0 +1,139 @@ +package repo + +import ( + "context" + "testing" + "time" + + "github.com/agentmail/gateway/internal/db" + "github.com/google/uuid" +) + +/* +★ 第三道防线:数**无人决策的 Agent↔Agent 连续往返**(与 relay 无关)。 + +# 这个缺陷的形状(2026-09-26 生产实测) + +`pi ↔ dsh` 一天 175 封、正文 612KB,其中 **relay = 0**:全部是模型**主动** +调 send_mail。而既有两道防线都只覆盖 `relay != ""`: + + · 会话预算 —— relay 才扣 + · maxRelayHops —— `if relay != "" {` 根本没进 + · 插件自动转发守卫 —— 日志明说"本轮不自动转发",守卫正常,但不管模型主动发 + +⇒ 三条全绕开。这三条判据钉住第四道的**判据形状**(连续性、人类参与即归零), + 而不是钉"有没有这个函数"。 +*/ + +/* +seedPingPong 造一封 Agent 之间的信。 + +★ **必须显式给 created_at,且逐封递增。** + +`mails.created_at` 的默认值是 `CURRENT_TIMESTAMP`(**秒**精度),而 +`CountTrailingAgentPingPong` 按 `created_at DESC, mail_id DESC` 从后往前扫。 +同一秒内插的多封,排序实际由随机 UUID 决定 ⇒ 连续段数**不确定**。 + +我第一版就是这么写的,当场被这条判据自己抓到:同一段数据一次读出 2、一次读出 1 +("人类插话后应只数它之后的 2 封,实际 1")。这与本仓记录过的老坑同一个形状 +("SQLite 时间戳只有秒精度:同秒多封排序不确定")。 + +生产侧不受影响:`mails` 的 INSERT 显式传 `NOW()`(微秒)。这里补上同样的事。 +*/ +var pingPongSeq int + +func seedPingPong(t *testing.T, sessionID uuid.UUID, from, to string) { + t.Helper() + pingPongSeq++ + // 以固定基准 + 递增秒数,确保先后顺序与插入顺序一致 + ts := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC).Add(time.Duration(pingPongSeq) * time.Second) + if _, err := db.DB.ExecContext(context.Background(), + `INSERT INTO mails (session_id, from_name, to_name, subject, body, created_at) + VALUES ($1, $2, $3, 'pingpong', 'b', $4)`, sessionID, from, to, ts); err != nil { + t.Fatal(err) + } +} + +func TestAgentPingPongCountsConsecutive(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, err := CreateSession(ctx, nil, "human", "回路", "") + if err != nil { + t.Fatal(err) + } + + // 三个来回 = 6 封,全是 Agent↔Agent(from/to 都不是人类用户名) + for i := 0; i < 3; i++ { + seedPingPong(t, sid, "pi", "dsh") + seedPingPong(t, sid, "dsh", "pi") + } + if n, err := CountTrailingAgentPingPong(ctx, sid); err != nil || n != 6 { + t.Fatalf("三个来回应为 6,实际 %d(err=%v)", n, err) + } +} + +// ★ 核心语义:**人类插一句话,计数归零**。 +// +// 这一条决定了这个判据会不会误伤正常用法: +// · 人类说一句 → Agent 回十封 → 人类再一句 ⇒ 每次都在阈值以下 +// · 没有人说话 → 一直涨 ⇒ 到阈值拒绝 +// 与 maxRelayHops 的"连续"同构 —— 两条防线在这个语义上必须一致, +// 否则同一个现象在两处得到不同的结论。 +func TestAgentPingPongResetsOnHuman(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, err := CreateSession(ctx, nil, "human", "回路", "") + if err != nil { + t.Fatal(err) + } + /* + ★ 判据依赖 "jianf 是人类" —— 而测试库是空的(`setupTestDB` 只跑迁移, + 不建用户)。不种的话 `IsHumanUser` 查不到,这封信会被当成 Agent 发的, + 于是一条本应绿的判据变红。 + + 我第一版就漏了这一步,当场被这条判据自己抓到("应只数 2 封,实际 3")。 + ⇒ 判据里任何"某名字是人类/Agent"的前提都必须**显式种下**, + 不能靠库里碰巧有。 + */ + if _, err := db.DB.ExecContext(ctx, + `INSERT INTO users (username, display_name, password_hash, role) + VALUES ('jianf', 'jianf', 'x', 'admin')`); err != nil { + t.Fatal(err) + } + + // 先造一封人类的信(计数应当从它之后才开始) + seedPingPong(t, sid, "jianf", "pi") + seedPingPong(t, sid, "pi", "dsh") + seedPingPong(t, sid, "dsh", "pi") + + if n, err := CountTrailingAgentPingPong(ctx, sid); err != nil || n != 2 { + t.Fatalf("人类插话后应只数它之后的 2 封,实际 %d(err=%v)", n, err) + } + + // 人类**收**信也算打断(to 是人类)—— 不是只认发件方 + seedPingPong(t, sid, "pi", "jianf") + if n, err := CountTrailingAgentPingPong(ctx, sid); err != nil || n != 0 { + t.Fatalf("人类收信也应打断连续性(应 0),实际 %d(err=%v)", n, err) + } +} + +// 反向对照:**没有人类参与**时计数不该被重置 —— 否则这道防线永远不触发 +// (把"人类参与才 reset"实现成"每封都 reset",上面那条仍会绿)。 +func TestAgentPingPongNotResetByAgents(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, err := CreateSession(ctx, nil, "human", "回路", "") + if err != nil { + t.Fatal(err) + } + seedPingPong(t, sid, "pi", "dsh") + seedPingPong(t, sid, "opencode", "homeagent") // 换一对 Agent 也仍是"无人决策" + seedPingPong(t, sid, "dsh", "pi") + + if n, err := CountTrailingAgentPingPong(ctx, sid); err != nil || n != 3 { + t.Fatalf("纯 Agent 往返应连续计数(应 3),实际 %d(err=%v)", n, err) + } + if MaxAgentPingPong() <= 0 { + t.Fatal("阈值必须为正,否则守卫恒真或恒假") + } +}