diff --git a/internal/imap_server/notify_test.go b/internal/imap_server/notify_test.go index b6e0081..8d32df5 100644 --- a/internal/imap_server/notify_test.go +++ b/internal/imap_server/notify_test.go @@ -247,3 +247,53 @@ func TestPushExpunged(t *testing.T) { t.Fatalf("seqs = %v, want [2 5]", seqs) } } + +// TestBroadcastUpdateIsolatedPerListener 回归测试:同一更新广播到多个监听器 +// 时,每个监听器必须持有独立的 Update 对象(独立 Done channel),否则 +// 多个 listenUpdates 会对同一 channel 二次 close 导致 +// panic: close of closed channel。 +func TestBroadcastUpdateIsolatedPerListener(t *testing.T) { + gdb, err := gorm.Open(sqlite.Open(filepath.Join(t.TempDir(), "test.db")), &gorm.Config{}) + if err != nil { + t.Fatalf("open sqlite: %v", err) + } + if err := gdb.AutoMigrate(&db.Message{}); err != nil { + t.Fatalf("migrate: %v", err) + } + stores := store.NewStores(gdb) + + hub := connhub.New() + srv := NewIMAPServer(config.IMAPConfig{}, stores, nil, config.BanConfig{}, hub) + srv.newServer("127.0.0.1:143", nil) + srv.newServer("127.0.0.1:993", nil) + + srv.PushExpunged("alice@example.com", "INBOX", []uint32{1}) + + srv.beMu.Lock() + bes := append([]*imapBackend(nil), srv.bes...) + srv.beMu.Unlock() + + // 每个监听器各收到一条更新 + var updates []backend.Update + for i, b := range bes { + select { + case upd := <-b.updates: + updates = append(updates, upd) + case <-time.After(time.Second): + t.Fatalf("backend %d: no update received", i) + } + } + if len(updates) != 2 { + t.Fatalf("updates = %d, want 2", len(updates)) + } + + // 关键断言:两条更新必须拥有独立的 Done channel + if updates[0].Done() == updates[1].Done() { + t.Fatal("listeners share the same Done channel: double close would panic") + } + + // 模拟两个 listenUpdates 各自执行 close(update.Done()):修复前必 panic + for _, upd := range updates { + close(upd.Done()) + } +} diff --git a/internal/imap_server/server.go b/internal/imap_server/server.go index 8bef425..74e73c1 100644 --- a/internal/imap_server/server.go +++ b/internal/imap_server/server.go @@ -95,6 +95,9 @@ func (s *IMAPServer) PushExpunged(userEmail, mailbox string, seqNums []uint32) { } // broadcastUpdate 把一条更新非阻塞地投递到所有监听器的推送通道。 +// 每个监听器必须收到独立的 Update 对象(各自独立的 Done channel): +// 每个监听器的 listenUpdates 都会对 update.Done() 执行 close,共享 +// 同一对象会导致对同一 channel 二次 close 而 panic。 func (s *IMAPServer) broadcastUpdate(update backend.Update, userEmail string, msgID uint) { s.beMu.Lock() bes := append([]*imapBackend(nil), s.bes...) @@ -102,13 +105,33 @@ func (s *IMAPServer) broadcastUpdate(update backend.Update, userEmail string, ms for _, b := range bes { select { - case b.updates <- update: + case b.updates <- cloneUpdate(update): default: log.Printf("IMAP: 推送通道已满,丢弃 %s 的更新 (msg=%d)", userEmail, msgID) } } } +// cloneUpdate 按类型复制一条 backend.Update:载荷(消息/序号)共享, +// 但 Username/Mailbox/Done channel 重置为独立实例。 +func cloneUpdate(u backend.Update) backend.Update { + switch u := u.(type) { + case *backend.MessageUpdate: + return &backend.MessageUpdate{ + Update: backend.NewUpdate(u.Username(), u.Mailbox()), + Message: u.Message, + } + case *backend.ExpungeUpdate: + return &backend.ExpungeUpdate{ + Update: backend.NewUpdate(u.Username(), u.Mailbox()), + SeqNum: u.SeqNum, + } + default: + // 防御:未知类型原样传递(当前不存在此类更新) + return u + } +} + // registerBackend 记录新建的 backend(用于新邮件推送)。 func (s *IMAPServer) registerBackend(be *imapBackend) { s.beMu.Lock()