Files
MailUI4Agents/plugins/homeagent-mail-bridge/ledger.go
JianFeeeee 8e501f041e pi 桥改为工作进程池 + homeagent 落盘投递账本
## pi 桥:模型工作下到子进程(并发模型重构)

主进程原来自己跑模型,而 pi 的会话装载是同步的:`SessionManager.open()` 走
`openSync` + `readSync` 循环把整个 `.jsonl` 读进内存并逐行 JSON.parse。实测本机
最大那条会话 23MB,`open` 一次**阻塞事件循环 118ms**;模型跑起来后 SDK 内部还有
更多同步工作。SSE 读循环在那期间完全停住 → 后续邮件卡在 TCP 缓冲区 → 久到
Gateway 认为连接死了 → 重连 → 重放。

上一轮我在几个调用点前加 `setImmediate` 是无效的仪式(让出一次之后同步工作照样
占满线程),已回退。这一轮把模型工作整体搬进子进程:实测同样的活在 fork 出的
子进程里跑,主进程事件循环阻塞 **0ms**。

- 新增 `src/worker.mjs`:一封邮件一个进程,跑完就退。权限询问期间的挂起只影响
  那一个 worker(原来 `await new Promise(...)` 等人决策,整座桥不再收信)。
- 新增 `src/pool.mjs`:**不同会话并发**(上限 3,每个 worker 约 140MB RSS)、
  **同一会话严格串行**(pi 假定「一文件一持有者」,两个进程同时装载同一条会话
  文件会让各自的内存索引看不见对方追加的行 → 会话树分叉)、满载排队不丢邮件、
  硬超时 SIGKILL 回收卡死进程。
- `src/index.mjs` 只剩 I/O 与调度:SSE、心跳、去重、分派。
- 选进程而不是 `worker_threads`:模型会跑 bash/write/edit,一次 OOM 不该带走
  整座桥。两者实测都能建起 AgentSession,但线程与主线程共享堆和生命周期。
- 「接管会话短暂持有」那套机制(adopted / adoptTimers / releaseAdopted + 兜底
  计时器)整个删掉 —— worker 退出**就是**释放,且普通会话与接管会话一视同仁。
- 轮次超时 60s → 10 分钟:60s 那个数字是「主进程要腾出手收下一封」的产物,
  worker 没有这个理由,等真结论更准(带工具调用的一轮跑几分钟很正常)。
- IPC 只传路径与标量(sessionFile / cwd / grants / 命名指纹)—— AgentSession
  跨不了进程边界,worker 每次从 sessionFile 重新装载。

`test/pool.test.mjs` +19 例,真 fork 子进程、用桩 worker(不装 SDK)跑毫秒级:
并发上限、同会话串行、不同会话真并发(判据是两个进程的心跳交错,不是 running
map 里有两个条目)、sessionFile/grants/命名指纹跨 worker 传递、config() 每次重取、
硬超时回收、权限决策路由、决策原文透传、worker 退出后清路由、mailDrivenIDs、
kind 透传、stop 先发 shutdown 再杀。跑过三组负向对照确认用例真能抓回归:
拆掉串行守卫 / 不传 sessionFile+grants / 硬超时不杀,对应用例分别失败。

## homeagent:投递去重必须落盘

用户报的重复投递不是上一轮那个 bug。两段提示词的措辞差异指出了来源:
SSE 那段写「你把本轮工作做完」,补投那段写「你把结论说出来就行」。

`deliveredMails` 是进程内的 map,而 homeagent 的插件跑在**子进程**里:

  1. 18:59:38 邮件落库,旧插件进程的 SSE 收到,注入第一次
  2. 同一秒 homed 被重启,那一轮被掐断(`context canceled`)
  3. 18:59:45 新进程起来,`deliveredMails` 是空的
  4. 心跳报 `pending_mails: 1`(第一轮没跑完 → read_inbox 没执行 → 仍未读)
     → catchUp 注入第二次

**不能只记「投过没有」**:那会把「重复」换成「丢件」—— 第 2 步里发件人没收到
回信,而记录说「已投过」→ 永远跳过。丢件比重复严重,重复至少人能看出来。

新增 `ledger.go`:JSONL 账本记两个状态。`completed` 才跳过;`delivered` 但未
`completed` 的仍然重投,但提示词前面插一段说明「上一轮被中断,别把同一件事做
两次」。落在 SDK 的 `Settings().DataDir()`;拿不到时退回 key 文件目录;目录不可
写时退化为纯内存(不比修复前差,也不该让插件起不来)。

- 判定与记录在同一把锁里:SSE 与 catchUp 两个 goroutine 的竞态
- 每行写完 fsync:这个文件的全部意义就是「进程死了之后还算数」
- 坏行跳过而不是报错退出(崩溃时最后一行可能写残)→ 那封退化为重投,安全
- 14 天保留期;过期过半时「临时文件 + rename」压实
- `shortID()` 替代 `id[:8]`:日志不该有能力 panic 掉投递协程

`ledger_test.go` +14 例,含两组负向对照(只记「投过」→ 丢件用例失败;不读账本
→ 跨进程用例失败)。

## 契约文档

`B-7.7`(MUST):子进程形式的插件去重必须落盘且区分「投过」与「跑完」,含事故
时序、两状态表、何时标 completed。已知取舍那节标注投递账本是唯一必须落盘的状态。
验收清单加「模型跑到一半重启宿主」一项。

## 生产验证

- pi 三封 → 三条会话:三个 worker PID 并存,回信「收到 1/2/3」各落自己线索
- pi 同一会话两封:严格串行(收到A 19:26:39 → 收到B 19:26:48,全程单 worker)
- pi 主进程事件循环阻塞 1ms(旧版单进程 open 23MB 一次就 118ms)
- homeagent 正常一封:账本 `c:false` → `c:true`,一封回信
- homeagent 处理中重启:日志「上一轮被中断,带说明重投」,**只有一封 Re:**
- homeagent 再次重启:账本 2 条 completed,不再投递,会话邮件数不变
- gateway 7 包 / web 176+26 / pi 269 / dsh 219 / opencode 201 / homeagent 14
2026-09-04 20:33:44 +08:00

304 lines
10 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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()``<daemon data>/plugin_data/<plugin>`
// 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 后 completedcompleted 只增不减。
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 后再 renamerename 本身是原子的,但内容没落盘时掉电会得到空文件。
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]
}