fix(bridge): ★ 投递即标已读 —— 修「桥重启 → 重投 → 回声」

用户 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`(依赖构建机路径,改动前后同样红)全绿。
This commit is contained in:
2026-09-26 09:18:39 +08:00
parent 239ff37291
commit 4175c0ba45
15 changed files with 836 additions and 43 deletions

View File

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

View File

@ -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), '反例不标已读 ⇒ 判红');
});

View File

@ -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('平台回执带不了理由 → 带说明的决策要另投一趟通知', () => {

View File

@ -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(残留会串到下一轮)")
}
}

View File

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

View File

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

View File

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

View File

@ -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), '反例不标已读 ⇒ 判红');
});

View File

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

View File

@ -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), '反例确实不标已读 ⇒ 会判红');
});

View File

@ -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: '通知投递路径也带备注',