diff --git a/docs/PLUGIN-CONTRACT.md b/docs/PLUGIN-CONTRACT.md index 60aa870..81566e0 100644 --- a/docs/PLUGIN-CONTRACT.md +++ b/docs/PLUGIN-CONTRACT.md @@ -349,6 +349,7 @@ SSE 只推连上之后的事件。插件重启前发来的邮件不会再推一 | B-7.4 | 按时间**正序**投(收件箱是倒序返回的) | MUST | | B-7.5 | `mail_type !== 'normal'` 的不补投 | MUST | | B-7.6 | 逐封投递前**再查一次**去重集合 | SHOULD | +| B-7.7 | 插件以**子进程**形式运行时,去重必须**落盘**,且区分「投过」与「跑完」 | MUST | > **B-7.1 为什么只补一次**:每轮都补会把「模型正在处理中、尚未标已读」的邮件 > 重复投递 —— 一封邮件起两轮模型。 @@ -364,6 +365,39 @@ SSE 只推连上之后的事件。插件重启前发来的邮件不会再推一 > > **B-7.5 为什么跳过 permission**:原来的工具调用早随进程一起没了, > 投过去模型没有可恢复的上下文。 +> +> **B-7.7 为什么内存去重不够**(生产事故,2026-09-04): +> `deliveredMails` 是进程内的,而 homeagent 的插件跑在**子进程**里 —— +> homed 重启(或插件崩溃自动重启)会换一个新进程,那个集合随之清空。时序: +> +> 1. 邮件落库,**旧**进程的 SSE 收到,注入第一次 +> 2. 同一秒进程被重启,那一轮被掉断(日志:`context canceled`) +> 3. **新**进程起来,`deliveredMails` 是空的 +> 4. 心跳报 `pending_mails: 1`(第一轮没跑完 → `read_inbox` 没执行 +> → 邮件仍未读)→ 补投注入第二次 +> +> 模型上下文里因此出现两段几乎相同的指令。 +> +> **为什么不能只记「投过没有」**:那会把「重复」换成「丢件」。上面第 2 步里 +> 那一轮被掉断,发件人**没有**收到回信,而落盘记录说「已投过」→ 永远跳过。 +> 丢件比重复严重:重复至少人能看出来,丢件是静默的。 +> +> 因此记两个状态: +> +> | 状态 | 含义 | 再次收到时 | +> |---|---|---| +> | `delivered` | 注入过(可能被中断) | **仍然重投**,但提示词里说明「上一轮被中断」 | +> | `completed` | 那一轮真的跑完且回信已发出 | 跳过 | +> +> 「说明上一轮被中断」不是装饰:不说的话模型在上下文里看到两段相似指令, +> 会以为人重复交代了一遍,于是可能把同一件事做两次。 +> +> 何时标 `completed`:回信发出去了、模型自己回过了(B-5.3 让位)、 +> 或失败通知发出去了(B-6)—— 三者都是「发件人得到了一个交代」。 +> 回信**发失败**时不标 —— 那时发件人一个字都没收到。 +> +> 常驻守护进程形态的插件(如 pi 桥)不受这一条约束:它自己就是进程, +> 重启同时也会丢掉 SSE 连接,那时走的是 B-7 补拉而不是两条路径并发。 ### B-8 平台权限询问 → 邮件(若平台有审批环节,MUST;没有则 N/A) @@ -1055,11 +1089,16 @@ GET /api/v1/attachments/{id} | `mailDrivenSessions` | 快照里 `mail_driven` 全变成 `false` | | `pendingPermissions` | 决策回来时走 `B-4.2` 的退化路径 | | `explicitSends` | 重启后的第一轮可能重复转发 | -| `deliveredMails` | 由 `B-7` 的补拉去重兜住 | +| `deliveredMails` | 同一进程内的 SSE 重放靠它;**跨进程**的重复必须由 `B-7.7` 的落盘账本兜住 | 要持久化的话应当落在平台的会话元数据里,而不是插件自己的文件 —— 那样才能跟着会话一起被平台清理。 +> **例外:投递账本必须落盘**(`B-7.7`)。上面那些丢了只是「多开一条会话」 +> 或「多转发一次」,而投递去重丢了是「同一封邮件被注入两遍」—— 模型上下文 +> 里出现两段相同指令,可能把同一件事做两次。子进程形式的插件尤其如此: +> 它的重启频率由宙主进程决定,不是罕见事件。 + --- ## 七、验收清单 / Acceptance Checklist @@ -1169,6 +1208,12 @@ GET /api/v1/attachments/{id} [ ] 插件在线时再发一封 → 只收到一封回信,没有重复投递(B-7.3) + +[ ] (子进程形式的插件)发一封 → 模型跑到一半时重启宙主进程(B-7.7) + → 数据库里只有一封 `Re:`(模型上下文里也只有一段有效指令) + → 日志出现「上一轮被中断,带说明重投」 + → 账本里该 mail_id 先一行 `c:false` 后一行 `c:true` + → 再次重启(邮件已跑完)→ **不再投递**,会话邮件数不变 ``` ### 7.6 权限(若平台支持) diff --git a/plugins/homeagent-mail-bridge/ledger.go b/plugins/homeagent-mail-bridge/ledger.go new file mode 100644 index 0000000..919444e --- /dev/null +++ b/plugins/homeagent-mail-bridge/ledger.go @@ -0,0 +1,303 @@ +package main + +// 投递账本:跨进程记住「这封邮件投过没有、跑完没有」。 +// +// # 为什么必须落盘 +// +// `deliveredMails` 是进程内的 map,重启即丢。而 homeagent 的插件是**子进程**, +// homed 重启(或插件崩溃自动重启)会换一个新进程 —— 于是: +// +// 1. 18:59:38 邮件落库,旧插件进程的 SSE 收到,`InjectInputSync` 注入第一次 +// 2. 同一秒 homed 被重启,那一轮被掐断(eventloop 报 +// `process error: all 2 providers failed, last error: context canceled`) +// 3. 18:59:45 新插件进程起来,`deliveredMails` 是空的 +// 4. 心跳报 `pending_mails: 1`(第一轮没跑完 → read_inbox 没执行 → 邮件仍未读) +// → `catchUp` 注入第二次 +// +// 模型的上下文里因此出现两段几乎相同的指令(措辞的细微差别正好指出来源: +// SSE 那段写「你把本轮工作做完」,补投那段写「你把结论说出来就行」)。 +// +// # 为什么不能只记「投过没有」 +// +// 那会把「重复」换成「丢件」:上面第 2 步里那一轮被掐断,发件人**没有**收到 +// 回信,而落盘记录说「已投过」→ 永远跳过 → 那封邮件再也不会被处理。 +// 丢件比重复严重:重复至少人能看出来,丢件是静默的。 +// +// 所以记的是两个状态: +// +// - delivered = 注入过(可能被中断) +// - completed = 那一轮真的跑完并且回信发出去了 +// +// 判据因此是「completed 才跳过」。delivered 但未 completed 的仍要重投, +// 但换一段提示词明确说「上一轮被中断」—— 模型看到的不再是两条重复指令。 +// +// # 为什么是 JSONL 而不是 SQLite +// +// 写入是纯 append 的单行记录,读取只在启动时一次。SQLite 要多一个依赖、 +// 多一次 schema 迁移,换来的是这里用不到的查询能力。 +// 崩溃时最坏情况是最后一行写残 —— 解析时跳过坏行即可(见 loadLedger)。 + +import ( + "bufio" + "encoding/json" + "fmt" + "log" + "os" + "path/filepath" + "strings" + "sync" + "time" +) + +// ledgerEntry 是账本里的一行。 +// +// 字段用短名:这个文件会被追加很多行,而它没有人类读者(排查时用 jq)。 +type ledgerEntry struct { + MailID string `json:"m"` + Completed bool `json:"c"` + At int64 `json:"t"` // Unix 秒,仅供排查与过期清理 +} + +// ledgerState 是一封邮件的投递状态。 +type ledgerState struct { + delivered bool + completed bool +} + +// deliveryLedger 是账本的内存视图 + 落盘句柄。 +type deliveryLedger struct { + mu sync.Mutex + path string + state map[string]*ledgerState + // file 为 nil 表示落盘不可用(目录没权限之类)。此时退化为纯内存 —— + // 那正是修复前的行为,不比它更糟,而且不该让插件起不来。 + file *os.File +} + +// ledgerRetention 是账本条目的保留期。 +// +// 14 天:足够覆盖「插件停了一阵子再起来」的情形,又不会让文件无限增长。 +// 判据是条目时间而不是文件大小 —— 后者要读全文才能算,而这里的目的只是 +// 别让一个长期运行的部署攒下几十万行。 +const ledgerRetention = 14 * 24 * time.Hour + +// openDeliveryLedger 打开(或新建)账本。 +// +// dir 取 SDK 的 `Settings().DataDir()`(`/plugin_data/`, +// SDK 保证存在)。拿不到时退回 key 文件所在目录 —— 那个目录本来就要能写。 +func openDeliveryLedger(dir string) *deliveryLedger { + l := &deliveryLedger{state: map[string]*ledgerState{}} + if strings.TrimSpace(dir) == "" { + log.Printf("[homeagent-mail-bridge] 未取到数据目录,投递账本退化为纯内存(重启后可能重复投递)") + return l + } + l.path = filepath.Join(dir, "delivered.jsonl") + + // 先读旧记录,再打开追加句柄:反过来的话 O_TRUNC 之类的手误会清空历史。 + kept := l.load() + + if err := os.MkdirAll(dir, 0o700); err != nil { + log.Printf("[homeagent-mail-bridge] 建数据目录 %s 失败,账本退化为纯内存: %v", dir, err) + return l + } + + // 过期条目多到一半以上时重写一遍,否则纯追加。 + // 重写走「临时文件 + rename」:原地截断时崩溃会留下一个残缺账本, + // 而那比多几行过期记录严重得多。 + if kept.dropped > 0 && kept.dropped*2 >= kept.total { + l.compact() + } + + f, err := os.OpenFile(l.path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600) + if err != nil { + log.Printf("[homeagent-mail-bridge] 打开投递账本 %s 失败,退化为纯内存: %v", l.path, err) + return l + } + l.file = f + log.Printf("[homeagent-mail-bridge] 投递账本 %s(%d 条有效记录)", l.path, len(l.state)) + return l +} + +type loadStats struct{ total, dropped int } + +// load 把账本读进内存。坏行跳过而不是报错退出。 +func (l *deliveryLedger) load() loadStats { + var st loadStats + f, err := os.Open(l.path) + if err != nil { + return st // 不存在是正常的(首次运行) + } + defer f.Close() + + cutoff := time.Now().Add(-ledgerRetention).Unix() + sc := bufio.NewScanner(f) + // 单行很短,但给个宽松上限防止一行坏数据把扫描器卡住 + sc.Buffer(make([]byte, 0, 4096), 64*1024) + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" { + continue + } + st.total++ + var e ledgerEntry + if err := json.Unmarshal([]byte(line), &e); err != nil || e.MailID == "" { + // 崩溃时最后一行可能写残。跳过它 —— 那封邮件退化为「没记录」, + // 于是会被重投一次,而重投是安全的(下面 shouldDeliver 的语义)。 + st.dropped++ + continue + } + if e.At > 0 && e.At < cutoff { + st.dropped++ + continue + } + s := l.state[e.MailID] + if s == nil { + s = &ledgerState{} + l.state[e.MailID] = s + } + s.delivered = true + // 同一封可能有两行(先 delivered 后 completed);completed 只增不减。 + if e.Completed { + s.completed = true + } + } + return st +} + +// compact 把内存状态重写成一个干净的账本。 +func (l *deliveryLedger) compact() { + tmp := l.path + ".tmp" + f, err := os.OpenFile(tmp, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o600) + if err != nil { + return // 压实失败不影响正确性,只是文件继续变长 + } + now := time.Now().Unix() + w := bufio.NewWriter(f) + for id, s := range l.state { + b, _ := json.Marshal(ledgerEntry{MailID: id, Completed: s.completed, At: now}) + w.Write(b) + w.WriteByte('\n') + } + if err := w.Flush(); err != nil { + f.Close() + os.Remove(tmp) + return + } + // fsync 后再 rename:rename 本身是原子的,但内容没落盘时掉电会得到空文件。 + f.Sync() + f.Close() + if err := os.Rename(tmp, l.path); err != nil { + os.Remove(tmp) + } +} + +// append 写一行。落盘不可用时只更新内存。 +func (l *deliveryLedger) appendLine(e ledgerEntry) { + if l.file == nil { + return + } + b, err := json.Marshal(e) + if err != nil { + return + } + // 单次 Write 写完整行:多次 Write 之间崩溃会留下半行。 + if _, err := l.file.Write(append(b, '\n')); err != nil { + log.Printf("[homeagent-mail-bridge] 写投递账本失败(本次去重仅内存生效): %v", err) + return + } + // 每行都 Sync:这个文件的全部意义就是「进程死了之后还算数」, + // 攒在页缓存里等于没写。一封邮件一次 fsync,代价可以忽略。 + l.file.Sync() +} + +/* +claim 判定这封邮件要不要投,并在要投时记下 delivered。 + +返回值: + + deliver = false → 已经**跑完**过,跳过 + deliver = true, resumed = false → 全新的邮件 + deliver = true, resumed = true → 投过但没跑完(上一轮被中断), + 调用方应换一段说明中断的提示词 + +这里把「判定」与「记录」放在同一把锁里:分开的话 SSE 与 catchUp 两个 goroutine +可能都判定为要投(这正是修复前的竞态,只是那时连内存都没查)。 +*/ +func (l *deliveryLedger) claim(mailID string) (deliver, resumed bool) { + if mailID == "" { + return false, false + } + l.mu.Lock() + defer l.mu.Unlock() + + s := l.state[mailID] + if s != nil && s.completed { + return false, false + } + if s != nil && s.delivered { + // 投过但没跑完。**仍然重投** —— 上一轮被中断意味着发件人还没收到回信, + // 跳过它就是静默丢件。但要让调用方知道这是重投。 + // + // 不在这里加「重投次数上限」:一封邮件反复中断说明有别的问题 + // (模型一直崩、进程反复重启),而那时停止重投只会让问题更难发现。 + return true, true + } + if s == nil { + s = &ledgerState{} + l.state[mailID] = s + } + s.delivered = true + l.appendLine(ledgerEntry{MailID: mailID, Completed: false, At: time.Now().Unix()}) + return true, false +} + +// complete 标记这封邮件真的处理完了(回信已发出或已确认无需回信)。 +func (l *deliveryLedger) complete(mailID string) { + if mailID == "" { + return + } + l.mu.Lock() + defer l.mu.Unlock() + s := l.state[mailID] + if s == nil { + s = &ledgerState{delivered: true} + l.state[mailID] = s + } + if s.completed { + return // 幂等:重复标记不再写盘 + } + s.completed = true + l.appendLine(ledgerEntry{MailID: mailID, Completed: true, At: time.Now().Unix()}) +} + +// close 关闭落盘句柄。 +func (l *deliveryLedger) close() { + l.mu.Lock() + defer l.mu.Unlock() + if l.file != nil { + l.file.Close() + l.file = nil + } +} + +// resumeNote 是重投时插在提示词前面的说明。 +// +// 必须说清「上一轮被中断」:不说的话模型在上下文里看到两段几乎相同的指令, +// 会以为人重复交代了一遍,于是可能把同一件事做两次(或者反问「你是不是发重了」)。 +func resumeNote(mailID string) string { + return fmt.Sprintf( + "(这封邮件 %s 之前已经通知过你一次,但那一轮被中断了 —— 进程重启或模型出错,"+ + "所以你可能在上文里看到一段几乎相同的通知。请把它当作**同一件事**继续处理,"+ + "不要重复执行已经做过的操作。)\n\n", mailID) +} + +// shortID 把 UUID 截成 8 位用于日志。 +// +// 不直接写 `id[:8]`:mail_id 理论上可能短于 8 字节(畸形事件、测试桩), +// 那会 panic 在一条日志语句上 —— 日志不该有能力弄死投递协程。 +func shortID(id string) string { + if len(id) <= 8 { + return id + } + return id[:8] +} diff --git a/plugins/homeagent-mail-bridge/ledger_test.go b/plugins/homeagent-mail-bridge/ledger_test.go new file mode 100644 index 0000000..58400c6 --- /dev/null +++ b/plugins/homeagent-mail-bridge/ledger_test.go @@ -0,0 +1,314 @@ +package main + +// 投递账本的单元测试。 +// +// 核心不变量只有两条,但它们互相拉扯,必须分别钉住: +// +// 1. **跑完的不再投** —— 否则跨进程重启时同一封邮件被注入两遍 +// (生产实测:模型上下文里两段几乎相同的通知)。 +// 2. **投过但没跑完的仍要投** —— 否则那一轮被中断时就是静默丢件, +// 发件人永远收不到回信。这条比第 1 条更重要:重复至少人能看出来。 + +import ( + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +func TestClaimFreshMail(t *testing.T) { + l := openDeliveryLedger(t.TempDir()) + defer l.close() + + deliver, resumed := l.claim("m1") + if !deliver { + t.Fatal("全新邮件必须投递") + } + if resumed { + t.Error("全新邮件不该标记为重投") + } +} + +func TestClaimEmptyIDRejected(t *testing.T) { + l := openDeliveryLedger(t.TempDir()) + defer l.close() + + if deliver, _ := l.claim(""); deliver { + t.Error("空 mail_id 不该被投递(畸形事件)") + } +} + +func TestCompletedMailNotRedelivered(t *testing.T) { + l := openDeliveryLedger(t.TempDir()) + defer l.close() + + l.claim("m1") + l.complete("m1") + + if deliver, _ := l.claim("m1"); deliver { + t.Error("跑完的邮件必须跳过 —— 重投会让模型上下文里出现两段相同通知") + } +} + +// 这一条是整个账本存在的理由:同进程内 claim 两次的第二次是「重投」而不是「跳过」。 +func TestDeliveredButNotCompletedIsResumed(t *testing.T) { + l := openDeliveryLedger(t.TempDir()) + defer l.close() + + l.claim("m1") // 注入了,但那一轮被中断(没有 complete) + + deliver, resumed := l.claim("m1") + if !deliver { + t.Fatal("投过但没跑完的必须重投 —— 跳过它就是静默丢件,发件人收不到回信") + } + if !resumed { + t.Error("必须标记为重投,好让调用方在提示词里说明「上一轮被中断」") + } +} + +// 模拟真实事故:进程 A 注入后被杀,进程 B 起来重读账本。 +func TestSurvivesProcessRestart_ResumesIncomplete(t *testing.T) { + dir := t.TempDir() + + a := openDeliveryLedger(dir) + a.claim("m1") + a.close() // 进程 A 死了,那一轮没跑完 + + b := openDeliveryLedger(dir) + defer b.close() + + deliver, resumed := b.claim("m1") + if !deliver { + t.Fatal("新进程必须重投未完成的邮件") + } + if !resumed { + t.Error("新进程必须知道这是重投(上一个进程注入过一次)") + } +} + +// 反过来:跑完的那封,新进程必须跳过 —— 这正是修复的目标。 +func TestSurvivesProcessRestart_SkipsCompleted(t *testing.T) { + dir := t.TempDir() + + a := openDeliveryLedger(dir) + a.claim("m1") + a.complete("m1") + a.close() + + b := openDeliveryLedger(dir) + defer b.close() + + if deliver, _ := b.claim("m1"); deliver { + t.Error("跨进程也必须跳过已跑完的邮件(这就是那次重复投递的根因)") + } +} + +func TestCompleteIsIdempotent(t *testing.T) { + dir := t.TempDir() + l := openDeliveryLedger(dir) + l.claim("m1") + l.complete("m1") + sizeAfterFirst := fileSize(t, filepath.Join(dir, "delivered.jsonl")) + l.complete("m1") + l.complete("m1") + l.close() + + if got := fileSize(t, filepath.Join(dir, "delivered.jsonl")); got != sizeAfterFirst { + t.Errorf("重复 complete 不该再写盘:%d → %d 字节", sizeAfterFirst, got) + } +} + +// complete 一封从没 claim 过的邮件(理论上不该发生)也不能让状态错乱。 +func TestCompleteWithoutClaim(t *testing.T) { + dir := t.TempDir() + a := openDeliveryLedger(dir) + a.complete("ghost") + a.close() + + b := openDeliveryLedger(dir) + defer b.close() + if deliver, _ := b.claim("ghost"); deliver { + t.Error("已标记完成的邮件不该被投递,即使它没走过 claim") + } +} + +// 崩溃时最后一行可能写残。坏行必须被跳过,而不是让整个账本失效。 +func TestCorruptLineSkipped(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "delivered.jsonl") + + good := `{"m":"m-good","c":true,"t":` + itoa(time.Now().Unix()) + "}\n" + // 半行 + 完全不是 JSON 的一行 + if err := os.WriteFile(path, []byte(good+`{"m":"m-half","c":fal`+"\nnot json at all\n"), 0o600); err != nil { + t.Fatal(err) + } + + l := openDeliveryLedger(dir) + defer l.close() + + if deliver, _ := l.claim("m-good"); deliver { + t.Error("完好的那行必须生效") + } + // 写残的那封退化为「没记录」→ 重投。重投是安全的,跳过才危险。 + deliver, resumed := l.claim("m-half") + if !deliver { + t.Error("坏行对应的邮件应当重投(宁可重复也不能丢)") + } + if resumed { + t.Error("坏行没留下任何有效记录,不该被当成重投") + } +} + +func TestExpiredEntriesDropped(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "delivered.jsonl") + + old := time.Now().Add(-ledgerRetention - time.Hour).Unix() + fresh := time.Now().Unix() + body := `{"m":"m-old","c":true,"t":` + itoa(old) + "}\n" + + `{"m":"m-fresh","c":true,"t":` + itoa(fresh) + "}\n" + if err := os.WriteFile(path, []byte(body), 0o600); err != nil { + t.Fatal(err) + } + + l := openDeliveryLedger(dir) + defer l.close() + + if deliver, _ := l.claim("m-fresh"); deliver { + t.Error("保留期内的记录必须生效") + } + if deliver, _ := l.claim("m-old"); !deliver { + t.Error("超过保留期的记录应被丢弃(否则账本无限增长)") + } +} + +// 过期条目过半时压实,且压实不能弄丢仍然有效的记录。 +func TestCompactPreservesLiveEntries(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "delivered.jsonl") + + old := time.Now().Add(-ledgerRetention - time.Hour).Unix() + fresh := time.Now().Unix() + var sb strings.Builder + for i := 0; i < 10; i++ { + sb.WriteString(`{"m":"old-` + itoa(int64(i)) + `","c":true,"t":` + itoa(old) + "}\n") + } + sb.WriteString(`{"m":"keep","c":true,"t":` + itoa(fresh) + "}\n") + if err := os.WriteFile(path, []byte(sb.String()), 0o600); err != nil { + t.Fatal(err) + } + + before := fileSize(t, path) + l := openDeliveryLedger(dir) + l.close() + + if after := fileSize(t, path); after >= before { + t.Errorf("过期条目过半时应压实:%d → %d 字节", before, after) + } + + // 压实后重开,仍然认得那条有效记录 + l2 := openDeliveryLedger(dir) + defer l2.close() + if deliver, _ := l2.claim("keep"); deliver { + t.Error("压实不该弄丢仍然有效的记录") + } +} + +// 目录不可写时退化为纯内存,而不是让插件起不来。 +func TestUnwritableDirDegradesToMemory(t *testing.T) { + l := openDeliveryLedger("") + defer l.close() + + if deliver, _ := l.claim("m1"); !deliver { + t.Fatal("退化为纯内存后仍要能投递") + } + l.complete("m1") + if deliver, _ := l.claim("m1"); deliver { + t.Error("纯内存模式下同进程内的去重仍要生效") + } +} + +// 并发 claim 同一封:只有一个能拿到「全新」,另一个必须是「重投」。 +// 这正是 SSE 与 catchUp 两个 goroutine 的竞态。 +func TestConcurrentClaimSameMail(t *testing.T) { + l := openDeliveryLedger(t.TempDir()) + defer l.close() + + const n = 20 + type res struct{ deliver, resumed bool } + out := make(chan res, n) + start := make(chan struct{}) + for i := 0; i < n; i++ { + go func() { + <-start + d, r := l.claim("hot") + out <- res{d, r} + }() + } + close(start) + + fresh := 0 + for i := 0; i < n; i++ { + r := <-out + if !r.deliver { + t.Error("未完成的邮件在任何一次 claim 上都该返回要投递") + } + if !r.resumed { + fresh++ + } + } + if fresh != 1 { + t.Errorf("恰好一次 claim 该被判为全新,实际 %d 次 —— 判定与记录必须在同一把锁里", fresh) + } +} + +func TestResumeNoteMentionsInterruption(t *testing.T) { + note := resumeNote("abc-123") + for _, want := range []string{"abc-123", "中断", "同一件事"} { + if !strings.Contains(note, want) { + t.Errorf("重投说明必须包含 %q,实际:%s", want, note) + } + } +} + +func TestShortIDDoesNotPanicOnShortInput(t *testing.T) { + // mail_id 理论上可能短于 8 字节(畸形事件、测试桩)。 + // 日志不该有能力弄死投递协程。 + for _, in := range []string{"", "ab", "12345678", "123456789"} { + got := shortID(in) + if len(got) > 8 { + t.Errorf("shortID(%q) = %q,超过 8 字节", in, got) + } + } +} + +// ─── 辅助 ─── + +func fileSize(t *testing.T, path string) int64 { + t.Helper() + st, err := os.Stat(path) + if err != nil { + t.Fatalf("stat %s: %v", path, err) + } + return st.Size() +} + +func itoa(n int64) string { + if n == 0 { + return "0" + } + neg := n < 0 + if neg { + n = -n + } + var b []byte + for n > 0 { + b = append([]byte{byte('0' + n%10)}, b...) + n /= 10 + } + if neg { + return "-" + string(b) + } + return string(b) +} diff --git a/plugins/homeagent-mail-bridge/plugin.go b/plugins/homeagent-mail-bridge/plugin.go index 5fe283a..199ad51 100644 --- a/plugins/homeagent-mail-bridge/plugin.go +++ b/plugins/homeagent-mail-bridge/plugin.go @@ -84,8 +84,17 @@ type Plugin struct { // B-7.3 邮件级去重:SSE 重放会重发同一批事件,没有这层去重 // 每封邮件会被注入 agent 两遍。契约 B-7.3 要求:每封只注入一次。 + // + // 这只挡得住**本进程内**的重复。跨进程(homed 重启、插件子进程被换) + // 靠 ledger —— 它落盘,且区分「投过」与「跑完」。 deliveredMails map[string]bool + // 跨进程投递账本(见 ledger.go)。 + // + // 它与 deliveredMails 不是重复:后者是同一进程内 SSE 重放的快速路径, + // 前者回答的是「上一个进程有没有已经把这封跑完」。 + ledger *deliveryLedger + // 单调递增的 last-seen-ID:被重放的旧事件不会让它回退。 // 原来直接赋值(p.lastEventID = eid),Gateway 重放时发旧 ID, // 于是 lastEventID 从 123 退回 116 → 下次重连又报 116 → 又重放。 @@ -185,6 +194,20 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { // B-1.1 解析密钥(环境变量 → 本地文件 → 生成) p.key, p.keyFile = resolveKey() + // 跨进程投递账本:子进程被换时(homed 重启 / 插件崩溃自动重启) + // 内存里的 deliveredMails 全丢,只靠它防不住重复注入。 + // + // DataDir() 是 SDK 保证存在的插件专属目录;拿不到时退回 key 文件所在目录 + // (那个目录本来就要能写)。 + dataDir := "" + if sett := s.Settings(); sett != nil { + dataDir = sett.DataDir() + } + if strings.TrimSpace(dataDir) == "" { + dataDir = filepath.Dir(p.keyFile) + } + p.ledger = openDeliveryLedger(dataDir) + // ─── 注册工具 ─── // registerTool 包一层只为计数:日志里的工具数必须与实际注册数一致。 @@ -413,7 +436,14 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { } func (p *Plugin) Stop() error { - p.stopOnce.Do(func() { close(p.stopCh) }) + p.stopOnce.Do(func() { + close(p.stopCh) + // 关账本句柄。每行写入都 Sync 过,所以不关也不丢数据 —— + // 关只是为了不把 fd 泄给下一个插件实例。 + if p.ledger != nil { + p.ledger.close() + } + }) return nil } @@ -532,8 +562,31 @@ func (p *Plugin) catchUp(pending int) { continue } - // 构造注入消息(与 handleNewMail 一致) - prompt := fmt.Sprintf( + // 跨进程去重:这次重启前那个进程可能已经把这封跑完了。 + // + // 这才是那次事故的真正修法:死掉的那个进程已经注入过一次, + // 而 deliveredMails 随它一起消失了。ledger 落盘,能说出区别: + // - 已跑完 → 真的跳过 + // - 投过但未跑完(上一轮被中断)→ 仍然重投,但带上说明 + // 后一条很要紧:那一轮被中断意味着发件人没收到回信,跳过它就是静默丢件。 + deliver, resumed := p.ledger.claim(m.MailID) + if !deliver { + log.Printf("[homeagent-mail-bridge] 补投跳过 %s:已在之前的进程里处理完毕", shortID(m.MailID)) + continue + } + if resumed { + log.Printf("[homeagent-mail-bridge] 补投 %s(上一轮被中断,带说明重投)", shortID(m.MailID)) + } + + // 注入消息。与 handleNewMail 那份的差异只在一句措辞上(这里不说 + // 「把本轮工作做完」)—— 那个差异正好是上次定位重复投递的线索: + // 两段提示词同时出现在上下文里,一看措辞就知道一段来自 SSE、 + // 一段来自补投。 + prefix := "" + if resumed { + prefix = resumeNote(m.MailID) + } + prompt := prefix + fmt.Sprintf( "你收到一封新邮件(AgentMail)。\n\n"+ "发件人:%s\n主题:%s\n邮件 ID:%s\n身份:你是 %s\n\n"+ "请先调用 read_inbox 读取完整正文,然后处理其中的请求。\n\n"+ @@ -545,8 +598,9 @@ func (p *Plugin) catchUp(pending int) { reply := p.sdk.InjectInputSync(p.name, p.name, prompt) if reply == "" { - // B-6:模型没回,发一封告知 + // B-6:模型没回,发一封告知。发出去就算处理完(理由同 handleNewMail)。 p.sendFailureReply(m.FromName, m.Subject, m.MailID, "模型未产生回复") + p.ledger.complete(m.MailID) continue } // B-5.3:检查模型是否已经自己发过信 @@ -557,11 +611,17 @@ func (p *Plugin) catchUp(pending int) { if sent { // 模型已经在这一轮里自己回了这封信,不再重复 relay + p.ledger.complete(m.MailID) continue } // B-5.2:自动回信带 relay:"summary" —— 搬运不算模型自主发信,不扣配额 - p.sendMailRelay(m.FromName, "Re: "+m.Subject, reply, m.MailID, "homeagent:"+m.MailID) + if err := p.sendMailRelay(m.FromName, "Re: "+m.Subject, reply, m.MailID, rk); err != nil { + // 回信没发出去 —— 不标完成,下次重启重试。 + log.Printf("[homeagent-mail-bridge] 补投回信失败(不标完成): %v", err) + continue + } + p.ledger.complete(m.MailID) } } @@ -707,12 +767,22 @@ func (p *Plugin) parseSSELine(line string) { p.deliveredMails[evt.MailID] = true p.sseMu.Unlock() + // 跨进程去重:上一个插件子进程可能已经把这封跑完了。 + // deliveredMails 只在本进程内有效,homed 重启会把它清空 —— + // 实测过一次两段几乎相同的通知堆在模型上下文里(一段来自这里的 + // SSE 路径,一段来自重启后的 catchUp)。 + deliver, resumed := p.ledger.claim(evt.MailID) + if !deliver { + log.Printf("[homeagent-mail-bridge] 邮件 %s 已在之前的进程里处理完毕,跳过", shortID(evt.MailID)) + return + } + // InjectInputSync 会阻塞几十秒(查日志、调工具、转发 QQ), // 而它跑在 readSSE 的读循环里 —— 循环卡住期间 SSE 事件积压在 // TCP 缓冲区,卡到超时断线重连后 Gateway 全部重放一遍。 // 把处理丢到独立 goroutine:parseSSELine 立刻返回,读循环继续。 // homeagent 是单事件循环,InjectInputSync 自己会排队。 - go p.handleNewMail(evt) + go p.handleNewMail(evt, resumed) } } @@ -822,8 +892,18 @@ func (p *Plugin) handleNewMail(evt struct { Workspace string `json:"to_workspace"` Alias string `json:"session_alias"` ReplyAddr string `json:"reply_address"` -}) { - prompt := fmt.Sprintf( +}, resumed bool) { + // resumed = 上一个进程注入过这封但那一轮被中断了。 + // + // 必须把这件事告诉模型:不说的话它在上下文里看到两段几乎相同的指令, + // 会以为人重复交代了一遍,于是可能把同一件事做两次。 + prefix := "" + if resumed { + prefix = resumeNote(evt.MailID) + log.Printf("[homeagent-mail-bridge] 邮件 %s 重投(上一轮被中断)", shortID(evt.MailID)) + } + + prompt := prefix + fmt.Sprintf( "你收到一封新邮件(AgentMail)。\n\n"+ "发件人:%s\n主题:%s\n邮件 ID:%s\n身份:你是 %s\n\n"+ "请先调用 read_inbox 读取完整正文,然后处理其中的请求。\n\n"+ @@ -839,8 +919,12 @@ func (p *Plugin) handleNewMail(evt struct { // B-6:模型没回(空 = turn/end 信号 kind=error,或模型没说话) if reply == "" { log.Printf("[homeagent-mail-bridge] 邮件 %s(来自 %s:%s)agent 无回复,发失败通知", - evt.MailID[:8], evt.FromName, evt.Subject) + shortID(evt.MailID), evt.FromName, evt.Subject) p.sendFailureReply(evt.FromName, evt.Subject, evt.MailID, "模型未产生回复") + // 失败通知发出去了就算**处理完**:发件人得到了一个明确的交代。 + // 不标的话下次重启会把同一封再投一遍 —— 而模型上一次就没回, + // 重投只会再发一封相同的失败通知。 + p.ledger.complete(evt.MailID) return } @@ -855,15 +939,19 @@ func (p *Plugin) handleNewMail(evt struct { if sent { // 模型已经在这一轮里自己回了这封信,让位 - log.Printf("[homeagent-mail-bridge] 邮件 %s 模型已自行回复,跳过自动 relay", evt.MailID[:8]) + log.Printf("[homeagent-mail-bridge] 邮件 %s 模型已自行回复,跳过自动 relay", shortID(evt.MailID)) + p.ledger.complete(evt.MailID) return } // B-5.2:自动回信带 relay:"summary" + relay_key if err := p.sendMailRelay(evt.FromName, "Re: "+evt.Subject, reply, evt.MailID, rk); err != nil { - log.Printf("[homeagent-mail-bridge] 自动回信失败: %v", err) + // 回信没发出去 —— **不标完成**,让下次重启能重试。 + // 发件人至今一个字都没收到,这时标「已完成」就是静默丢件。 + log.Printf("[homeagent-mail-bridge] 自动回信失败(不标完成,下次会重试): %v", err) } else { log.Printf("[homeagent-mail-bridge] 已自动回信给 %s(%d 字)", evt.FromName, len(reply)) + p.ledger.complete(evt.MailID) } // 清理过期的 explicitSends 记录 diff --git a/plugins/pi-mail-bridge/src/index.mjs b/plugins/pi-mail-bridge/src/index.mjs index 44f0116..1d61f9b 100644 --- a/plugins/pi-mail-bridge/src/index.mjs +++ b/plugins/pi-mail-bridge/src/index.mjs @@ -4,18 +4,32 @@ * * 形态是**常驻守护进程**,不是 pi 扩展。原因见 src/session-pool.mjs 顶部: * 扩展被加载进一条已存在的会话,cwd 由启动 pi 的人决定;而 B-3.1 要求每封邮件的 - * to_workspace 成为会话 cwd。桥用 SDK 的 createAgentSession 按邮件起会话, - * 一个进程里并存多条不同 cwd 的会话(实测可行)。 + * to_workspace 成为会话 cwd。 + * + * # 进程结构 + * + * 主进程**只做 I/O 与调度**:SSE、心跳、去重、把邮件派给子进程。模型工作全部 + * 下到 `src/worker.mjs`(一封邮件一个进程,跑完就退),由 `src/pool.mjs` 调度。 + * + * 这不是为了并行度,是为了**不阻塞事件循环**。pi 的会话装载是同步的: + * `SessionManager.open()` 走 `openSync` + `readSync` 循环把整个 `.jsonl` 读进内存 + * 并逐行 JSON.parse。实测本机最大那条会话 23MB,`open` 一次阻塞事件循环 118ms; + * 模型跑起来之后 SDK 内部还有更多同步工作。原来这些都在主线程上 —— SSE 读循环 + * 在那期间完全停住,后续邮件卡在 TCP 缓冲区,久到 Gateway 认为连接死了, + * 重连又触发重放。实测同样的活在 fork 出的子进程里跑,主进程阻塞 0ms。 + * + * 权限询问期间的挂起也随之只影响那一个 worker:原来 `await new Promise(...)` + * 等人做决定,整座桥在那段时间不再收信。 * * 契约实现对照(docs/PLUGIN-CONTRACT.md): * B-1 启动 → main() * B-2 心跳 → beat(),30 秒 - * B-3 new_mail → deliverMail() + * B-3 new_mail → pool.submit()(投递本体在 worker.mjs) * B-4 决策 → handlePermissionDecision() - * B-5 转发 → relaySummary(),挂在 agent_end 上 - * B-6 失败回信 → deliverMail() 末尾的 renderFailureReport + * B-5 转发 → worker.mjs 的 relaySummary() + * B-6 失败回信 → worker.mjs 末尾的 renderFailureReport * B-7 补拉 → catchUp() - * B-8 权限 → permissionExtension() 的 tool_call 钩子 + * B-8 权限 → worker.mjs 的 permissionExtension() * B-9 关停 → shutdown() */ @@ -26,16 +40,11 @@ import { ModelRuntime } from '@earendil-works/pi-coding-agent'; import { GatewayClient, readLocalKey, generateLocalKey, saveConfig } from './gateway.mjs'; import { createMailTools } from './tools.mjs'; -import { openSession, runTurn } from './session-pool.mjs'; -import { buildMailPrompt, lastAssistantText, replySubject, relayKeyFor, describeError } from './turn.mjs'; -import { planNamingSync, planWriteBack } from './naming.mjs'; -import { resolveWorkspaceCwd, ensureCwd } from '../lib/workspace.js'; -import { modelAttemptOrder, renderFailureReport, snapshotPiModels } from '../lib/model-scope.js'; +import { createWorkerPool } from './pool.mjs'; +import { describeError } from './turn.mjs'; +import { snapshotPiModels } from '../lib/model-scope.js'; import { snapshotPiSessions } from '../lib/session-snapshot.js'; import { selectCatchup } from '../lib/catchup.js'; -import { explicitSends, shouldSkipAutoRelay } from '../lib/relay-dedup.js'; -import { adoptedSessionID, adoptMissingMessage } from '../lib/adopt.js'; -import { createGrantStore, isApproval } from '../lib/permission-grants.js'; // ─── 配置 ─── @@ -44,7 +53,34 @@ const AGENT_NAME = process.env.AGENTMAIL_AGENT_NAME || 'pi'; const AGENT_SECRET = process.env.AGENTMAIL_AGENT_SECRET || ''; const REPLY_PROVIDER = process.env.AGENTMAIL_REPLY_PROVIDER || ''; const REPLY_MODEL = process.env.AGENTMAIL_REPLY_MODEL || ''; -const TURN_TIMEOUT_MS = Number(process.env.AGENTMAIL_TURN_TIMEOUT_MS || 60_000); + +/** + * 单轮超时。 + * + * 从 60 秒放宽到 10 分钟:60 秒那个数字是「主进程要腾出手来收下一封邮件」的 + * 产物 —— 超时按成功返回,好让 deliverMail 早点结束。worker 没有这个理由, + * 它只为这封邮件活着,等真结论更准。带工具调用的一轮跑几分钟很正常, + * 60 秒返回会让转发落在一个还没说完的结论上。 + */ +const TURN_TIMEOUT_MS = Number(process.env.AGENTMAIL_TURN_TIMEOUT_MS || 600_000); + +/** + * 并发上限。 + * + * 每个 worker 约 140MB RSS(实测),并且每个都会对上游 provider 发请求。 + * 3 是内存与吞吐的折中;同一条会话无论如何都是串行的(见 pool.mjs)。 + */ +const MAX_WORKERS = Number(process.env.AGENTMAIL_MAX_WORKERS || 3); + +/** + * worker 硬超时。 + * + * 比轮次超时留出余量:正常情况下 worker 自己会在轮次超时后收尾退出, + * 这个数字兜的是「连收尾都没做」(进程卡死、权限等不到决策而决策事件也丢了)。 + * 到点 SIGKILL —— 否则那条会话的后续邮件永远排队。 + */ +const WORKER_MAX_MS = Number(process.env.AGENTMAIL_WORKER_MAX_MS || TURN_TIMEOUT_MS + 120_000); + const LOCK_FILE = join(process.env.AGENTMAIL_CONFIG_DIR || join(homedir(), '.agentmail'), 'pi-bridge.lock'); /** 日志一律 console.error:它一定进 journalctl(契约 9.8)。 */ @@ -52,23 +88,16 @@ const log = (...args) => console.error('[pi-mail-bridge]', ...args); // ─── 进程内状态 ─── // -// 全部只在内存,重启即丢 —— 这是契约第六节列明的已知取舍。 -// 要持久化的话该落在 pi 的会话元数据里,而不是桥自己的文件。 +// 主进程只留「调度需要的」那几样,全部在内存(重启即丢,契约第六节的已知取舍)。 +// 会话映射、权限授权、命名指纹都下沉到 pool 里按邮件会话存 —— 主进程不再持有 +// AgentSession 对象(那东西跨不了进程边界)。 -const sessions = new Map(); // agentmail session_id -> { session, sessionManager, cwd } -const reverseMap = new Map(); // pi session id -> agentmail session_id -const mailDriven = new Set(); // pi session id -const mailContexts = new Map(); // agentmail session_id -> { replyTo, subject, mailID } -const relayedSummaries = new Map(); // pi session id -> 已转发过的 relay_key -const syncedNames = new Map(); // pi session id -> 上次提交给 Gateway 的名字 -const pendingPermissions = new Map(); // relay_key -> { resolve, piSessionId } -// 人点过「一直同意」的 (会话, 工具)。作用域与清理语义见 lib/permission-grants.js。 -const permissionGrants = createGrantStore(); const deliveredMails = new Set(); // 已投过的 mail_id(SSE 与补拉共用,B-7.3) let allowedModels = []; let modelRuntime = null; let client = null; +let pool = null; let heartbeatTimer = null; let shuttingDown = false; @@ -113,648 +142,31 @@ function releaseLock() { } catch { /* 已经没了 */ } } -// ─── 权限钩子(B-8)─── - -/** - * 内联 pi 扩展:把 pi 拦下的危险工具调用转成一封邮件问人。 - * - * 这是 `I-1` 最直接的体现 —— 被平台真正拦下的那一次才是事实, - * 不依赖模型「记得」调 request_permission(它会忘,也会在不需要时乱调)。 - * - * pi 的 `tool_call` 钩子**可以 await**(C-9 实测成立:处理器里 await 300ms - * 再返回 {block:true},pi 会等),所以这里能真的等人做决定, - * 不必走「先拒一次再重试」的退化路径。 - * - * @param {string} piSessionIdRef 用一个 getter 拿会话 id:扩展工厂在 - * createAgentSession **内部**被调用,那时 session 对象还没返回给桥。 - */ -function permissionExtension(getMailContext) { - // pi 默认放行内建工具;桥只拦真正有副作用的那几个。 - // read/grep/ls 之类不拦:每一步都问人会让 Agent 什么也做不成, - // 而人也会很快开始无脑点同意(那比不问更危险)。 - const GUARDED = new Set(['bash', 'write', 'edit']); - - return (pi) => { - pi.on('tool_call', async (event, ctx) => { - if (!GUARDED.has(event.toolName)) return; - - const piSessionId = ctx?.sessionManager?.getSessionId?.() || ''; - const mailSessionId = reverseMap.get(piSessionId); - // 不是邮件驱动的会话 → 让位给 pi 自己的本地 UI(B-8.2)。 - // 占着钩子不放会让人在 TUI 里干活时每一步都卡住等邮件。 - if (!mailSessionId) return; - - // 人对这条会话的这个工具点过「一直同意」→ 直接放行,不再发邮件。 - // 这一步必须在 POST 之前:否则每条命令都生成一封邮件,人点过的 - // 「一直同意」形同虚设(实测同一条会话被问了 15 次 bash)。 - if (permissionGrants.isGranted(piSessionId, event.toolName)) { - return; - } - - // relay_key 用 pi 给的 toolCallId(B-8.1):服务端会随决策事件回传它, - // 桥重启丢了 pendingPermissions 也能对上(B-4.2)。自造随机 id 做不到。 - const relayKey = `${piSessionId}:${event.toolCallId}`; - const ctxInfo = getMailContext(mailSessionId); - - try { - // **不传 `to`**(这里曾经传 `ctxInfo.replyTo`,那是个死锁 bug)。 - // - // replyTo 是来信人的名字,而来信人可能是另一个 Agent —— Agent 把任务 - // 分派给自己的另一条会话时(pi→pi),权限邮件就发给了 `pi` 自己。 - // 后果是死锁而不是报错:Agent 不可能在 Web 界面上点「同意」, - // 服务端的 SendToUser 又投进一个不存在的用户通道(没有任何人被提醒), - // 于是下面那个 await 永不 resolve —— 会话永久挂死,没有超时、没有日志。 - // - // 决策人交给服务端定:它按 会话 owner → 线索里最近的人类 → 无人可问则 - // 409 的顺序解析,那是唯一能看到整条线索的地方。插件只有本地那点上下文, - // 猜不出「这条 Agent 链最初是谁派的活」。 - await client.post('/permission/request', { - question: `是否允许执行 ${event.toolName}?`, - options: ['同意', '一直同意', '拒绝'], - // 带上触发这次询问的来信(B-8.4)。决策人未必是这条会话的参与者 —— - // Agent 转派出来的会话,人从没见过它,只给一句「是否允许执行 bash」 - // 是无从判断的:得知道这活是谁派的、为的什么事。 - context: [ - describeToolCall(event), - ctxInfo?.subject ? `\n触发任务:${ctxInfo.subject}` : '', - ctxInfo?.replyTo ? `任务来自:${ctxInfo.replyTo}` : '', - ].filter(Boolean).join('\n'), - session_id: mailSessionId, - relay_key: relayKey, - }); - } catch (e) { - // 409 = 服务端已判定这条任务链上没有人类,永远不会有人来点头。 - // - // 不能「让位给本地 UI」:邮件驱动的会话没有 TUI,让位之后 pi 按默认 - // 策略处置,而默认策略在没有交互端时就是等 —— 又一次无声挂死。 - // - // 直接 block 并把服务端的建议原文当作 reason:模型从工具报错里 - // 看到「这条链上没人可问,换不需要权限的方式」才能自己改道, - // 而挂死时它连重试的机会都没有。 - if (e?.status === 409) { - const b = e.body || {}; - const reason = [ - b.error || '权限询问无法送达:这条任务链上没有人类用户', - b.detail || '', - b.suggestion || '', - ].filter(Boolean).join('\n'); - log(`权限询问无人可投,当场拒绝 ${relayKey}:${b.error || ''}`); - return { block: true, reason }; - } - - // 其余失败(网络抖动、Gateway 重启)是暂时的 → 让位给 pi 本地 UI(B-8.2)。 - // 返回 undefined 表示「这个钩子不表态」,pi 会走它自己的批准流程。 - log(`权限转发失败,让位给本地决策: ${describeError(e)}`); - return; - } - - log(`权限询问已发出(${event.toolName},key=${relayKey}),等待决策…`); - const decision = await new Promise((resolve) => { - pendingPermissions.set(relayKey, { resolve, piSessionId }); - }); - - // fail closed(B-9.2 / N-9):只有明确的同意才放行。 - // 关停时 shutdown() 会用 'shutdown' 唤醒所有等待者,落到这里的 else。 - if (isApproval(decision)) { - // 「一直同意」要真的记住,否则这个选项是在骗人:人点了它, - // 下一条命令照样来一封邮件。grant() 内部只认精确的 always 文本 —— - // 「同意」是单次授权,把它当 always 会放行人没看过的后续命令。 - if (permissionGrants.grant(piSessionId, event.toolName, decision)) { - log(`本会话的 ${event.toolName} 已获「一直同意」,后续不再询问`); - } - log(`权限 ${relayKey} 获批(${decision}),放行 ${event.toolName}`); - return; - } - return { block: true, reason: `用户${decision === 'shutdown' ? '未及决策(桥已关停)' : `拒绝了这次 ${event.toolName} 调用`}` }; - }); - }; -} - -/** 把一次工具调用摘要成人能判断的文本(B-8.4)。 */ -function describeToolCall(event) { - const input = event?.input ?? {}; - if (event.toolName === 'bash') { - return `命令:\n${String(input.command ?? '').slice(0, 800)}`; - } - if (event.toolName === 'write' || event.toolName === 'edit') { - return `文件:${input.file_path ?? input.path ?? '(未给出)'}`; - } - return JSON.stringify(input).slice(0, 800); -} - -// ─── 会话解析(B-3)─── - -/** - * 没有可用 `to_workspace` 时的兜底目录。 - * - * 与 DSH 的 `mailSessionFallback` 同构,但目录名是 `.pi`:那个函数在 - * lib/ 下(三平台逐字节相同),写死了 `.dsh`,不能为 pi 改。 - * 让 pi 的会话落进 `~/.dsh/` 会让人以为是 DSH 在干活。 - */ -function piMailFallback(sessionKey) { - return join(homedir(), '.pi', 'mail-sessions', String(sessionKey || 'default')); -} - -/** - * 接管过的 pi 会话(pi session id)。用完即释放,见 releaseAdopted。 - * - * 与 `sessions` 的区别:那里存的是桥自己起的、长期持有的会话;这里是 - * 「借用一下磁盘上人家的会话」,一轮结束就还回去。 - */ -const adopted = new Set(); - -/** 接管会话的兜底释放计时器(pi session id -> Timeout)。 */ -const adoptTimers = new Map(); - -/** - * 接管会话最长持有多久。 - * - * 取轮次超时的两倍:`runTurn` 60 秒就按成功返回(长任务很正常,判成失败会 - * 换模型重跑一遍),但会话仍在跑。正常结束走 agent_end 提前释放,这个数字 - * 只兜「事件永远不来」的底。 - */ -const ADOPT_MAX_HOLD_MS = TURN_TIMEOUT_MS * 2; - -/** - * 接管一条磁盘上已经存在的 pi 会话,把这封邮件投进去。 - * - * # 为什么必须**短暂持有** - * - * pi 没有任何锁机制,它假定「一个文件一个持有者」。活着的 SessionManager - * 不 watch 文件:外部(TUI)追加的行它看不见,之后它自己的写入算出的 parentId - * 指向一个对方不知道的 entry —— 文件不会坏(写入是纯 append),但会话树分叉。 - * - * 所以这里 open → 跑一轮 → 丢弃,**不放进 sessions 长期缓存**。下一封邮件 - * 再来时重新 open,那一次读到的就是 TUI 期间写的全部内容。 - * - * 窗口是一轮对话的时长。人正好在这期间也在 TUI 里发消息仍会分叉,但那需要 - * 两边同时动手,且后果是历史看起来少一段,不是数据损坏。 - * - * # 为什么不校验「TUI 是否正开着这条会话」 - * - * pi 不提供这个信息(没有 lockfile、没有 pid 记录)。能做的只有猜 mtime, - * 而任何阈值都是猜。与其用一个猜出来的数字拒掉合法投递,不如让窗口尽量短。 - */ -async function adoptSession(platformID, data, mailTools) { - const mailSessionID = data.session_id; - const { SessionManager } = await import('@earendil-works/pi-coding-agent'); - - // listAll 而不是 list(cwd):桥的进程 cwd 与会话 cwd 无关。 - const all = await SessionManager.listAll(); - const info = all.find((e) => e?.id === platformID); - if (!info?.path) { - // 镜像是快照,可以过期:那条会话可能已经被删了。 - // **不能**退回「新建一条」—— 那会让人在 TUI 里看不到这封邮件带来的对话, - // 而那正是接管的目的(N-8 同理:静默改语义比报错糟)。 - throw new Error(adoptMissingMessage(platformID, '磁盘上已无这个会话文件')); - } - - // cwd 取会话自己的(SessionInfo.cwd 来自持久化 header)。 - // 老会话的 cwd 是空串,那种情况退回地址里的 path 位。 - const { cwd } = resolveWorkspaceCwd( - info.cwd || data.to_workspace, piMailFallback(mailSessionID)); - - const opened = await openSession({ - cwd, - modelRuntime, - customTools: mailTools, - extension: permissionExtension((id) => mailContexts.get(id)), - sessionFile: info.path, - }); - for (const d of opened.diagnostics) { - log(`扩展诊断: ${d?.message ?? JSON.stringify(d)}`); - } - - const piSessionId = opened.session.sessionId; - const entry = { session: opened.session, sessionManager: opened.sessionManager, cwd }; - - if (mailSessionID) { - sessions.set(mailSessionID, entry); - reverseMap.set(piSessionId, mailSessionID); - // 加进 mailDriven:接管之后这条会话**开始**参与邮件往来,轮次结束要把 - // 总结转回发件人。不加的话邮件投进去了却永远没有回音。 - mailDriven.add(piSessionId); - adopted.add(piSessionId); - } - - // 兜底释放:`agent_end` 不来就永远握着这个文件,而握着它的期间 TUI 那边 - // 的写入对我们不可见 —— 正是要避免的分叉窗口。会话跑挂、事件丢失、 - // 模型一直不结束都属于这种情形。 - // - // 时长取轮次超时的两倍:runTurn 自己 60 秒就按成功返回了(长任务很正常), - // 那之后会话仍在跑,正常结束时 agent_end 会照常触发并提前释放。 - const safety = setTimeout(() => { - if (!adopted.has(piSessionId)) return; - log(`接管会话 ${piSessionId} 超过 ${ADOPT_MAX_HOLD_MS / 1000}s 未结束,强制释放`); - releaseAdopted(piSessionId, mailSessionID); - }, ADOPT_MAX_HOLD_MS); - if (typeof safety.unref === 'function') safety.unref(); - adoptTimers.set(piSessionId, safety); - - // 只挂 agent_end,不挂 session_info_changed:改名同步会把 Gateway 侧的别名 - // 覆盖成 pi 的标题,而接管会话的别名是人从补全里选的那个 slug —— - // 改掉会让他找不到自己刚发的信。 - opened.session.subscribe((event) => { - if (event?.type !== 'agent_end') return; - if (event.willRetry) return; - // 转发完再释放:relaySummary 要读 sessions 里的 entry。 - relaySummary(piSessionId) - .catch((e) => log(`自动转发失败: ${describeError(e)}`)) - .finally(() => { - // **还在跑就不能释放。** - // - // 同一条会话可能已经排了下一封邮件:runTurn 在 isStreaming 时走 - // `streamingBehavior: 'followUp'`,那封信排在当轮之后。此时 dispose - // 会把排着的那一轮一起杀掉 —— 发件人只看到信发出去后再无音讯。 - // 排着的那轮结束时会再触发一次 agent_end,由它来释放。 - if (opened.session.isStreaming) { - log(`接管会话 ${piSessionId} 仍有排队轮次,暂不释放`); - return; - } - releaseAdopted(piSessionId, mailSessionID); - }); - }); - - log(`接管 pi 会话 ${piSessionId}(cwd=${cwd},文件 ${info.path})`); - // reused: true —— 这条会话有历史,提示词不该重新自我介绍, - // 且 deliverMail 的续谈支不做模型降级(换模型要换会话,会丢掉整条上下文)。 - return { ...entry, reused: true }; -} - -/** - * 还回一条接管来的会话:dispose + 清缓存。 - * - * mailDriven 不清:它同时喂给心跳快照的 mail_driven 标记,那条平台会话 - * 确实已经在邮件往来里了。reverseMap 也不清 —— 留着让迟到的事件能找到线索, - * 而 relaySummary 在 entry 缺失时会自己早退。 - */ -function releaseAdopted(piSessionId, mailSessionID) { - if (!adopted.has(piSessionId)) return; - adopted.delete(piSessionId); - const timer = adoptTimers.get(piSessionId); - if (timer) { - clearTimeout(timer); - adoptTimers.delete(piSessionId); - } - const entry = mailSessionID ? sessions.get(mailSessionID) : undefined; - if (entry?.session?.sessionId === piSessionId) { - sessions.delete(mailSessionID); - } - try { - entry?.session?.dispose?.(); - } catch (e) { - log(`释放接管会话失败(不影响后续): ${describeError(e)}`); - } - log(`释放接管会话 ${piSessionId}(文件已交还,下一封邮件重新打开)`); -} - -/** - * 找到(或建立)这封邮件该落进的 pi 会话。 - * - * Gateway 已经按三维地址的 session 位做完了「复用默认 / 新建 / 具名必须存在」 - * 的判定,推来的 session_id 就是判定结果 —— 桥只负责忠实映射, - * 不自己决定开不开新会话(N-8:404 后自动改用 .new 是禁止的)。 - */ -async function resolveSession(data, mailTools) { - const mailSessionID = data.session_id; - const bound = mailSessionID ? sessions.get(mailSessionID) : undefined; - if (bound) return { ...bound, reused: true }; - - // 服务端说这条邮件会话**接管了平台上已经存在的那条会话**(人在 TUI 里开的 - // 那种)—— 投进它而不是新建。TUI 与邮箱是同一个 Agent 的两个入口。 - const adoptedID = adoptedSessionID(data); - if (adoptedID) { - return await adoptSession(adoptedID, data, mailTools); - } - - // cwd 取寻址里的 path 位(B-3.1)。校验走共用模块:目录不存在时**不创建** - // (N-2:笔误会在磁盘上落下真目录,而 Agent 在里面一无所获),拒绝相对路径(N-3)。 - // - // 兜底用 `~/.pi/mail-sessions/<会话>` 而不是共用模块里的 mailSessionFallback —— - // 后者写死了 `.dsh` 目录名(那是 DSH 的家),pi 的会话落进去会让人以为 - // DSH 在干活。lib/ 里的函数三平台逐字节相同,不能为 pi 改它。 - const { cwd, grouped } = resolveWorkspaceCwd(data.to_workspace, piMailFallback(mailSessionID)); - if (!grouped && data.to_workspace) { - log(`工作目录 ${data.to_workspace} 不可用,回退到 ${cwd}`); - } - ensureCwd(cwd, grouped); - - const opened = await openSession({ - cwd, - modelRuntime, - customTools: mailTools, - extension: permissionExtension((id) => mailContexts.get(id)), - }); - for (const d of opened.diagnostics) { - log(`扩展诊断: ${d?.message ?? JSON.stringify(d)}`); - } - - const piSessionId = opened.session.sessionId; - const entry = { session: opened.session, sessionManager: opened.sessionManager, cwd }; - - if (mailSessionID) { - sessions.set(mailSessionID, entry); - reverseMap.set(piSessionId, mailSessionID); - mailDriven.add(piSessionId); - } - - // 一轮结束就转发总结(B-5)。挂 agent_end 而不是 message_end: - // 后者在流式生成中反复触发,转出去的是半截话。 - // subscribe 收的是一个普通函数(AgentSessionEventListener),不是 {onEvent}。 - opened.session.subscribe((event) => { - if (event?.type === 'agent_end') { - // willRetry 为真表示 pi 自己要重试(auto_retry),这一轮还没定论 —— 不转。 - if (event.willRetry) return; - relaySummary(piSessionId).catch((e) => log(`自动转发失败: ${describeError(e)}`)); - } - // pi 侧改名(pi-web 生成标题、人在 TUI 里 /name)→ 同步给 Gateway - if (event?.type === 'session_info_changed') { - syncNaming(piSessionId, event.name).catch((e) => log(`命名同步失败: ${describeError(e)}`)); - } - }); - - log(`新建 pi 会话 ${piSessionId}(cwd=${cwd})`); - return { ...entry, reused: false }; -} - -// ─── 命名一致(C-11 / W-7)─── - -/** - * pi 的名字 → Gateway → 定稿别名回写进 pi。 - * - * 完整推理见 src/naming.mjs 顶部。这里只是把那套决策接上 I/O。 - */ -async function syncNaming(piSessionId, platformName) { - const mailSessionID = reverseMap.get(piSessionId); - if (!mailSessionID) return; // 不是邮件驱动的会话,不碰 - - const plan = planNamingSync({ - platformName, - mailSubject: mailContexts.get(mailSessionID)?.subject, - lastSynced: syncedNames.get(piSessionId), - }); - if (plan.skip) return; - - // 先记下指纹再发请求:响应回来时 setSessionName 会再次触发 - // session_info_changed,这一步是防自激循环的关键。 - syncedNames.set(piSessionId, plan.signature); - - const res = await client.post(`/sessions/${mailSessionID}/sync`, { - alias: plan.alias, - title: plan.title, - }); - - const entry = sessions.get(mailSessionID); - const back = planWriteBack({ - finalAlias: res?.alias, - currentPiName: entry?.session?.sessionName, - }); - log(`命名同步 ${piSessionId}: alias=${res?.alias || '(未变)'} 来源=${plan.source}`); - - if (back.write && entry?.session) { - // 顺序要紧:先更新指纹,再改名。 - // - // setSessionName **同步**触发 session_info_changed(实测),于是本函数会在 - // 这一行里被重入。指纹在改名之后才更新的话,重入那次看到的还是旧指纹, - // 于是又打一次 sync —— 每条会话两次请求,内容完全相同。 - // - // 记的是「把定稿别名当作平台名字」会算出的指纹:重入那次的 platformName - // 正是 back.name,来源判定成 platform,算出来的就是这个值。 - syncedNames.set(piSessionId, `platform:${back.name}|${back.name}`); - // 只用 setSessionName(走 pi 自己的写入路径)。绝不自己拼路径写会话文件: - // 首条 assistant 消息落盘前文件还不存在,pi 首次落盘用 openSync(file,"wx"), - // 抢先创建会让它抛 EEXIST(实测)。 - entry.session.setSessionName(back.name); - log(`别名回写 pi:${back.name}(${back.reason})`); - } -} - -// ─── 自动转发(B-5)─── - -async function relaySummary(piSessionId) { - const mailSessionID = reverseMap.get(piSessionId); - if (!mailSessionID) return; - // 只对邮件驱动的会话转发(B-5.5):人在 pi 里正常干活时不该往邮箱灌总结 - if (!mailDriven.has(piSessionId)) return; - - const entry = sessions.get(mailSessionID); - if (!entry) return; - - // 一轮结束是命名的自然时机(C-11 / D-5)。 - // - // 这一步不能只挂在 session_info_changed 上:桥用 SDK 起的会话**永远不会** - // 触发那个事件 —— pi 的标题生成器在 pi-web 里,不在内核里,SDK 路径上没有它。 - // 只等事件的话别名永远是空的,于是 `name@path.<别名>` 续谈无从下手 - // (实测过:第一封邮件跑通了,sessions.session_alias 仍是空串)。 - // - // 放在转发**之前**:回信里会带上会话别名,收件人看到的第一封回信就能用它续谈。 - // - // **接管来的会话跳过这一步。** - // - // 它的名字是人在 TUI 里定的,也是他从补全里选中的那个 slug。同步会双向改坏它: - // 别名撞上本侧已有会话时 Gateway 加后缀(`agent-only-chain` → - // `agent-only-chain-2`),而定稿别名又会**回写进 pi 的会话文件** —— - // 于是下一次心跳上报的 slug 变成加了后缀那个,人从补全里选的名字凭空消失。 - // 实测撞出来过一次。 - // - // 接管会话的别名由服务端在接管时按 slug 定好(AdoptPlatformSession), - // 这里不需要也不应该再动它。 - if (!adopted.has(piSessionId)) { - await syncNaming(piSessionId, entry.session.sessionName) - .catch((e) => log(`命名同步失败: ${describeError(e)}`)); - } - - // 只取 type==='text' 的块(B-5.1 / N-6):thinking 是思考过程,不是结论 - const text = lastAssistantText(entry.session.messages); - if (!text) return; // 空文本不发空邮件(B-5.4) - - const ctx = mailContexts.get(mailSessionID); - if (!ctx?.replyTo) return; // 不知道回给谁 - - // 幂等键用 pi 的会话 id + 会话树叶子 id:两者都落盘,重启重放也是同一个键。 - const relayKey = relayKeyFor(piSessionId, entry.sessionManager.getLeafId?.()); - if (relayedSummaries.get(piSessionId) === relayKey) return; - - // 模型这一轮已亲手回过这条线索 → 让位(B-5.3)。 - // 否则收件箱里是两封说同一件事的邮件(生产实测过)。 - if (shouldSkipAutoRelay(explicitSends.get(piSessionId), ctx.replyTo, ctx.mailID)) { - explicitSends.delete(piSessionId); - relayedSummaries.set(piSessionId, relayKey); - log(`本轮模型已主动回信 ${ctx.replyTo},跳过自动转发`); - return; - } - - await client.post('/mail/send', { - to: ctx.replyTo, - subject: replySubject(ctx.subject), - body: text, - reply_to: ctx.mailID || '', - // relay + relay_key 走免配额通道(I-2):模型已经把话说完了, - // 桥只是把它搬到邮件里。对搬运收费会让配额用尽时 Agent 连交代都做不了。 - relay: 'summary', - relay_key: relayKey, - }); - relayedSummaries.set(piSessionId, relayKey); - explicitSends.delete(piSessionId); // 一轮结束,窗口关闭 - log(`已转发本轮总结给 ${ctx.replyTo}(${text.length} 字)`); -} - -// ─── 投递(B-3 / B-6)─── - -async function deliverMail(data, kind, mailTools) { - const { session, reused } = await resolveSession(data, mailTools); - const piSessionId = session.sessionId; - - // 新一轮开始:清掉上一轮「模型主动发过信」的记录。不清的话, - // 上一轮亲手回过信会永久压掉这个会话之后所有的自动转发。 - explicitSends.delete(piSessionId); - - if (kind === 'mail' && data.session_id) { - // 一个会话里可能来过多封信,只留最近那封 —— 回信要落回最新的线索 - mailContexts.set(data.session_id, { - replyTo: data.from_name || '', - subject: data.subject || '', - mailID: data.mail_id || '', - }); - } - - const prompt = buildMailPrompt({ agentName: AGENT_NAME, data, kind, reused }); - - // 续谈:会话已经存在,模型也已经定了(pi 的模型在 createAgentSession 时绑定), - // 所以这一支不做模型降级。runTurn 内部按 isStreaming 分流: - // 空闲就直接起一轮,正在跑就排到当轮之后(不打断上一封邮件的工作)。 - if (reused) { - const outcome = await runTurn(session, prompt, TURN_TIMEOUT_MS); - log(`续谈 ${piSessionId}(mail ${data.mail_id}${outcome.queued ? ',已排队' : ''})`); - // 续谈失败不换模型重试(换模型要换会话,会丢掉整条上下文 —— - // 而上下文正是发件人指定这条会话的原因),但要让失败可见。 - if (!outcome.ok) throw new Error(`续谈失败: ${outcome.error}`); - return; - } - - // 按管理员划定的范围逐个尝试(D-3)。 - // 关键点:`prompt()` resolve **不代表模型跑成功了** —— 无凭证的 provider - // 会让它 reject(实测 `No API key found for amazon-bedrock.`), - // 而上游报错走 stopReason==='error'。判定交给 classifyTurnOutcome。 - const attempts = modelAttemptOrder(allowedModels, { - provider: REPLY_PROVIDER, - model: REPLY_MODEL, - }); - const failures = []; - - for (const route of attempts) { - const label = route ? `${route.provider}/${route.model}` : '(平台默认)'; - if (route) { - const model = modelRuntime.getModel(route.provider, route.model); - if (!model) { - // 目录里根本没有这个路由:同步就能判定,不必起一轮 - failures.push({ ...route, error: `平台目录里没有 ${label}` }); - log(`模型 ${label} 不存在,跳过`); - continue; - } - // 换模型要换会话:pi 的模型在 createAgentSession 时绑定。 - // 上一次尝试失败的会话没有任何 assistant 消息,丢掉不损失内容。 - const cwd = sessions.get(data.session_id)?.cwd; - const current = sessions.get(data.session_id)?.session; - current?.dispose?.(); - const retried = await openSession({ - cwd, - modelRuntime, - model, - customTools: mailTools, - extension: permissionExtension((id) => mailContexts.get(id)), - }); - rebind(data.session_id, current?.sessionId ?? piSessionId, retried, cwd); - const outcome = await runTurn(retried.session, prompt, TURN_TIMEOUT_MS); - if (outcome.ok) { - if (failures.length) log(`${label} 成功(前 ${failures.length} 个失败)`); - return; - } - failures.push({ ...route, error: outcome.error }); - log(`模型 ${label} 失败: ${outcome.error}`); - continue; - } - - const outcome = await runTurn(session, prompt, TURN_TIMEOUT_MS); - if (outcome.ok) { - if (failures.length) log(`${label} 成功(前 ${failures.length} 个失败)`); - return; - } - failures.push({ error: outcome.error }); - log(`模型 ${label} 失败: ${outcome.error}`); - } - - // 全部失败 → 必须回信(B-6):模型一次都没跑起来,会话里没有任何 - // assistant 消息,自动转发因此什么也不会发 —— 发件人只会看到再无音讯。 - if (kind === 'mail' && data.from_name) { - try { - await client.post('/mail/send', { - to: data.from_name, - subject: `处理失败: ${data.subject || '(无主题)'}`, - body: renderFailureReport(failures, data.subject), - reply_to: data.mail_id || '', - relay: 'summary', - relay_key: `model-failure:${data.mail_id || piSessionId}`, - }); - log(`已回报模型调用失败给 ${data.from_name}`); - } catch (e) { - log(`失败回报也发不出去: ${describeError(e)}`); - } - } - // 发完仍要 throw(B-6.4):静默会让这次失败只存在于邮件里,日志上看不出来 - throw new Error(`范围内 ${failures.length} 个模型全部失败:${failures.map(f => f.error).join(' | ')}`); -} - -/** 换模型重开会话后,把三张映射表指向新会话。 */ -function rebind(mailSessionID, oldPiId, opened, cwd) { - reverseMap.delete(oldPiId); - mailDriven.delete(oldPiId); - // 免批授权跟着旧的 pi 会话作废:它是人对**那次**上下文的判断, - // 换模型意味着重开一条会话、重跑一遍提示,不该继承上一条的授权。 - permissionGrants.revokeSession(oldPiId); - const piSessionId = opened.session.sessionId; - // cwd 由调用方传:AgentSession 上没有 cwd getter(只有 sessionId / - // sessionFile / sessionName),从 sessionManager.getCwd() 也行, - // 但这里本来就有那个值,多绕一层没有意义。 - const entry = { session: opened.session, sessionManager: opened.sessionManager, cwd }; - if (mailSessionID) { - sessions.set(mailSessionID, entry); - reverseMap.set(piSessionId, mailSessionID); - mailDriven.add(piSessionId); - } - opened.session.subscribe((event) => { - if (event?.type === 'agent_end' && !event.willRetry) { - relaySummary(piSessionId).catch((e) => log(`自动转发失败: ${describeError(e)}`)); - } - if (event?.type === 'session_info_changed') { - syncNaming(piSessionId, event.name).catch((e) => log(`命名同步失败: ${describeError(e)}`)); - } - }); -} - // ─── 权限决策回来(B-4)─── -async function handlePermissionDecision(data, mailTools) { +/** + * 把决策路由给发起询问的那个 worker。 + * + * 找不到 worker 有两种情形,都不该新开会话(B-4.3): + * - 桥重启了:那次工具调用早已随进程消失。但人刚刚点了「同意」—— + * 什么都不做的话人以为自己批准了、Agent 却毫无反应,所以退化为把决策 + * 当一封通知投进原会话(B-4.2)。 + * - 那条邮件会话从没被处理过:连通知都无处可投,只能记一行日志。 + */ +function handlePermissionDecision(data) { const relayKey = data.relay_key || ''; - const pending = relayKey ? pendingPermissions.get(relayKey) : undefined; - - if (pending) { - pendingPermissions.delete(relayKey); - pending.resolve(String(data.decision || '拒绝')); - log(`权限 ${relayKey} 决策 ${data.decision}(决策人 ${data.decided_by || '?'})`); + if (relayKey && pool.routePermission(relayKey, String(data.decision || '拒绝'))) { + log(`权限 ${relayKey} 决策 ${data.decision}(决策人 ${data.decided_by || '?'})已转交 worker`); return; } - // 找不到挂起项(桥重启丢了内存映射)→ 退化为把决策当一封通知投进原会话(B-4.2)。 - // 此时 pi 侧那次工具调用早已随进程消失,但人刚刚点了「同意」—— - // 什么都不做的话人以为自己批准了、Agent 却毫无反应。 - if (!data.session_id || !sessions.has(data.session_id)) { + if (!data.session_id || !pool.hasSession(data.session_id)) { // **不得凭空新开会话**(B-4.3) log(`权限决策 ${relayKey} 无对应会话,忽略`); return; } log(`权限 ${relayKey} 无挂起项,退化为通知投递`); - await deliverMail(data, 'permission', mailTools); + pool.submit('permission', data); } // ─── 心跳(B-2)─── @@ -769,7 +181,8 @@ async function reportSessions() { // 用 listAll 而不是 list(cwd):桥的进程 cwd 与会话 cwd 无关, // 按前者过滤会漏掉所有真正在干活的会话。 const all = await SessionManager.listAll(); - return snapshotPiSessions(all, (id) => mailDriven.has(id)); + const driven = pool.mailDrivenIDs(); + return snapshotPiSessions(all, (id) => driven.has(id)); } catch (e) { // 拉不到就**省略字段**而不是传 [](N-7 / W-3): // 空数组的语义是「平台确实一条会话都没有」,会把服务端镜像抹掉。 @@ -790,22 +203,26 @@ async function reportModels() { } } -async function catchUp(pending, mailTools) { +/** + * 补投离线期间积压的未读邮件(B-7)。 + * + * SSE 只推连上之后的事件,插件重启前发来的邮件不会再推一次。 + * + * 与旧版的差别:**不再 await 每一封**。旧版串行是因为「每封都要起一轮模型, + * 并发放出去等于对上游打 N 个并发请求」—— 那个约束现在由 pool 的 maxWorkers + * 承担,而且它比串行更好:同一条会话仍然串行,不同会话可以并行。 + */ +async function catchUp(pending) { if (!pending) return; try { const box = await client.get('/mail/inbox?status=unread&limit=20'); const tasks = selectCatchup(box?.mails ?? box, deliveredMails); if (!tasks.length) return; log(`补投 ${tasks.length} 封离线期间的邮件(共 ${pending} 封未读)`); - // 串行(B-7.2):每封都要起一轮模型,并发放出去等于对上游打 N 个并发请求 for (const ev of tasks) { if (deliveredMails.has(ev.mail_id)) continue; // 逐封再查(B-7.6) deliveredMails.add(ev.mail_id); - try { - await deliverMail(ev, 'mail', mailTools); - } catch (e) { - log(`补投 ${ev.mail_id} 失败: ${describeError(e)}`); - } + pool.submit('mail', ev); } } catch (e) { log(`补投失败: ${describeError(e)}`); @@ -828,8 +245,8 @@ async function main() { agentSecret: AGENT_SECRET, }); - // ModelRuntime 建一次全进程共用:它要读 auth.json / models.json 并做 - // 可用性探测,每条会话建一个既慢又会重复打 provider 的探测请求。 + // ModelRuntime 主进程也要一个:心跳的 reportModels 用它。worker 各自再建 + // 一个(跨进程传不了),代价是每个 worker 多 ~20ms(实测 11–23ms)。 // // allowModelNetwork 保持默认的 false:桥启动时不去网上拉模型目录。 // 拉了也没用 —— 上报给 Gateway 的是 getAvailable()(有凭证、真能调起来的), @@ -839,16 +256,43 @@ async function main() { const runtimeErr = modelRuntime.getError?.(); if (runtimeErr) log(`模型运行时告警: ${runtimeErr}`); - // onReconnect:connect_to_server 换了 Gateway 地址/密钥后,旧 SSE 长连仍连着 - // 旧地址(或已被旧密钥打断),必须在这里重建,模型调完工具才真正「切过去」。 - // reconfigure 已清空 lastEventID,所以 startSSE 会以「首次连接」姿态 - // (不带 Last-Event-ID,N-11)连上新地址 —— 拿旧序号问新 Gateway 只会 - // 命中一段无关的历史。 - const mailTools = createMailTools({ + // 工作进程池。config() 每次派活时取一次 —— allowedModels 随心跳变, + // 取快照会让 worker 用上一轮的模型范围。 + pool = createWorkerPool({ + log, + config: () => ({ + gatewayURL: client.baseURL, + agentName: AGENT_NAME, + agentKey: client.agentKey, + agentSecret: AGENT_SECRET, + allowedModels, + replyProvider: REPLY_PROVIDER, + replyModel: REPLY_MODEL, + turnTimeoutMs: TURN_TIMEOUT_MS, + }), + // worker 里 connect_to_server 换了坐标:worker 马上就退了,改在它自己身上 + // 等于没改。主进程据此重建 SSE,后续 worker 的 job 也会带上新坐标。 + onReconfigure: (url, key) => { + if (!client.reconfigure({ url, agentKey: key })) return; + log(`Gateway 坐标已更新为 ${client.baseURL},重建 SSE`); + client.stopSSE(); + client.startSSE(handleSSEEvent, log); + }, + maxWorkers: MAX_WORKERS, + workerMaxMs: WORKER_MAX_MS, + }); + + // 主进程仍要一套邮件工具:它自己不跑模型,但 connect_to_server 的 + // onReconnect 语义要在这里闭环(worker 侧那套只负责回报给主进程)。 + // + // 这些工具不会被任何模型调用 —— 主进程没有会话。留着是因为 + // createMailTools 同时承担「校验工具 schema」的职责(test/tool-schema.test.mjs), + // 而工具总数是契约里核对过的数字。 + createMailTools({ client, log, agentName: AGENT_NAME, onReconnect: () => { client.stopSSE(); - client.startSSE((type, data) => handleSSEEvent(type, data, mailTools), log); + client.startSSE(handleSSEEvent, log); }, }); @@ -874,7 +318,7 @@ async function main() { if (Array.isArray(res?.allowed_models)) allowedModels = res.allowed_models; // B-2.2 if (!caughtUp) { // B-7.1:只在首个成功心跳后补一次 caughtUp = true; - await catchUp(res?.pending_mails, mailTools); + await catchUp(res?.pending_mails); } } catch { // B-2.1:心跳失败不重试不报错。真连不上时 Gateway 会把它判成离线, @@ -884,7 +328,7 @@ async function main() { await beat(); // B-1.3:不等第一个 30 秒周期 heartbeatTimer = setInterval(beat, 30_000); // B-1.5 - client.startSSE((type, data) => handleSSEEvent(type, data, mailTools), log); + client.startSSE(handleSSEEvent, log); for (const sig of ['SIGINT', 'SIGTERM']) process.on(sig, () => shutdown(sig)); } @@ -892,13 +336,15 @@ async function main() { /** * SSE 事件分派。 * + * 这个函数**必须保持廉价**:它跑在读循环上。派活给 pool 是同步的(fork 是 + * 异步的,pool.submit 只是入队),所以读循环不会因为一封邮件停下。 + * * 提成命名函数是因为 connect_to_server 换地址后要用同一个处理器重建长连 —— * 内联箭头函数在那里拿不到,只能复制一遍,而复制出来的两份迟早会分叉。 */ -function handleSSEEvent(type, data, mailTools) { +function handleSSEEvent(type, data) { if (type === 'permission_decision') { - handlePermissionDecision(data, mailTools).catch((e) => - log(`权限决策处理失败: ${describeError(e)}`)); + handlePermissionDecision(data); return; } if (type !== 'new_mail') return; @@ -906,7 +352,7 @@ function handleSSEEvent(type, data, mailTools) { const id = data?.mail_id; if (!id || deliveredMails.has(id)) return; // B-3 第 1 步:去重 deliveredMails.add(id); - deliverMail(data, 'mail', mailTools).catch((e) => log(`投递 ${id} 失败: ${describeError(e)}`)); + pool.submit('mail', data); } function shutdown(reason) { @@ -917,21 +363,17 @@ function shutdown(reason) { if (heartbeatTimer) clearInterval(heartbeatTimer); // B-9.1 client?.stopSSE(); - // B-9.2 / N-9:所有未决权限询问 fail closed。 - // 不唤醒的话 pi 侧那些 await 永不返回,整条会话挂死; - // 而默认放行一个没人批准的危险操作,比让它失败严重得多。 - for (const [key, p] of pendingPermissions) { - log(`未决权限 ${key} fail closed`); - p.resolve('shutdown'); - } - pendingPermissions.clear(); + // pool.stop 先给每个 worker 发 shutdown(让它把未决权限询问 fail closed, + // B-9.2 / N-9),再给 2 秒自己退,然后 SIGKILL。 + // + // 不直接杀:pi 侧那些 await 不会返回,而 worker 里可能正握着会话文件。 + pool?.stop(); - for (const { session } of sessions.values()) { - try { session.dispose?.(); } catch { /* 关停期的报错没有价值 */ } - } releaseLock(); + // 留出 pool.stop 的宽限窗口再退:主进程先死会让子进程变成孤儿 + // (systemd 的 KillMode 会兜住,但那时 fail closed 已经来不及做了)。 + setTimeout(() => process.exit(0), 2500).unref?.(); // B-9.3:不发「插件下线」通知邮件 - process.exit(0); } main().catch((e) => { diff --git a/plugins/pi-mail-bridge/src/pool.mjs b/plugins/pi-mail-bridge/src/pool.mjs new file mode 100644 index 0000000..378bf1b --- /dev/null +++ b/plugins/pi-mail-bridge/src/pool.mjs @@ -0,0 +1,277 @@ +/** + * 投递工作进程池 —— 主进程侧的调度逻辑。 + * + * # 它解决的问题 + * + * 桥的主进程唯一的实时职责是读 SSE。模型工作放在主进程里跑会占满事件循环 + * (pi 的会话装载是同步的:23MB 的会话文件 `SessionManager.open` 一次阻塞 + * 118ms,实测;模型跑起来之后 SDK 内部还有大量同步工作),SSE 读循环停住, + * 后续邮件卡在 TCP 缓冲区,久到 Gateway 认为连接死了 → 重连 → 重放。 + * + * 所以:**收到事件就派给一个子进程,主进程立刻回去读 SSE。** + * + * # 并发与串行的边界 + * + * - **不同邮件会话并发**,上限 `maxWorkers`(默认 3)。上限的理由是内存 + * (每个 worker 约 140MB RSS,实测)和对上游 provider 的并发请求数。 + * - **同一邮件会话串行**。这是正确性要求,不是限流:pi 没有任何锁机制, + * 它假定「一个文件一个持有者」。两个 worker 同时装载同一条会话文件,各自的 + * 内存索引都看不见对方追加的行,算出的 parentId 指向对方不知道的 entry + * → 会话树分叉。串行还顺带保证了同一条线索里两封邮件的先后顺序。 + * + * # 排队而不是拒绝 + * + * 满载时邮件进 `queue`,有 worker 空出来就派。丢掉邮件是不可接受的: + * 发件人只会看到信发出去后再无音讯。队列无上限 —— 有上限就得决定丢哪封, + * 而任何丢弃策略都比「慢一点」糟。 + * + * # 主进程持有什么 + * + * 只有**路径与标量**:sessionFile / cwd / piSessionId / 「一直同意」表 / + * 命名同步指纹。AgentSession 对象跨不了进程边界,worker 每次从 sessionFile + * 重新装载 —— 拿到的是包含 TUI 期间写入的全部历史(这也让「短暂持有」 + * 从一套需要计时器兜底的机制退化成「worker 退出就是释放」)。 + * + * # IPC 协议 + * + * 主进程 → worker: + * `{type:'job', kind, data, session:{sessionFile,cwd}, grants, lastSyncedName, config}` + * `{type:'permission_decision', relayKey, decision}` + * `{type:'shutdown'}` + * worker → 主进程: + * `{type:'ready'}` 进程起来了,可以派活 + * `{type:'log', line}` 日志(主进程加 pid 前缀) + * `{type:'session_opened', piSessionId, sessionFile, cwd, reused}` + * `{type:'permission_pending', relayKey}` 主进程记下路由表 + * `{type:'permission_grant', toolName}` 「一直同意」要跨 worker 活下来 + * `{type:'name_synced', signature}` 命名指纹,防下一个 worker 重复 sync + * `{type:'reconfigure', url, agentKey}` connect_to_server 换了坐标 + * `{type:'done', ok, error}` 这封处理完了 + */ + +import { fork } from 'node:child_process'; +import { fileURLToPath } from 'node:url'; + +const WORKER_PATH = fileURLToPath(new URL('./worker.mjs', import.meta.url)); + +/** + * @param {object} deps + * @param {(...a: any[]) => void} deps.log + * @param {() => object} deps.config 每次派活时取一次(allowedModels 会随心跳变) + * @param {(url: string, key: string) => void} deps.onReconfigure + * @param {number} [deps.maxWorkers] + * @param {number} [deps.workerMaxMs] worker 硬超时:卡死的进程必须能被回收 + * @param {string} [deps.workerPath] 只为测试存在:换成不装 pi SDK 的桩 worker, + * 让调度不变量(并发上限、同会话串行、硬超时)能在毫秒级验证。 + */ +export function createWorkerPool({ + log, config, onReconfigure, + maxWorkers = 3, workerMaxMs = 600_000, workerPath = WORKER_PATH, +}) { + /** 正在跑的 worker:mailSessionKey -> {child, mailID, startedAt, timer} */ + const running = new Map(); + /** 等着派的活,先进先出。 */ + const queue = []; + /** relay_key -> mailSessionKey,把决策路由回发起询问的那个 worker。 */ + const permissionRoutes = new Map(); + /** + * 跨 worker 存活的会话状态:mailSessionKey -> {sessionFile, cwd, piSessionId, + * grants:Set, lastSyncedName}。 + * + * 这是 worker 一封一进程之后仍需在主进程留存的全部东西 —— 下一封邮件靠 + * sessionFile 接着谈,靠 grants 不重复问已经「一直同意」过的工具。 + */ + const sessionState = new Map(); + /** + * 被模型降级换掉的旧 pi 会话 id。 + * + * 仍要计入 mail_driven:它们已经参与过邮件往来,而磁盘上的会话文件 + * 不会因为换模型而消失 —— 心跳快照仍会上报它们。 + */ + const retired = new Set(); + let stopped = false; + + /** + * 邮件会话 id 作为串行化的键。 + * + * 没有 session_id 的事件(理论上不该有)退回 mail_id:那样每封各占一个 + * worker,不会串行 —— 但它们本来也不属于同一条会话。 + */ + const keyOf = (data) => data?.session_id || `mail:${data?.mail_id || Math.random()}`; + + function submit(kind, data) { + if (stopped) return; + queue.push({ kind, data, key: keyOf(data) }); + pump(); + } + + function pump() { + if (stopped) return; + for (let i = 0; i < queue.length; i++) { + const job = queue[i]; + // 同一会话已有 worker 在跑 → 跳过它,看后面有没有别的会话可以先跑。 + // 不能 break:那会让一条慢会话把所有别的会话都堵住(正是要修的病)。 + if (running.has(job.key)) continue; + if (running.size >= maxWorkers) return; + queue.splice(i, 1); + i--; + spawn(job); + } + } + + function spawn(job) { + const state = sessionState.get(job.key) || { grants: new Set(), lastSyncedName: '' }; + const child = fork(workerPath, [], { + // stdio 继承:worker 里 pi SDK 自己打的东西直接进 journalctl。 + // 'ipc' 必须显式列出,否则 process.send 不存在。 + stdio: ['ignore', 'inherit', 'inherit', 'ipc'], + }); + + // 硬超时:worker 卡死(模型不返回、权限等不到决策而主进程也没收到事件) + // 时必须能回收,否则那条会话的后续邮件永远排队。 + const timer = setTimeout(() => { + log(`worker ${child.pid} 处理 ${job.data?.mail_id} 超过 ${workerMaxMs / 1000}s,强杀`); + try { child.kill('SIGKILL'); } catch { /* 已经死了 */ } + }, workerMaxMs); + if (typeof timer.unref === 'function') timer.unref(); + + const entry = { child, mailID: job.data?.mail_id || '', key: job.key, startedAt: Date.now(), timer }; + running.set(job.key, entry); + + child.on('message', (msg) => onWorkerMessage(entry, msg)); + + child.on('exit', (code, signal) => { + clearTimeout(timer); + running.delete(job.key); + for (const [rk, k] of permissionRoutes) if (k === job.key) permissionRoutes.delete(rk); + if (code !== 0) { + log(`worker ${child.pid}(mail ${entry.mailID})异常退出 code=${code} signal=${signal || '-'}`); + } + pump(); + }); + + child.on('error', (e) => log(`worker ${child.pid} 出错: ${e?.message || e}`)); + + // 等 worker 说 ready 再派活:fork 返回时子进程的 import 还没跑完, + // 此时 send 的消息会排在 IPC 队列里(能收到,但 ready 让顺序确定)。 + child.once('message', function first(msg) { + if (msg?.type !== 'ready') return; + child.send({ + type: 'job', + kind: job.kind, + data: job.data, + session: { + sessionFile: state.sessionFile || '', + cwd: state.cwd || '', + }, + grants: [...state.grants], + lastSyncedName: state.lastSyncedName || '', + config: config(), + }); + }); + } + + function onWorkerMessage(entry, msg) { + const state = sessionState.get(entry.key) || { grants: new Set(), lastSyncedName: '' }; + switch (msg?.type) { + case 'log': + log(`[w${entry.child.pid}] ${msg.line}`); + return; + case 'session_opened': + // 一条会话可能先后用过多个 pi 会话 id(模型降级会换会话)。 + // 旧 id 仍计入 mail_driven,理由见 retired 的注释。 + if (state.piSessionId && state.piSessionId !== msg.piSessionId) { + retired.add(state.piSessionId); + } + state.piSessionId = msg.piSessionId; + state.sessionFile = msg.sessionFile; + state.cwd = msg.cwd; + sessionState.set(entry.key, state); + return; + case 'permission_pending': + permissionRoutes.set(msg.relayKey, entry.key); + return; + case 'permission_grant': + // 「一直同意」必须跨 worker 活着:worker 一封一进程,不存的话下一封 + // 邮件又问一遍,那个选项就是在骗人。 + state.grants.add(msg.toolName); + sessionState.set(entry.key, state); + return; + case 'name_synced': + state.lastSyncedName = msg.signature; + sessionState.set(entry.key, state); + return; + case 'reconfigure': + onReconfigure?.(msg.url, msg.agentKey); + return; + case 'done': + if (!msg.ok) log(`投递 ${entry.mailID} 失败: ${msg.error}`); + return; + default: + return; + } + } + + /** + * 把权限决策路由到发起询问的那个 worker。 + * + * @returns {boolean} 有没有找到对应的 worker。找不到说明那个 worker 已经退了 + * (桥重启、硬超时被杀、或者处理已经结束)—— 调用方据此走 B-4.2 的 + * 降级路径(把决策当一封通知投进原会话)。 + */ + function routePermission(relayKey, decision) { + const key = permissionRoutes.get(relayKey); + if (!key) return false; + const entry = running.get(key); + if (!entry) { + permissionRoutes.delete(relayKey); + return false; + } + permissionRoutes.delete(relayKey); + entry.child.send({ type: 'permission_decision', relayKey, decision }); + return true; + } + + /** 这条邮件会话有 worker 在跑吗(B-4.2 判断降级路径用)。 */ + const hasSession = (mailSessionID) => sessionState.has(mailSessionID); + + /** + * 邮件驱动过的 pi 会话 id,喂给心跳快照的 `mail_driven` 标记。 + * + * 不随 worker 退出而清:worker 退了不代表那条会话不再参与邮件往来 —— + * 下一封邮件还会接着谈,而人在补全里需要看到它带着这个标记。 + * 重启丢是已知取舍(契约第六节)。 + */ + const mailDrivenIDs = () => { + const out = new Set(retired); + for (const st of sessionState.values()) { + if (st.piSessionId) out.add(st.piSessionId); + } + return out; + }; + + function stop() { + stopped = true; + queue.length = 0; + for (const { child, timer } of running.values()) { + clearTimeout(timer); + // 先 shutdown 让 worker 把未决权限 fail closed(B-9.2),再给它一点 + // 时间自己退。不直接 SIGKILL:那样 pi 侧的 await 不会返回,而 worker + // 里可能正握着会话文件。 + try { child.send({ type: 'shutdown' }); } catch { /* 通道已断 */ } + setTimeout(() => { try { child.kill('SIGKILL'); } catch { /* 已经死了 */ } }, 2000).unref?.(); + } + } + + /** 观测用:现在跑着几个、排了几个。 */ + const stats = () => ({ + running: running.size, + queued: queue.length, + sessions: sessionState.size, + workers: [...running.values()].map((e) => ({ + pid: e.child.pid, mailID: e.mailID, ageMs: Date.now() - e.startedAt, + })), + }); + + return { submit, routePermission, hasSession, mailDrivenIDs, stop, stats }; +} diff --git a/plugins/pi-mail-bridge/src/worker.mjs b/plugins/pi-mail-bridge/src/worker.mjs new file mode 100644 index 0000000..f067484 --- /dev/null +++ b/plugins/pi-mail-bridge/src/worker.mjs @@ -0,0 +1,555 @@ +#!/usr/bin/env node +/** + * pi-mail-bridge 的**投递工作进程** —— 一封邮件一个,跑完就退。 + * + * # 为什么必须是独立进程 + * + * 桥的主进程要一直读 SSE。而 pi 的会话装载是**同步**的: + * `SessionManager.open()` 走 `openSync` + `readSync` 循环把整个 `.jsonl` + * 读进内存并逐行 JSON.parse(SDK core/session-manager.js 的 + * loadEntriesFromFile)。实测本机最大那条会话 23MB,`open` 一次 + * **阻塞事件循环 118ms**;模型跑起来之后 SDK 内部还有大量同步工作。 + * 全都发生在主线程上,SSE 读循环在那期间完全停住 —— 后续邮件卡在 TCP + * 缓冲区里,久到 Gateway 认为连接死了,重连又触发重放。 + * + * 把模型工作搬进子进程后,主进程只剩「收事件 → 去重 → 分派」, + * 实测同样的活在 fork 出的子进程里跑,主进程事件循环阻塞 0ms。 + * + * # 为什么是 child_process 而不是 worker_threads + * + * 两者实测都能建起 AgentSession。选进程的理由是**隔离**: + * 模型会跑 bash/write/edit,一次 OOM 或 uncaughtException 不该带走 + * 整座桥;而 `worker_threads` 与主线程共享堆和进程生命周期。 + * 代价是每个 worker 约 140MB RSS 和 ~500ms 启动,由主进程的并发上限约束。 + * + * # 一封邮件一个进程带来的简化 + * + * 原来的「接管会话短暂持有」(adopted / adoptTimers / releaseAdopted) + * 整套机制没有了:worker 退出**就是**释放,而且对普通会话和接管会话 + * 一视同仁 —— 每封邮件都是「open → 跑一轮 → 还回去」。 + * pi 假定「一个文件一个持有者」,主进程按会话串行分派保证了这一点。 + * + * # 进程边界上传什么 + * + * 主进程持有的是**路径与标量**(sessionFile / cwd / piSessionId / 授权表), + * 不是 AgentSession 对象 —— 那东西跨不了进程。worker 每次从 sessionFile + * 重新装载,拿到的是包含 TUI 期间写入的全部历史。 + * + * 协议见 src/pool.mjs 顶部。 + */ + +import { existsSync } from 'node:fs'; +import { homedir } from 'node:os'; +import { join } from 'node:path'; +import { ModelRuntime } from '@earendil-works/pi-coding-agent'; + +import { GatewayClient } from './gateway.mjs'; +import { createMailTools } from './tools.mjs'; +import { openSession, runTurn } from './session-pool.mjs'; +import { buildMailPrompt, lastAssistantText, replySubject, relayKeyFor, describeError } from './turn.mjs'; +import { planNamingSync, planWriteBack } from './naming.mjs'; +import { resolveWorkspaceCwd, ensureCwd } from '../lib/workspace.js'; +import { modelAttemptOrder, renderFailureReport } from '../lib/model-scope.js'; +import { explicitSends, shouldSkipAutoRelay } from '../lib/relay-dedup.js'; +import { adoptedSessionID, adoptMissingMessage } from '../lib/adopt.js'; +import { isApproval, isAlwaysDecision } from '../lib/permission-grants.js'; + +// ─── 与主进程的通道 ─── + +/** 日志一律回传主进程:worker 的 stdout 会混在一起,加 pid 前缀才分得清。 */ +const log = (...args) => send({ type: 'log', line: args.join(' ') }); + +function send(msg) { + try { + process.send?.(msg); + } catch { + /* 主进程已经走了,这条日志没有去处 */ + } +} + +/** relay_key -> resolve;决策由主进程的 SSE 收到后路由进来。 */ +const pending = new Map(); + +/** + * 这一轮里人点过「一直同意」的工具。 + * + * 初值由主进程在 job 里给(跨 worker 持久),新增的回报给主进程。 + * 不用 lib/permission-grants.js 的 store:那份按 (会话, 工具) 存, + * 而 worker 只服务一条会话,一个 Set 就够,且要能整体回传。 + */ +const grants = new Set(); + +// ─── 单封邮件的全部状态 ─── +// +// worker 只处理一封邮件、只碰一条会话,所以这些原来在主进程里 +// 按 sessionId 分桶的 Map 在这里都退化成单个变量。 + +let job = null; +let client = null; +let modelRuntime = null; +let piSessionId = ''; +let mailContext = { replyTo: '', subject: '', mailID: '' }; +let lastSyncedName = ''; +let relayedKey = ''; +let finished = false; + +// ─── 权限钩子(B-8)─── + +/** + * 与主进程版本逐条对应,差别只有两处: + * - 「是不是邮件驱动的会话」不必查表:worker 只为邮件而存在。 + * - 等决策的 promise 由主进程通过 IPC 唤醒,而不是本进程的 SSE。 + * + * 决策等待期间**只有这个 worker 停住**,主进程照常读 SSE、照常给别的会话 + * 派活 —— 这正是原来最难受的一处:权限询问会让整座桥不再收信。 + */ +function permissionExtension() { + const GUARDED = new Set(['bash', 'write', 'edit']); + + return (pi) => { + pi.on('tool_call', async (event, ctx) => { + if (!GUARDED.has(event.toolName)) return; + + const sid = ctx?.sessionManager?.getSessionId?.() || ''; + // 只管自己那条会话。worker 里不该出现第二条,出现了说明有 bug —— + // 让位(返回 undefined)比拦错一个安全。 + if (piSessionId && sid && sid !== piSessionId) return; + + // 人点过「一直同意」→ 直接放行。必须在 POST 之前:否则每条命令都发一封 + // 邮件,那个选项形同虚设(生产实测同一条会话被问了 15 次 bash)。 + if (grants.has(event.toolName)) return; + + // relay_key 用 pi 给的 toolCallId(B-8.1):服务端随决策事件回传它。 + const relayKey = `${sid || piSessionId}:${event.toolCallId}`; + + try { + // 不传 `to`:决策人由服务端按 会话 owner → 线索里最近的人类 → 409 + // 解析。插件只有本地上下文,猜不出「这条 Agent 链最初是谁派的活」。 + await client.post('/permission/request', { + question: `是否允许执行 ${event.toolName}?`, + options: ['同意', '一直同意', '拒绝'], + context: [ + describeToolCall(event), + mailContext.subject ? `\n触发任务:${mailContext.subject}` : '', + mailContext.replyTo ? `任务来自:${mailContext.replyTo}` : '', + ].filter(Boolean).join('\n'), + session_id: job.data?.session_id || '', + relay_key: relayKey, + }); + } catch (e) { + // 409 = 服务端判定这条任务链上没有人类,永远不会有人来点头。 + // 当场 block 并把服务端建议原文当 reason:模型从工具报错里看到 + // 「没人可问,换不需要权限的方式」才能自己改道,挂死时连重试机会都没有。 + if (e?.status === 409) { + const b = e.body || {}; + const reason = [ + b.error || '权限询问无法送达:这条任务链上没有人类用户', + b.detail || '', + b.suggestion || '', + ].filter(Boolean).join('\n'); + log(`权限询问无人可投,当场拒绝 ${relayKey}:${b.error || ''}`); + return { block: true, reason }; + } + // 其余失败是暂时的 → 让位给 pi 本地决策(B-8.2)。 + log(`权限转发失败,让位给本地决策: ${describeError(e)}`); + return; + } + + log(`权限询问已发出(${event.toolName},key=${relayKey}),等待决策…`); + send({ type: 'permission_pending', relayKey }); + const decision = await new Promise((resolve) => pending.set(relayKey, resolve)); + + // fail closed(B-9.2 / N-9):只有明确同意才放行。 + if (isApproval(decision)) { + // 「一直同意」要真的记住,否则这个选项在骗人。判定交给 isAlwaysDecision —— + // 「同意」是单次授权,把它当 always 会放行人没看过的后续命令。 + if (isAlwaysDecision(decision)) { + grants.add(event.toolName); + send({ type: 'permission_grant', toolName: event.toolName }); + log(`本会话的 ${event.toolName} 已获「一直同意」,后续不再询问`); + } + log(`权限 ${relayKey} 获批(${decision}),放行 ${event.toolName}`); + return; + } + return { + block: true, + reason: `用户${decision === 'shutdown' ? '未及决策(桥已关停)' : `拒绝了这次 ${event.toolName} 调用`}`, + }; + }); + }; +} + +/** 把一次工具调用摘要成人能判断的文本(B-8.4)。 */ +function describeToolCall(event) { + const input = event?.input ?? {}; + if (event.toolName === 'bash') return `命令:\n${String(input.command ?? '').slice(0, 800)}`; + if (event.toolName === 'write' || event.toolName === 'edit') { + return `文件:${input.file_path ?? input.path ?? '(未给出)'}`; + } + return JSON.stringify(input).slice(0, 800); +} + +// ─── 会话装载 ─── + +/** + * 没有可用 `to_workspace` 时的兜底目录。 + * + * 与 DSH 的 `mailSessionFallback` 同构,但目录名是 `.pi`:那个函数在 lib/ 下 + * (三平台逐字节相同),写死了 `.dsh`,不能为 pi 改 —— pi 的会话落进 `~/.dsh/` + * 会让人以为是 DSH 在干活。 + */ +function piMailFallback(sessionKey) { + return join(homedir(), '.pi', 'mail-sessions', String(sessionKey || 'default')); +} + +/** + * 找到这封邮件该落进的会话文件,装载它。 + * + * 三条路,优先级从高到低: + * 1. 主进程给了 sessionFile —— 这条邮件会话之前有 worker 跑过,接着谈。 + * 2. 服务端说接管了平台会话(人在 TUI 里开的那条)→ 按 platform id 找文件。 + * 3. 都没有 → 新建,cwd 取地址的 path 位(B-3.1)。 + * + * 返回 `reused` 供提示词与模型降级判断:有历史的会话不重新自我介绍, + * 也不做模型降级(换模型要换会话,会丢掉整条上下文,而上下文正是 + * 发件人指定这条会话的原因)。 + */ +async function loadSession(mailTools) { + const data = job.data; + const given = job.session?.sessionFile; + + if (given && existsSync(given)) { + const cwd = job.session.cwd || resolveWorkspaceCwd( + data.to_workspace, piMailFallback(data.session_id)).cwd; + const opened = await openSession({ + cwd, modelRuntime, customTools: mailTools, + extension: permissionExtension(), sessionFile: given, + }); + return { ...opened, cwd, reused: true }; + } + + const adoptID = adoptedSessionID(data); + if (adoptID) { + const { SessionManager } = await import('@earendil-works/pi-coding-agent'); + // listAll 而不是 list(cwd):worker 的进程 cwd 与会话 cwd 无关。 + const all = await SessionManager.listAll(); + const info = all.find((e) => e?.id === adoptID); + if (!info?.path) { + // 镜像是快照,可以过期。**不能**退回「新建一条」—— 那会让人在 TUI 里 + // 看不到这封邮件带来的对话,而那正是接管的目的(N-8:静默改语义比报错糟)。 + throw new Error(adoptMissingMessage(adoptID, '磁盘上已无这个会话文件')); + } + // cwd 取会话自己的(SessionInfo.cwd 来自持久化 header); + // 老会话的 cwd 是空串,那种情况退回地址里的 path 位。 + const { cwd } = resolveWorkspaceCwd( + info.cwd || data.to_workspace, piMailFallback(data.session_id)); + const opened = await openSession({ + cwd, modelRuntime, customTools: mailTools, + extension: permissionExtension(), sessionFile: info.path, + }); + log(`接管 pi 会话 ${adoptID}(cwd=${cwd},文件 ${info.path})`); + return { ...opened, cwd, reused: true }; + } + + // 目录不存在时**不创建**(N-2:笔误会在磁盘上落下真目录,而 Agent 在 + // 里面一无所获),拒绝相对路径(N-3)。 + const { cwd, grouped } = resolveWorkspaceCwd( + data.to_workspace, piMailFallback(data.session_id)); + if (!grouped && data.to_workspace) { + log(`工作目录 ${data.to_workspace} 不可用,回退到 ${cwd}`); + } + ensureCwd(cwd, grouped); + const opened = await openSession({ + cwd, modelRuntime, customTools: mailTools, extension: permissionExtension(), + }); + log(`新建 pi 会话 ${opened.session.sessionId}(cwd=${cwd})`); + return { ...opened, cwd, reused: false }; +} + +// ─── 命名一致(C-11 / W-7)─── + +async function syncNaming(session, platformName) { + const mailSessionID = job.data?.session_id; + if (!mailSessionID) return; + + const plan = planNamingSync({ + platformName, + mailSubject: mailContext.subject, + lastSynced: lastSyncedName, + }); + if (plan.skip) return; + + // 先记指纹再发请求:响应回来时 setSessionName 会再次触发 + // session_info_changed,这一步是防自激循环的关键。 + lastSyncedName = plan.signature; + send({ type: 'name_synced', signature: plan.signature }); + + const res = await client.post(`/sessions/${mailSessionID}/sync`, { + alias: plan.alias, title: plan.title, + }); + + const back = planWriteBack({ finalAlias: res?.alias, currentPiName: session.sessionName }); + log(`命名同步 ${piSessionId}: alias=${res?.alias || '(未变)'} 来源=${plan.source}`); + + if (back.write) { + // 顺序要紧:先更新指纹,再改名。setSessionName **同步**触发 + // session_info_changed(实测),本函数会被重入;指纹后更新的话重入那次 + // 看到旧指纹,于是又打一次 sync —— 每条会话两次内容相同的请求。 + lastSyncedName = `platform:${back.name}|${back.name}`; + send({ type: 'name_synced', signature: lastSyncedName }); + // 只用 setSessionName(走 pi 自己的写入路径)。绝不自己拼路径写会话文件: + // 首条 assistant 消息落盘前文件还不存在,pi 首次落盘用 openSync(file,"wx"), + // 抢先创建会让它抛 EEXIST(实测)。 + session.setSessionName(back.name); + log(`别名回写 pi:${back.name}(${back.reason})`); + } +} + +// ─── 自动转发(B-5)─── + +/** + * 把这一轮的结论转回发件人。 + * + * 与主进程版本的差别:不再挂在 `agent_end` 订阅上,而是在 `runTurn` 返回后 + * **确定性地**调一次。原来必须靠事件是因为主进程的 60 秒超时会先返回、 + * 会话还在跑;worker 没有那个约束(它就为这封邮件活着),等真结束再转发。 + * + * @param {boolean} adopted 接管会话跳过命名同步 —— 它的别名是人从补全里 + * 选中的平台 slug,同步会双向改坏:撞名时 Gateway 加后缀,定稿别名又回写进 + * pi 会话文件,于是下次心跳上报的 slug 变成带后缀那个。实测撞出来过一次。 + */ +async function relaySummary(session, sessionManager, adopted) { + if (!adopted) { + await syncNaming(session, session.sessionName) + .catch((e) => log(`命名同步失败: ${describeError(e)}`)); + } + + // 只取 type==='text' 的块(B-5.1 / N-6):thinking 是思考过程,不是结论。 + const text = lastAssistantText(session.messages); + if (!text) return; // 空文本不发空邮件(B-5.4) + if (!mailContext.replyTo) return; // 不知道回给谁 + + // 幂等键用 pi 会话 id + 会话树叶子 id:两者都落盘,重放也是同一个键。 + const relayKey = relayKeyFor(piSessionId, sessionManager.getLeafId?.()); + if (relayedKey === relayKey) return; + + // 模型这一轮已亲手回过这条线索 → 让位(B-5.3)。 + // 否则收件箱里是两封说同一件事的邮件(生产实测过)。 + if (shouldSkipAutoRelay(explicitSends.get(piSessionId), mailContext.replyTo, mailContext.mailID)) { + relayedKey = relayKey; + log(`本轮模型已主动回信 ${mailContext.replyTo},跳过自动转发`); + return; + } + + await client.post('/mail/send', { + to: mailContext.replyTo, + subject: replySubject(mailContext.subject), + body: text, + reply_to: mailContext.mailID || '', + // relay + relay_key 走免配额通道(I-2):模型已经把话说完了,桥只是把它 + // 搬到邮件里。对搬运收费会让配额用尽时 Agent 连交代都做不了。 + relay: 'summary', + relay_key: relayKey, + }); + relayedKey = relayKey; + log(`已转发本轮总结给 ${mailContext.replyTo}(${text.length} 字)`); +} + +// ─── 主流程 ─── + +async function run() { + const data = job.data; + const kind = job.kind; + + client = new GatewayClient({ + url: job.config.gatewayURL, + agentName: job.config.agentName, + agentKey: job.config.agentKey, + agentSecret: job.config.agentSecret, + }); + + // ModelRuntime 每个 worker 建一次。allowModelNetwork 保持默认 false: + // 上报给 Gateway 的是 getAvailable()(有凭证、真能调起来的), + // 那取决于本机 auth.json,不取决于目录里有多少条。 + modelRuntime = await ModelRuntime.create(); + const runtimeErr = modelRuntime.getError?.(); + if (runtimeErr) log(`模型运行时告警: ${runtimeErr}`); + + // connect_to_server 在 worker 里换了坐标要让主进程知道:worker 马上就退了, + // 改在自己身上等于没改。主进程收到后重建 SSE 并写进后续 worker 的 job。 + const mailTools = createMailTools({ + client, log, agentName: job.config.agentName, + onReconnect: () => send({ + type: 'reconfigure', url: client.baseURL, agentKey: client.agentKey, + }), + }); + + const { session, sessionManager, diagnostics, cwd, reused } = await loadSession(mailTools); + for (const d of diagnostics) log(`扩展诊断: ${d?.message ?? JSON.stringify(d)}`); + + piSessionId = session.sessionId; + // 「这条会话是接管来的吗」以 Gateway 给的 platform_session_id 为准, + // **不能**从 `reused && !job.session.sessionFile` 推断:第二封邮件进同一条 + // 接管会话时主进程给了 sessionFile,那样推断会得出 false,于是命名同步跑起来 + // 把人从补全里选的 slug 冲掉 —— 那正是上一轮修掉的那个 bug。 + const adopted = Boolean(adoptedSessionID(data)); + send({ + type: 'session_opened', + piSessionId, + sessionFile: session.sessionFile || '', + cwd, + reused, + }); + + // pi 侧改名(人在 TUI 里 /name)→ 同步给 Gateway。接管会话不同步,理由见 + // relaySummary 的 adopted 参数注释。 + session.subscribe((event) => { + if (event?.type === 'session_info_changed' && !adopted) { + syncNaming(session, event.name).catch((e) => log(`命名同步失败: ${describeError(e)}`)); + } + }); + + const prompt = buildMailPrompt({ agentName: job.config.agentName, data, kind, reused }); + + // 轮次超时给得很宽(默认 10 分钟,由主进程的 workerMaxMs 派生)。 + // + // 原来是 60 秒「超时按成功返回」,因为主进程要腾出手来收下一封邮件; + // worker 没有这个理由 —— 它只为这封邮件活着,等真结论更准。带工具调用的 + // 一轮跑几分钟很正常,60 秒返回会让转发落在一个还没说完的结论上。 + const turnTimeout = job.config.turnTimeoutMs; + + if (reused) { + // 续谈不做模型降级:换模型要换会话,会丢掉整条上下文 —— 而上下文正是 + // 发件人指定这条会话的原因。但失败要可见。 + const outcome = await runTurn(session, prompt, turnTimeout); + log(`续谈 ${piSessionId}(mail ${data.mail_id}${outcome.queued ? ',已排队' : ''})`); + if (!outcome.ok) throw new Error(`续谈失败: ${outcome.error}`); + await relaySummary(session, sessionManager, adopted); + return; + } + + // 按管理员划定的范围逐个尝试(D-3)。 + // 关键点:`prompt()` resolve **不代表模型跑成功了** —— 无凭证的 provider + // 会让它 reject(实测 `No API key found for amazon-bedrock.`), + // 而上游报错走 stopReason==='error'。判定交给 classifyTurnOutcome。 + const attempts = modelAttemptOrder(job.config.allowedModels, { + provider: job.config.replyProvider, + model: job.config.replyModel, + }); + const failures = []; + let live = { session, sessionManager }; + + for (const route of attempts) { + const label = route ? `${route.provider}/${route.model}` : '(平台默认)'; + if (route) { + const model = modelRuntime.getModel(route.provider, route.model); + if (!model) { + // 目录里根本没有这个路由:同步就能判定,不必起一轮。 + failures.push({ ...route, error: `平台目录里没有 ${label}` }); + log(`模型 ${label} 不存在,跳过`); + continue; + } + // 换模型要换会话:pi 的模型在 createAgentSession 时绑定。 + // 上一次尝试失败的会话没有任何 assistant 消息,丢掉不损失内容。 + try { live.session.dispose?.(); } catch { /* 已经没了 */ } + const retried = await openSession({ + cwd, modelRuntime, model, customTools: mailTools, + extension: permissionExtension(), + }); + live = { session: retried.session, sessionManager: retried.sessionManager }; + piSessionId = retried.session.sessionId; + send({ + type: 'session_opened', + piSessionId, + sessionFile: retried.session.sessionFile || '', + cwd, + reused: false, + }); + } + const outcome = await runTurn(live.session, prompt, turnTimeout); + if (outcome.ok) { + if (failures.length) log(`${label} 成功(前 ${failures.length} 个失败)`); + await relaySummary(live.session, live.sessionManager, false); + return; + } + failures.push({ ...(route || {}), error: outcome.error }); + log(`模型 ${label} 失败: ${outcome.error}`); + } + + // 全部失败 → 必须回信(B-6):模型一次都没跑起来,会话里没有任何 assistant + // 消息,自动转发因此什么也不会发 —— 发件人只会看到再无音讯。 + if (kind === 'mail' && data.from_name) { + try { + await client.post('/mail/send', { + to: data.from_name, + subject: `处理失败: ${data.subject || '(无主题)'}`, + body: renderFailureReport(failures, data.subject), + reply_to: data.mail_id || '', + relay: 'summary', + relay_key: `model-failure:${data.mail_id || piSessionId}`, + }); + log(`已回报模型调用失败给 ${data.from_name}`); + } catch (e) { + log(`失败回报也发不出去: ${describeError(e)}`); + } + } + // 发完仍要 throw(B-6.4):静默会让这次失败只存在于邮件里,日志上看不出来。 + throw new Error(`范围内 ${failures.length} 个模型全部失败:${failures.map((f) => f.error).join(' | ')}`); +} + +// ─── 入口 ─── + +process.on('message', (msg) => { + if (msg?.type === 'job') { + if (job) return; // 一个 worker 只接一封 + job = msg; + mailContext = { + replyTo: msg.data?.from_name || '', + subject: msg.data?.subject || '', + mailID: msg.data?.mail_id || '', + }; + lastSyncedName = msg.lastSyncedName || ''; + for (const t of msg.grants || []) grants.add(t); + run() + .then(() => finish({ ok: true, error: '' })) + .catch((e) => finish({ ok: false, error: describeError(e) })); + return; + } + if (msg?.type === 'permission_decision') { + const resolve = pending.get(msg.relayKey); + if (!resolve) return; + pending.delete(msg.relayKey); + resolve(String(msg.decision || '拒绝')); + return; + } + if (msg?.type === 'shutdown') { + // fail closed(B-9.2 / N-9):唤醒所有未决询问,让 pi 侧那些 await 返回。 + // 不唤醒的话会话挂死;而默认放行一个没人批准的危险操作更糟。 + for (const [key, resolve] of pending) { + log(`未决权限 ${key} fail closed`); + resolve('shutdown'); + } + pending.clear(); + } +}); + +function finish(result) { + if (finished) return; + finished = true; + send({ type: 'done', ...result }); + // **一律 exit 0**:`done` 已经把成败说清楚了,用退出码再说一遍会让主进程 + // 把「模型全部失败」这种已处理的结果也打成「worker 异常退出」。 + // 非零退出码留给真正的崩溃(没来得及发 done 的那种)。 + // + // 给 IPC 一个 tick 把 done 送出去再退:未 flush 的消息会丢,而主进程靠它判成败。 + // unref 不影响触发(IPC 通道本身在保持事件循环存活),只是不由这个计时器兜着。 + const t = setTimeout(() => process.exit(0), 50); + if (typeof t.unref === 'function') t.unref(); +} + +// 未捕获异常也要回报:静默退出会让主进程只看到 exit code, +// 而那条邮件的失败原因就此丢失。 +process.on('uncaughtException', (e) => finish({ ok: false, error: `未捕获异常: ${describeError(e)}` })); +process.on('unhandledRejection', (e) => finish({ ok: false, error: `未处理拒绝: ${describeError(e)}` })); + +send({ type: 'ready' }); diff --git a/plugins/pi-mail-bridge/test/pool.test.mjs b/plugins/pi-mail-bridge/test/pool.test.mjs new file mode 100644 index 0000000..4441efa --- /dev/null +++ b/plugins/pi-mail-bridge/test/pool.test.mjs @@ -0,0 +1,389 @@ +/** + * 工作进程池的调度不变量。 + * + * 这些用例**真的 fork 子进程**,用一个极小的桩 worker(不装 pi SDK): + * 要验的是调度(并发上限、同会话串行、排队、硬超时、状态跨 worker 存活), + * 而不是模型怎么跑。用桩让每个用例在几百毫秒内完成。 + * + * 桩 worker 的协议与真 worker 一致:`ready` → 收 `job` → 按 `__hold` 停一会儿 + * → 发 `done`。它还把收到的 job 载荷原样回声成一行 `JOB {…}` 日志 —— + * 判据因此能落在「主进程真的把 sessionFile / grants / 命名指纹传下来了」上, + * 而不是一个间接的计数(计数在字段被丢掉时依然会给出绿色)。 + */ + +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { mkdtempSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; + +// ─── 桩 worker ─── +// +// 落到临时目录而不是仓库里:它是测试脚手架,不该被 check-shared-libs 之类的 +// 一致性脚本看到,也不该让人误以为是第二个真 worker。 +const STUB_DIR = mkdtempSync(join(tmpdir(), 'pi-pool-test-')); +const STUB = join(STUB_DIR, 'stub-worker.mjs'); +writeFileSync(STUB, ` +process.on('message', (msg) => { + if (msg?.type === 'job') { + const hold = msg.data?.__hold ?? 30; + process.send({ type: 'log', line: 'JOB ' + JSON.stringify({ + mailID: msg.data.mail_id, + kind: msg.kind, + sessionFile: msg.session?.sessionFile || '', + cwd: msg.session?.cwd || '', + grants: msg.grants || [], + lastSyncedName: msg.lastSyncedName || '', + turnTimeoutMs: msg.config?.turnTimeoutMs ?? null, + }) }); + process.send({ type: 'session_opened', piSessionId: 'pi-' + msg.data.mail_id, + sessionFile: '/tmp/f-' + msg.data.mail_id + '.jsonl', cwd: '/tmp', reused: false }); + if (msg.data?.__grant) { + process.send({ type: 'permission_grant', toolName: msg.data.__grant }); + } + if (msg.data?.__name) { + process.send({ type: 'name_synced', signature: msg.data.__name }); + } + if (msg.data?.__reopen) { + // 模型降级换会话:同一个 worker 里第二次 session_opened + process.send({ type: 'session_opened', piSessionId: msg.data.__reopen, + sessionFile: '/tmp/f2.jsonl', cwd: '/tmp', reused: false }); + } + if (msg.data?.__pending) { + process.send({ type: 'permission_pending', relayKey: msg.data.__pending }); + return; // 等决策,见下面的分支 + } + // 每 40ms 报一次心跳:并发的判据必须是「两个进程真的同时在干活」, + // 而不是「running map 里有两个条目」—— fork 返回后立即就有两个条目了。 + const beat = setInterval(() => process.send({ type: 'log', line: 'TICK ' + msg.data.mail_id }), 40); + if (hold < 0) return; // 永不结束,用来验硬超时 + setTimeout(() => { + clearInterval(beat); + process.send({ type: 'done', ok: true, error: '' }); + setTimeout(() => process.exit(0), 20); + }, hold); + return; + } + if (msg?.type === 'permission_decision') { + process.send({ type: 'log', line: 'DECISION ' + msg.relayKey + '=' + msg.decision }); + process.send({ type: 'done', ok: true, error: '' }); + setTimeout(() => process.exit(0), 20); + return; + } + if (msg?.type === 'shutdown') { + process.send({ type: 'log', line: 'SHUTDOWN' }); + setTimeout(() => process.exit(0), 10); + } +}); +process.send({ type: 'ready' }); +`); + +const { createWorkerPool } = await import('../src/pool.mjs'); + +/** 建一个用桩 worker 的池。 */ +function makePool(opts = {}) { + const lines = []; + const pool = createWorkerPool({ + log: (...a) => lines.push(a.join(' ')), + config: () => ({ turnTimeoutMs: 1000, ...(opts.config || {}) }), + onReconfigure: opts.onReconfigure || (() => {}), + maxWorkers: opts.maxWorkers ?? 2, + workerMaxMs: opts.workerMaxMs ?? 5000, + workerPath: opts.workerPath || STUB, + }); + return { pool, lines }; +} + +const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); + +/** 轮询到条件成立或超时 —— 比固定 sleep 稳。 */ +async function until(fn, timeoutMs = 4000) { + const t0 = Date.now(); + while (Date.now() - t0 < timeoutMs) { + if (fn()) return true; + await sleep(20); + } + return false; +} + +/** 从日志行里取出桩 worker 回声的 job 载荷。 */ +function jobs(lines) { + const out = []; + for (const l of lines) { + const i = l.indexOf('JOB '); + if (i === -1) continue; + try { out.push(JSON.parse(l.slice(i + 4))); } catch { /* 不是完整一行 */ } + } + return out; +} + +/** 某个 mailID 的心跳出现过几次。 */ +const ticks = (lines, id) => lines.filter((l) => l.includes(`TICK ${id}`)).length; + +test('并发上限被遵守:第三条会话要等前面空出来', async () => { + const { pool } = makePool({ maxWorkers: 2 }); + for (const id of ['a', 'b', 'c']) { + pool.submit('mail', { mail_id: id, session_id: `S-${id}`, __hold: 250 }); + } + + let peak = 0; + const t = setInterval(() => { peak = Math.max(peak, pool.stats().running); }, 15); + const sawQueue = await until(() => pool.stats().queued > 0, 1000); + await until(() => pool.stats().running === 0 && pool.stats().queued === 0); + clearInterval(t); + pool.stop(); + + assert.ok(peak <= 2, `同时跑的 worker 峰值 ${peak},不该超过 maxWorkers=2`); + assert.ok(sawQueue, '满载时第三封该进队列而不是被丢掉'); +}); + +test('同一会话串行:两个 worker 的心跳不得重叠', async () => { + const { pool, lines } = makePool({ maxWorkers: 3 }); + pool.submit('mail', { mail_id: 'm1', session_id: 'SAME', __hold: 220 }); + pool.submit('mail', { mail_id: 'm2', session_id: 'SAME', __hold: 60 }); + + // 判据一:任一时刻只有一个 worker 在跑。 + let everTwo = false; + const t = setInterval(() => { if (pool.stats().running > 1) everTwo = true; }, 10); + await until(() => jobs(lines).length === 2 && pool.stats().running === 0, 5000); + clearInterval(t); + pool.stop(); + + assert.equal(everTwo, false, '同一条会话不得有两个 worker 同时装载会话文件'); + // 判据二:m2 一次心跳都没能在 m1 结束前发出 —— m1 的心跳数应当远多于 m2。 + assert.ok(ticks(lines, 'm1') >= 3, `m1 该跑满 220ms,实际心跳 ${ticks(lines, 'm1')} 次`); + assert.equal(jobs(lines).length, 2, '两封都要被处理,不能因为串行而丢掉一封'); +}); + +test('不同会话真并发:两个进程的心跳在同一段时间里交错', async () => { + const { pool, lines } = makePool({ maxWorkers: 3 }); + pool.submit('mail', { mail_id: 'p', session_id: 'S-P', __hold: 300 }); + pool.submit('mail', { mail_id: 'q', session_id: 'S-Q', __hold: 300 }); + + const interleaved = await until(() => ticks(lines, 'p') >= 2 && ticks(lines, 'q') >= 2, 2500); + pool.stop(); + assert.ok(interleaved, + `两条不同会话应当并发,实际 p=${ticks(lines, 'p')} q=${ticks(lines, 'q')} 次心跳`); +}); + +test('sessionFile 与 cwd 跨 worker 传下去:第二封接着第一封的会话谈', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 'first', session_id: 'KEEP', __hold: 30 }); + await until(() => jobs(lines).length === 1 && pool.stats().running === 0); + + pool.submit('mail', { mail_id: 'second', session_id: 'KEEP', __hold: 30 }); + await until(() => jobs(lines).length === 2 && pool.stats().running === 0); + pool.stop(); + + const [j1, j2] = jobs(lines); + assert.equal(j1.sessionFile, '', '第一封时还没有会话文件'); + assert.equal(j2.sessionFile, '/tmp/f-first.jsonl', + '第二封必须带上第一封开出来的会话文件,否则每封邮件都从零开始、上下文全丢'); + assert.equal(j2.cwd, '/tmp', 'cwd 也要传下去'); +}); + +test('「一直同意」与命名指纹跨 worker 存活', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { + mail_id: 'g1', session_id: 'GRANT', __hold: 30, + __grant: 'bash', __name: 'platform:某名字|某名字', + }); + await until(() => jobs(lines).length === 1 && pool.stats().running === 0); + + pool.submit('mail', { mail_id: 'g2', session_id: 'GRANT', __hold: 30 }); + await until(() => jobs(lines).length === 2 && pool.stats().running === 0); + pool.stop(); + + const [j1, j2] = jobs(lines); + assert.deepEqual(j1.grants, [], '第一封时还没人点过「一直同意」'); + assert.deepEqual(j2.grants, ['bash'], + '「一直同意」不跨 worker 存活的话,下一封邮件又问一遍 —— 那个选项就是在骗人'); + assert.equal(j2.lastSyncedName, 'platform:某名字|某名字', + '命名指纹要传下去,否则每封邮件都重新 sync 一次'); +}); + +test('config() 每次派活时重取:allowedModels 随心跳变,不能用快照', async () => { + let turnTimeoutMs = 111; + const { pool, lines } = makePool({ maxWorkers: 1, config: {} }); + // makePool 的 config 是固定值,这里换成动态的 + pool.stop(); + + const lines2 = []; + const p2 = createWorkerPool({ + log: (...a) => lines2.push(a.join(' ')), + config: () => ({ turnTimeoutMs }), + onReconfigure: () => {}, + maxWorkers: 1, + workerMaxMs: 5000, + workerPath: STUB, + }); + p2.submit('mail', { mail_id: 'c1', session_id: 'C1', __hold: 20 }); + await until(() => jobs(lines2).length === 1 && p2.stats().running === 0); + turnTimeoutMs = 222; + p2.submit('mail', { mail_id: 'c2', session_id: 'C2', __hold: 20 }); + await until(() => jobs(lines2).length === 2 && p2.stats().running === 0); + p2.stop(); + + const [j1, j2] = jobs(lines2); + assert.equal(j1.turnTimeoutMs, 111); + assert.equal(j2.turnTimeoutMs, 222, 'config() 必须每次重取,否则 worker 用的是上一轮的模型范围'); + assert.equal(lines.length >= 0, true); +}); + +test('硬超时回收卡死的 worker,且不堵住同会话后续邮件', async () => { + const { pool, lines } = makePool({ maxWorkers: 2, workerMaxMs: 400 }); + pool.submit('mail', { mail_id: 'stuck', session_id: 'STUCK', __hold: -1 }); + await until(() => pool.stats().running === 1, 1500); + + const freed = await until(() => pool.stats().running === 0, 3000); + assert.ok(freed, '卡死的 worker 必须被硬超时回收,否则那条会话的后续邮件永远排队'); + assert.ok(lines.some((l) => l.includes('强杀')), `应记下强杀日志,实际:\n${lines.join('\n')}`); + + pool.submit('mail', { mail_id: 'after', session_id: 'STUCK', __hold: 30 }); + const ran = await until(() => jobs(lines).some((j) => j.mailID === 'after'), 2000); + await until(() => pool.stats().running === 0); + pool.stop(); + assert.ok(ran, '硬超时后同一会话的后续邮件必须能被处理'); +}); + +test('权限决策路由到发起询问的那个 worker', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 'perm', session_id: 'PERM', __pending: 'rk-1' }); + await until(() => pool.stats().running === 1, 1500); + await sleep(150); // 等 permission_pending 到主进程 + + assert.equal(pool.routePermission('rk-1', '同意'), true, '应当路由成功'); + await until(() => pool.stats().running === 0, 2000); + pool.stop(); + + assert.ok(lines.some((l) => l.includes('DECISION rk-1=同意')), + `worker 应收到决策原文,实际:\n${lines.join('\n')}`); +}); + +test('决策发的是选项原文而不是归一化的 allow/deny', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 'p2', session_id: 'P2', __pending: 'rk-2' }); + await until(() => pool.stats().running === 1, 1500); + await sleep(150); + pool.routePermission('rk-2', '一直同意'); + await until(() => pool.stats().running === 0, 2000); + pool.stop(); + + // 「同意」与「一直同意」语义不同,归一化会让后者退化成单次授权 + assert.ok(lines.some((l) => l.includes('DECISION rk-2=一直同意')), + `必须原文透传,实际:\n${lines.join('\n')}`); +}); + +test('决策找不到 worker 时返回 false(调用方据此走 B-4.2 降级)', async () => { + const { pool } = makePool(); + assert.equal(pool.routePermission('never-seen', '同意'), false); + pool.stop(); +}); + +test('worker 退出后它的权限路由被清掉,不会误投给下一个 worker', async () => { + const { pool } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 'gone', session_id: 'GONE', __pending: 'rk-gone' }); + await until(() => pool.stats().running === 1, 1500); + await sleep(150); + // 不给决策,直接停掉它 + pool.stop(); + await until(() => pool.stats().running === 0, 4000); + + assert.equal(pool.routePermission('rk-gone', '同意'), false, + 'worker 已退出,路由必须返回 false 让调用方走降级路径'); +}); + +test('mailDrivenIDs 报出所有跑过的 pi 会话,且不随 worker 退出而清', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 'd1', session_id: 'D1', __hold: 30 }); + pool.submit('mail', { mail_id: 'd2', session_id: 'D2', __hold: 30 }); + await until(() => jobs(lines).length === 2 && pool.stats().running === 0); + pool.stop(); + + const ids = pool.mailDrivenIDs(); + assert.ok(ids.has('pi-d1'), 'D1 的 pi 会话该被标记为邮件驱动'); + assert.ok(ids.has('pi-d2'), 'D2 的 pi 会话该被标记为邮件驱动'); +}); + +test('模型降级换掉的旧 pi 会话仍算邮件驱动', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 'r', session_id: 'RETIRE', __hold: 40, __reopen: 'pi-new' }); + await until(() => jobs(lines).length === 1 && pool.stats().running === 0); + pool.stop(); + + const ids = pool.mailDrivenIDs(); + assert.ok(ids.has('pi-new'), '新会话要在'); + assert.ok(ids.has('pi-r'), '被换掉的旧会话也参与过邮件往来,磁盘上的文件还在,快照该报它'); +}); + +test('hasSession 只对跑过的邮件会话为真', async () => { + const { pool } = makePool(); + assert.equal(pool.hasSession('NOPE'), false); + pool.submit('mail', { mail_id: 'h1', session_id: 'HAS', __hold: 30 }); + await until(() => pool.hasSession('HAS'), 2000); + await until(() => pool.stats().running === 0); + pool.stop(); + assert.equal(pool.hasSession('HAS'), true, 'worker 退出后仍该记着这条会话'); +}); + +test('kind 透传:权限通知走 permission 而不是 mail', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('permission', { mail_id: 'k1', session_id: 'K1', __hold: 20 }); + await until(() => jobs(lines).length === 1 && pool.stats().running === 0); + pool.stop(); + assert.equal(jobs(lines)[0].kind, 'permission', + 'kind 决定 worker 用哪套提示词,传错会让模型以为收到一封新邮件'); +}); + +test('stop 之后不再派活', async () => { + const { pool } = makePool(); + pool.stop(); + pool.submit('mail', { mail_id: 'late', session_id: 'LATE', __hold: 30 }); + await sleep(200); + assert.equal(pool.stats().running, 0, '关停后不该再起 worker'); + assert.equal(pool.stats().queued, 0, '关停后队列应为空'); +}); + +test('stop 会先给 worker 发 shutdown(让它 fail closed)再杀', async () => { + const { pool, lines } = makePool({ maxWorkers: 2 }); + pool.submit('mail', { mail_id: 's1', session_id: 'S1', __hold: -1 }); + await until(() => pool.stats().running === 1, 1500); + pool.stop(); + const gotShutdown = await until(() => lines.some((l) => l.includes('SHUTDOWN')), 2000); + assert.ok(gotShutdown, + '必须先发 shutdown:直接 SIGKILL 会让 pi 侧那些等权限的 await 永不返回'); +}); + +test('没有 session_id 的事件各占一个 key,不会互相串行', async () => { + const { pool, lines } = makePool({ maxWorkers: 3 }); + pool.submit('mail', { mail_id: 'n1', __hold: 300 }); + pool.submit('mail', { mail_id: 'n2', __hold: 300 }); + const both = await until(() => ticks(lines, 'n1') >= 2 && ticks(lines, 'n2') >= 2, 2500); + pool.stop(); + assert.ok(both, '无 session_id 的两封不属于同一条会话,应能并发'); +}); + +test('reconfigure 上报被转达给主进程', async () => { + const STUB2 = join(STUB_DIR, 'stub-reconf.mjs'); + writeFileSync(STUB2, ` +process.on('message', (msg) => { + if (msg?.type === 'job') { + process.send({ type: 'reconfigure', url: 'http://new:9999', agentKey: 'k2' }); + process.send({ type: 'done', ok: true, error: '' }); + setTimeout(() => process.exit(0), 20); + } +}); +process.send({ type: 'ready' }); +`); + let got = null; + const { pool } = makePool({ + maxWorkers: 1, + workerPath: STUB2, + onReconfigure: (url, key) => { got = { url, key }; }, + }); + pool.submit('mail', { mail_id: 'r1', session_id: 'R1' }); + await until(() => got !== null, 3000); + pool.stop(); + assert.deepEqual(got, { url: 'http://new:9999', key: 'k2' }, + 'worker 里 connect_to_server 换的坐标必须回到主进程 —— worker 马上就退了,改在它自己身上等于没改'); +});