fix(imap): 修复多监听器广播更新导致 close of closed channel panic
broadcastUpdate 此前把同一个 backend.Update 对象推送到明文/TLS 两个 监听器的推送通道,两个 go-imap listenUpdates goroutine 各自对 update.Done() 执行 close,二次关闭同一 channel 触发 panic (server.go:373 close of closed channel)。 修复:按监听器克隆 Update 对象(载荷共享、Done channel 独立), 新增回归测试:两个监听器收到的更新 Done channel 必须不同,且各自 close 不 panic;全量 -race 通过
This commit is contained in:
@@ -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())
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user