## 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
315 lines
8.4 KiB
Go
315 lines
8.4 KiB
Go
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)
|
||
}
|