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] }