## 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
304 lines
10 KiB
Go
304 lines
10 KiB
Go
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 后 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]
|
||
}
|