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