Files
MailUI4Agents/plugins/homeagent-mail-bridge/ledger_test.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

315 lines
8.4 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
// 投递账本的单元测试。
//
// 核心不变量只有两条,但它们互相拉扯,必须分别钉住:
//
// 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)
}