fix(imap): 升级 go-imap v2 + \Deleted 落库,修复删除邮件只复制到垃圾桶
删除流程(Apple Mail/iOS 等):UID COPY→Trash + UID STORE \Deleted + UID EXPUNGE。此前 go-imap v1.2.1 不支持 UID EXPUNGE(回 Command unsupported with UID),且 \Deleted 只存内存会话 map(重选/ 重连即丢失),导致 EXPUNGE 不删原邮件——垃圾桶只有副本、INBOX 原封 不动。 - 依赖:go-imap v1.2.1 → go-imap/v2 v2.0.0-beta.8(原生支持 UID EXPUNGE,且库为 race-clean,移除集成测试 !race 标签) - Message 新增 IsDeleted 字段持久化 \Deleted 标记(AutoMigrate 自动 加列,SQLite/MySQL 通用;migrate 工具随模型字段自动复制) - IMAP 层重写为 v2 imapserver.Session 架构(session.go 替代 backend.go):SELECT/STATUS/LIST/FETCH/SEARCH/STORE/COPY/MOVE/ APPEND/EXPUNGE/Poll/Idle;EXPUNGE 按数据库 IsDeleted 删除,UID EXPUNGE 只删集合内已标记消息 - server.go:显式能力集(IMAP4rev1+UIDPLUS+MOVE+LITERAL++CHILDREN+ SPECIAL-USE);自实现 mailboxHub 跨会话推送(含来源排除,避免 v2 库 MailboxTracker 对 EXPUNGE/EXISTS 的回声重复响应);Pusher 接口 签名不变,Web/SMTP/POP3 调用方零改动;DisconnectByAddr 改用会话 注册表 + conn.Bye() - 测试:集成测试改用 v2 imapclient,新增删除流程/跨会话持久化/MOVE 回归用例;notify_test 适配 hub 推送模型并覆盖来源排除与 allowExpunge 语义
This commit is contained in:
@@ -65,6 +65,11 @@ type Message struct {
|
||||
RawData string `gorm:"type:mediumtext" json:"raw_data"`
|
||||
IsRead bool `gorm:"default:false" json:"is_read"`
|
||||
IsFlagged bool `gorm:"default:false" json:"is_flagged"`
|
||||
// IsDeleted 持久化 IMAP \Deleted 标记(STORE +FLAGS \Deleted 写入,
|
||||
// EXPUNGE/UID EXPUNGE 时按此删除)。此前该标记只存于内存会话
|
||||
// (imapMailbox.deleted map),重选文件夹/重连即丢失,导致客户端
|
||||
// COPY→Trash 后 EXPUNGE 删不掉原邮件(垃圾桶只有副本)。
|
||||
IsDeleted bool `gorm:"default:false;index" json:"is_deleted"`
|
||||
Date time.Time `json:"date"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,10 +1,6 @@
|
||||
//go:build !race
|
||||
|
||||
// 集成测试:启动真实 IMAP 监听 + 脚本客户端(go-imap client)。
|
||||
// 注意:仅在非 -race 构建下运行——go-imap v1.2.1 存在库内数据竞争
|
||||
// (cmd_selected.go STORE 写 *conn.silent() vs listenUpdates 读),
|
||||
// 启用 backend 推送(Updates != nil)时必然触发,-race 下会误报。
|
||||
// 推送逻辑的竞态覆盖由单元测试(notify_test.go)承担。
|
||||
// 集成测试:启动真实 IMAP 监听 + 脚本客户端(go-imap v2 imapclient)。
|
||||
// v2 库为 race-clean(此前 v1.2.1 库内数据竞争导致本文件必须带
|
||||
// !race 构建标签,升级后已移除)。推送逻辑的单元覆盖由 notify_test.go 承担。
|
||||
package imap_server
|
||||
|
||||
import (
|
||||
@@ -18,8 +14,8 @@ import (
|
||||
"mail_go/internal/db"
|
||||
"mail_go/internal/store"
|
||||
|
||||
"github.com/emersion/go-imap"
|
||||
"github.com/emersion/go-imap/client"
|
||||
"github.com/emersion/go-imap/v2"
|
||||
"github.com/emersion/go-imap/v2/imapclient"
|
||||
"golang.org/x/crypto/bcrypt"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
@@ -32,7 +28,7 @@ func startIntegrationServer(t *testing.T) (*store.Stores, string) {
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
if err := gdb.AutoMigrate(&db.User{}, &db.Domain{}, &db.Message{}, &db.ProtocolLog{}, &db.BanEntry{}); err != nil {
|
||||
if err := gdb.AutoMigrate(&db.User{}, &db.Domain{}, &db.Message{}, &db.ProtocolLog{}, &db.BanEntry{}, &db.MailboxState{}); err != nil {
|
||||
t.Fatalf("migrate: %v", err)
|
||||
}
|
||||
stores := store.NewStores(gdb)
|
||||
@@ -83,23 +79,22 @@ func seedMailbox(t *testing.T, stores *store.Stores, userID uint, n int) []uint
|
||||
return ids
|
||||
}
|
||||
|
||||
func loginAndSelect(t *testing.T, addr string) *client.Client {
|
||||
func loginAndSelect(t *testing.T, addr string) *imapclient.Client {
|
||||
t.Helper()
|
||||
c, err := client.Dial(addr)
|
||||
c, err := imapclient.DialInsecure(addr, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { c.Logout() })
|
||||
if err := c.Login("alice@example.com", "secret123"); err != nil {
|
||||
t.Cleanup(func() { c.Logout().Wait() })
|
||||
if err := c.Login("alice@example.com", "secret123").Wait(); err != nil {
|
||||
t.Fatalf("login: %v", err)
|
||||
}
|
||||
if _, err := c.Select("INBOX", false); err != nil {
|
||||
if _, err := c.Select("INBOX", nil).Wait(); err != nil {
|
||||
t.Fatalf("select: %v", err)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// uidOf returns the message id of the newest message by date.
|
||||
func assertReadState(t *testing.T, stores *store.Stores, msgID uint, want bool) {
|
||||
t.Helper()
|
||||
msg, err := stores.Mails.GetByID(msgID)
|
||||
@@ -117,13 +112,13 @@ func TestUidStorePersists(t *testing.T) {
|
||||
ids := seedMailbox(t, stores, 1, 3)
|
||||
|
||||
c := loginAndSelect(t, addr)
|
||||
seqset := new(imap.SeqSet)
|
||||
seqset.AddNum(uint32(ids[1])) // UID = 第二条消息
|
||||
ch := make(chan *imap.Message, 1)
|
||||
if err := c.UidStore(seqset, imap.AddFlags, []interface{}{imap.SeenFlag}, ch); err != nil {
|
||||
cmd := c.Store(imap.UIDSetNum(imap.UID(ids[1])), &imap.StoreFlags{
|
||||
Op: imap.StoreFlagsAdd,
|
||||
Flags: []imap.Flag{imap.FlagSeen},
|
||||
}, nil)
|
||||
if _, err := cmd.Collect(); err != nil {
|
||||
t.Fatalf("uid store: %v", err)
|
||||
}
|
||||
<-ch
|
||||
|
||||
assertReadState(t, stores, ids[1], true)
|
||||
assertReadState(t, stores, ids[0], false)
|
||||
@@ -139,15 +134,13 @@ func TestSeqStoreServerIssued(t *testing.T) {
|
||||
c := loginAndSelect(t, addr)
|
||||
|
||||
// 拉取全部消息,找到 ids[2](最新一封)的服务器序号
|
||||
seqsetAll := new(imap.SeqSet)
|
||||
seqsetAll.AddRange(1, 3)
|
||||
messages := make(chan *imap.Message, 3)
|
||||
if err := c.Fetch(seqsetAll, []imap.FetchItem{imap.FetchFlags, imap.FetchUid}, messages); err != nil {
|
||||
msgs, err := c.Fetch(imap.SeqSetNum(1, 2, 3), &imap.FetchOptions{UID: true}).Collect()
|
||||
if err != nil {
|
||||
t.Fatalf("fetch: %v", err)
|
||||
}
|
||||
var targetSeq uint32
|
||||
for m := range messages {
|
||||
if m.Uid == uint32(ids[2]) {
|
||||
for _, m := range msgs {
|
||||
if m.UID == imap.UID(ids[2]) {
|
||||
targetSeq = m.SeqNum
|
||||
}
|
||||
}
|
||||
@@ -155,13 +148,13 @@ func TestSeqStoreServerIssued(t *testing.T) {
|
||||
t.Fatal("target message not found in fetch")
|
||||
}
|
||||
|
||||
seqset := new(imap.SeqSet)
|
||||
seqset.AddNum(targetSeq)
|
||||
ch := make(chan *imap.Message, 1)
|
||||
if err := c.Store(seqset, imap.AddFlags, []interface{}{imap.SeenFlag}, ch); err != nil {
|
||||
cmd := c.Store(imap.SeqSetNum(targetSeq), &imap.StoreFlags{
|
||||
Op: imap.StoreFlagsAdd,
|
||||
Flags: []imap.Flag{imap.FlagSeen},
|
||||
}, nil)
|
||||
if _, err := cmd.Collect(); err != nil {
|
||||
t.Fatalf("store: %v", err)
|
||||
}
|
||||
<-ch
|
||||
|
||||
assertReadState(t, stores, ids[2], true)
|
||||
}
|
||||
@@ -177,13 +170,13 @@ func TestSeqStoreClientSelfNumbered(t *testing.T) {
|
||||
c := loginAndSelect(t, addr)
|
||||
|
||||
// 客户端按日期倒序视图:最新一封 = seq 1
|
||||
seqset := new(imap.SeqSet)
|
||||
seqset.AddNum(1)
|
||||
ch := make(chan *imap.Message, 1)
|
||||
if err := c.Store(seqset, imap.AddFlags, []interface{}{imap.SeenFlag}, ch); err != nil {
|
||||
cmd := c.Store(imap.SeqSetNum(1), &imap.StoreFlags{
|
||||
Op: imap.StoreFlagsAdd,
|
||||
Flags: []imap.Flag{imap.FlagSeen},
|
||||
}, nil)
|
||||
if _, err := cmd.Collect(); err != nil {
|
||||
t.Fatalf("store: %v", err)
|
||||
}
|
||||
<-ch
|
||||
|
||||
// 客户端意图是标记最新一封(ids[2])为已读
|
||||
assertReadState(t, stores, ids[2], true)
|
||||
@@ -191,14 +184,12 @@ func TestSeqStoreClientSelfNumbered(t *testing.T) {
|
||||
|
||||
// TestFetchBodyMalformedMIME 回归:消息包含无法解析的 MIME(base64 编码的
|
||||
// message/rfc822 附件 / 截断的 multipart)时,FETCH BODY/BODYSTRUCTURE
|
||||
// 不得因 nil BodyStructure 触发服务器 panic(否则连接中断,客户端只收到
|
||||
// 部分邮件或一直卡在同步)。修复前 go-imap send() 协程会 nil 指针崩溃。
|
||||
// 不得 panic(v2 的 ExtractBodyStructure 对畸形 MIME 返回降级结构,
|
||||
// WriteBodyStructure 要求 Extended 非 nil,两者都需满足)。
|
||||
func TestFetchBodyMalformedMIME(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(t)
|
||||
|
||||
// 1) base64 编码的 message/rfc822 附件(转发邮件场景):
|
||||
// backendutil.FetchBodyStructure 不解码 base64,直接把编码文本
|
||||
// 当嵌套消息头解析 → "malformed MIME header line" 错误。
|
||||
// 1) base64 编码的 message/rfc822 附件(转发邮件场景)
|
||||
rfc822Body := "UmVjZWl2ZWQ6IGZyb20gb3V0Ym91bmQuY2kuaWNsb3VkLmNvbSAodW5rbm93biBbMTI3LjAuMC4yKVxuXHQgYnkgcDAwLWljbG91ZG10YS1hc210cC11cy1jZW50cmFsLTFrLTEwMC1wZXJjZW50LTggKFBvc3RmaXgpIHdpdGggRVNNVFBTIGlkIDIxRTlBMThDQURDRjM4MlxuXHQgZm9yIDxkc2hAbG12ZS5uZXQ+OyBTdW4sIDE2IEF1ZyAyMDI2IDEzOjU4OjIxICswMDAwIChVVEMpXG5YLUlDTC1SZXBJZDogRURWY1BlQ3RlWG4tZ0Z1T0xxUWhfSjZvcE9fN1B2OEtsOW1mMDg2VUFxZ29zXG5EYXRlOiBTdW4sIDE2IEF1ZyAyMDI2IDEzOjU4OjIxICswMDAwXG5Gcm9tOiBkYXZpZEB5YW5kZXguY29tXG5UbzogZHNoQGxtdmUubmV0XG5NZXNzYWdlLUlEOiA8QTIxNzBEMTEtMkI1MC00MTQwLTlEQTMtMkI3M0U2RUIwQTc4QHlhbmRleC5jb20+XG5TdWJqZWN0OiB0ZXN0XG5cbmhlbGxvXG4="
|
||||
msgWithRFC822 := &db.Message{
|
||||
UserID: 1,
|
||||
@@ -224,7 +215,7 @@ func TestFetchBodyMalformedMIME(t *testing.T) {
|
||||
rfc822Body + "\r\n" +
|
||||
"--==fwd==--\r\n",
|
||||
}
|
||||
// 2) 截断的 multipart(缺少结束边界):BODYSTRUCTURE(extended) 解析报错
|
||||
// 2) 截断的 multipart(缺少结束边界)
|
||||
msgTruncated := &db.Message{
|
||||
UserID: 1,
|
||||
Folder: "INBOX",
|
||||
@@ -251,32 +242,138 @@ func TestFetchBodyMalformedMIME(t *testing.T) {
|
||||
|
||||
c := loginAndSelect(t, addr)
|
||||
|
||||
seqset := new(imap.SeqSet)
|
||||
seqset.AddRange(1, 2)
|
||||
seqSet := imap.SeqSetNum(1, 2)
|
||||
|
||||
// BODY:历史上 message/rfc822 消息解析失败 → nil BodyStructure → panic
|
||||
msgs := make(chan *imap.Message, 10)
|
||||
if err := c.Fetch(seqset, []imap.FetchItem{imap.FetchBody}, msgs); err != nil {
|
||||
// BODY(非扩展)
|
||||
msgs, err := c.Fetch(seqSet, &imap.FetchOptions{BodyStructure: &imap.FetchItemBodyStructure{}}).Collect()
|
||||
if err != nil {
|
||||
t.Fatalf("fetch body: %v", err)
|
||||
}
|
||||
got := 0
|
||||
for range msgs {
|
||||
got++
|
||||
}
|
||||
if got != 2 {
|
||||
t.Fatalf("FETCH BODY 返回 %d/2 封", got)
|
||||
if len(msgs) != 2 {
|
||||
t.Fatalf("FETCH BODY 返回 %d/2 封", len(msgs))
|
||||
}
|
||||
|
||||
// BODYSTRUCTURE:截断 multipart 在 extended 解析时报错 → nil → panic
|
||||
msgs2 := make(chan *imap.Message, 10)
|
||||
if err := c.Fetch(seqset, []imap.FetchItem{imap.FetchBodyStructure}, msgs2); err != nil {
|
||||
// BODYSTRUCTURE(扩展)
|
||||
msgs2, err := c.Fetch(seqSet, &imap.FetchOptions{BodyStructure: &imap.FetchItemBodyStructure{Extended: true}}).Collect()
|
||||
if err != nil {
|
||||
t.Fatalf("fetch bodystructure: %v", err)
|
||||
}
|
||||
got2 := 0
|
||||
for range msgs2 {
|
||||
got2++
|
||||
}
|
||||
if got2 != 2 {
|
||||
t.Fatalf("FETCH BODYSTRUCTURE 返回 %d/2 封", got2)
|
||||
if len(msgs2) != 2 {
|
||||
t.Fatalf("FETCH BODYSTRUCTURE 返回 %d/2 封", len(msgs2))
|
||||
}
|
||||
}
|
||||
|
||||
// TestUidExpungeDeleteFlow 核心回归:模拟 Apple Mail / iOS Mail 的删除流程
|
||||
// (UID COPY → Trash + UID STORE \Deleted + UID EXPUNGE)。修复前:
|
||||
// \Deleted 只存内存会话、UID EXPUNGE 不被 go-imap v1 支持 → 原邮件留在
|
||||
// INBOX,垃圾桶只有副本。修复后:原邮件被真正删除,Trash 留有一份副本。
|
||||
func TestUidExpungeDeleteFlow(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(t)
|
||||
ids := seedMailbox(t, stores, 1, 3)
|
||||
uid := imap.UID(ids[1])
|
||||
|
||||
c := loginAndSelect(t, addr)
|
||||
|
||||
// 1) UID COPY → Trash
|
||||
if _, err := c.Copy(imap.UIDSetNum(uid), "Trash").Wait(); err != nil {
|
||||
t.Fatalf("uid copy: %v", err)
|
||||
}
|
||||
|
||||
// 2) UID STORE +FLAGS.SILENT (\Deleted)
|
||||
if _, err := c.Store(imap.UIDSetNum(uid), &imap.StoreFlags{
|
||||
Op: imap.StoreFlagsAdd,
|
||||
Silent: true,
|
||||
Flags: []imap.Flag{imap.FlagDeleted},
|
||||
}, nil).Collect(); err != nil {
|
||||
t.Fatalf("uid store deleted: %v", err)
|
||||
}
|
||||
|
||||
// 3) UID EXPUNGE
|
||||
seqs, err := c.UIDExpunge(imap.UIDSetNum(uid)).Collect()
|
||||
if err != nil {
|
||||
t.Fatalf("uid expunge: %v", err)
|
||||
}
|
||||
if len(seqs) != 1 {
|
||||
t.Fatalf("uid expunge seqs = %v, want 1 条", seqs)
|
||||
}
|
||||
|
||||
// 原邮件必须被真正删除
|
||||
if _, err := stores.Mails.GetByID(ids[1]); err == nil {
|
||||
t.Fatal("原邮件仍在数据库,UID EXPUNGE 未生效")
|
||||
}
|
||||
// 其余邮件保留
|
||||
assertReadState(t, stores, ids[0], false)
|
||||
assertReadState(t, stores, ids[2], false)
|
||||
// Trash 中有一份副本
|
||||
trashCount, err := stores.Mails.CountByUserAndFolder(1, "Trash")
|
||||
if err != nil || trashCount != 1 {
|
||||
t.Fatalf("Trash count = %d, want 1", trashCount)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeletedPersistsAcrossSessions 验证 \Deleted 落库:标记后断开重连
|
||||
// (模拟客户端切换文件夹/重连),普通 EXPUNGE 仍能删除。修复前标记只存
|
||||
// 内存会话对象,重连后 EXPUNGE 什么都不删。
|
||||
func TestDeletedPersistsAcrossSessions(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(t)
|
||||
ids := seedMailbox(t, stores, 1, 3)
|
||||
uid := imap.UID(ids[1])
|
||||
|
||||
c := loginAndSelect(t, addr)
|
||||
if _, err := c.Store(imap.UIDSetNum(uid), &imap.StoreFlags{
|
||||
Op: imap.StoreFlagsAdd,
|
||||
Silent: true,
|
||||
Flags: []imap.Flag{imap.FlagDeleted},
|
||||
}, nil).Collect(); err != nil {
|
||||
t.Fatalf("uid store deleted: %v", err)
|
||||
}
|
||||
if err := c.Logout().Wait(); err != nil {
|
||||
t.Fatalf("logout: %v", err)
|
||||
}
|
||||
|
||||
// 重新连接(新的会话对象,此前标记必须仍在库中)
|
||||
c2, err := imapclient.DialInsecure(addr, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("dial2: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { c2.Logout().Wait() })
|
||||
if err := c2.Login("alice@example.com", "secret123").Wait(); err != nil {
|
||||
t.Fatalf("login2: %v", err)
|
||||
}
|
||||
if _, err := c2.Select("INBOX", nil).Wait(); err != nil {
|
||||
t.Fatalf("select2: %v", err)
|
||||
}
|
||||
seqs, err := c2.Expunge().Collect()
|
||||
if err != nil {
|
||||
t.Fatalf("expunge: %v", err)
|
||||
}
|
||||
if len(seqs) != 1 {
|
||||
t.Fatalf("expunge seqs = %v, want 1 条", seqs)
|
||||
}
|
||||
if _, err := stores.Mails.GetByID(ids[1]); err == nil {
|
||||
t.Fatal("标记 \\Deleted 的邮件在重连后未被 EXPUNGE 删除")
|
||||
}
|
||||
}
|
||||
|
||||
// TestUidMoveFlow 验证 UID MOVE:目标文件夹出现该邮件且源文件夹删除。
|
||||
func TestUidMoveFlow(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(t)
|
||||
ids := seedMailbox(t, stores, 1, 2)
|
||||
uid := imap.UID(ids[1])
|
||||
|
||||
c := loginAndSelect(t, addr)
|
||||
if _, err := c.Move(imap.UIDSetNum(uid), "Trash").Wait(); err != nil {
|
||||
t.Fatalf("uid move: %v", err)
|
||||
}
|
||||
|
||||
if _, err := stores.Mails.GetByID(ids[1]); err != nil {
|
||||
t.Fatalf("moved message missing: %v", err)
|
||||
}
|
||||
msg, err := stores.Mails.GetByID(ids[1])
|
||||
if err != nil {
|
||||
t.Fatalf("get moved msg: %v", err)
|
||||
}
|
||||
if msg.Folder != "Trash" {
|
||||
t.Fatalf("moved msg folder = %q, want Trash", msg.Folder)
|
||||
}
|
||||
}
|
||||
+125
-217
@@ -10,18 +10,24 @@ import (
|
||||
"mail_go/internal/db"
|
||||
"mail_go/internal/store"
|
||||
|
||||
"github.com/emersion/go-imap/backend"
|
||||
"github.com/emersion/go-imap/v2"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// TestPushNewMessage 验证本地投递成功后推送的 MessageUpdate 内容正确。
|
||||
func TestPushNewMessage(t *testing.T) {
|
||||
// fakeSession 构造一个挂接在推送中心上的裸会话(无网络连接),用于
|
||||
// 验证 Pusher 的跨会话推送内容与来源排除语义。
|
||||
func fakeSession() *imapSession {
|
||||
return &imapSession{notify: make(chan struct{}, 1)}
|
||||
}
|
||||
|
||||
func newTestServer(t *testing.T) (*IMAPServer, *store.Stores) {
|
||||
t.Helper()
|
||||
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.User{}, &db.Domain{}, &db.Message{}); err != nil {
|
||||
if err := gdb.AutoMigrate(&db.User{}, &db.Domain{}, &db.Message{}, &db.MailboxState{}); err != nil {
|
||||
t.Fatalf("migrate: %v", err)
|
||||
}
|
||||
stores := store.NewStores(gdb)
|
||||
@@ -34,108 +40,51 @@ func TestPushNewMessage(t *testing.T) {
|
||||
if err := stores.Users.Create(user); err != nil {
|
||||
t.Fatalf("create user: %v", err)
|
||||
}
|
||||
email := "alice@example.com"
|
||||
|
||||
// 已有一封旧邮件(日期更早);规范排序最新在前,新邮件应为 INBOX 第 1 封
|
||||
old := &db.Message{UserID: user.ID, Folder: "INBOX", FromAddr: "x@y", Subject: "old", Date: time.Now().Add(-time.Hour)}
|
||||
srv := NewIMAPServer(config.IMAPConfig{}, stores, nil, config.BanConfig{}, connhub.New())
|
||||
return srv, stores
|
||||
}
|
||||
|
||||
// TestPushNewMessage 验证本地投递成功后推送给已选中会话的 EXISTS 更新。
|
||||
func TestPushNewMessage(t *testing.T) {
|
||||
srv, stores := newTestServer(t)
|
||||
|
||||
// 已有一封旧邮件;新邮件投递后 INBOX 计数应为 2
|
||||
old := &db.Message{UserID: 1, Folder: "INBOX", FromAddr: "x@y", Subject: "old", Date: time.Now().Add(-time.Hour)}
|
||||
if err := stores.Mails.Create(old); err != nil {
|
||||
t.Fatalf("create old message: %v", err)
|
||||
}
|
||||
inboxMsg := &db.Message{
|
||||
UserID: user.ID,
|
||||
Folder: "INBOX",
|
||||
FromAddr: "sender@other.com",
|
||||
ToAddr: email,
|
||||
Subject: "新邮件",
|
||||
RawData: "From: sender@other.com\r\nSubject: 新邮件\r\n\r\nhello",
|
||||
MessageID: "<new-1@other.com>",
|
||||
Date: time.Now(),
|
||||
IsRead: false,
|
||||
UserID: 1,
|
||||
Folder: "INBOX",
|
||||
FromAddr: "sender@other.com",
|
||||
ToAddr: "alice@example.com",
|
||||
Subject: "新邮件",
|
||||
Date: time.Now(),
|
||||
}
|
||||
if err := stores.Mails.Create(inboxMsg); err != nil {
|
||||
t.Fatalf("create message: %v", err)
|
||||
}
|
||||
|
||||
hub := connhub.New()
|
||||
srv := NewIMAPServer(config.IMAPConfig{}, stores, nil, config.BanConfig{}, hub)
|
||||
// 模拟明文 + TLS 两个监听器(生产环境由 Start/StartTLS 注册)
|
||||
srv.newServer("127.0.0.1:143", nil)
|
||||
srv.newServer("127.0.0.1:993", nil)
|
||||
srv.PushNewMessage(email, inboxMsg)
|
||||
hub := srv.hubForOrCreate("alice@example.com", "INBOX")
|
||||
sess := fakeSession()
|
||||
hub.add(sess)
|
||||
|
||||
// 两个监听器(明文/TLS)各有一个 backend 通道,都应收到同一更新
|
||||
srv.beMu.Lock()
|
||||
bes := append([]*imapBackend(nil), srv.bes...)
|
||||
srv.beMu.Unlock()
|
||||
if len(bes) == 0 {
|
||||
t.Fatal("no backends registered")
|
||||
srv.PushNewMessage("alice@example.com", inboxMsg)
|
||||
|
||||
updates := sess.takeUpdates(true)
|
||||
if len(updates) != 1 || updates[0].exists == nil {
|
||||
t.Fatalf("updates = %+v, want 1 条 EXISTS", updates)
|
||||
}
|
||||
|
||||
for i, b := range bes {
|
||||
select {
|
||||
case upd := <-b.updates:
|
||||
mu, ok := upd.(*backend.MessageUpdate)
|
||||
if !ok {
|
||||
t.Fatalf("backend %d: update type = %T, want *MessageUpdate", i, upd)
|
||||
}
|
||||
if mu.Username() != email {
|
||||
t.Fatalf("backend %d: username = %q, want %q", i, mu.Username(), email)
|
||||
}
|
||||
if mu.Mailbox() != "INBOX" {
|
||||
t.Fatalf("backend %d: mailbox = %q, want INBOX", i, mu.Mailbox())
|
||||
}
|
||||
if mu.Message.Uid != uint32(inboxMsg.ID) {
|
||||
t.Fatalf("backend %d: uid = %d, want %d", i, mu.Message.Uid, inboxMsg.ID)
|
||||
}
|
||||
if mu.Message.SeqNum != 1 {
|
||||
t.Fatalf("backend %d: seq = %d, want 1", i, mu.Message.SeqNum)
|
||||
}
|
||||
if mu.Message.Envelope == nil || mu.Message.Envelope.Subject != "新邮件" {
|
||||
t.Fatalf("backend %d: envelope missing subject", i)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatalf("backend %d: no update received", i)
|
||||
}
|
||||
if *updates[0].exists != 2 {
|
||||
t.Fatalf("EXISTS = %d, want 2", *updates[0].exists)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPushNewMessageChannelFull 验证通道满时推送不阻塞(非阻塞丢弃)。
|
||||
func TestPushNewMessageChannelFull(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.User{}, &db.Domain{}, &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)
|
||||
|
||||
msg := &db.Message{ID: 1, Folder: "INBOX", Date: time.Now()}
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
// 灌满所有 backend 通道(容量 256),再调用必须立即返回
|
||||
srv.beMu.Lock()
|
||||
bes := append([]*imapBackend(nil), srv.bes...)
|
||||
srv.beMu.Unlock()
|
||||
for _, b := range bes {
|
||||
for i := 0; i < cap(b.updates); i++ {
|
||||
b.updates <- backend.NewUpdate("a@b", "INBOX")
|
||||
}
|
||||
}
|
||||
srv.PushNewMessage("a@b", msg)
|
||||
close(done)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("PushNewMessage blocked on full channel")
|
||||
}
|
||||
// TestPushNewMessageNoSession 验证无会话选中时推送为 no-op(不 panic)。
|
||||
func TestPushNewMessageNoSession(t *testing.T) {
|
||||
srv, _ := newTestServer(t)
|
||||
srv.PushNewMessage("alice@example.com", &db.Message{ID: 1, UserID: 1, Folder: "INBOX", Date: time.Now()})
|
||||
}
|
||||
|
||||
// TestPushNewMessageNilSafe 验证空参数/空指针安全。
|
||||
@@ -145,155 +94,114 @@ func TestPushNewMessageNilSafe(t *testing.T) {
|
||||
srv = NewIMAPServer(config.IMAPConfig{}, nil, nil, config.BanConfig{}, nil)
|
||||
srv.PushNewMessage("", &db.Message{ID: 1}) // 空邮箱
|
||||
srv.PushNewMessage("a@b", nil) // 空消息
|
||||
srv.PushFlagsChanged("", "", nil)
|
||||
srv.PushExpunged("", "", nil)
|
||||
}
|
||||
|
||||
// TestPushFlagsChanged 验证标志变化(已读/星标)推送内容正确。
|
||||
func TestPushFlagsChanged(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.User{}, &db.Domain{}, &db.Message{}); err != nil {
|
||||
t.Fatalf("migrate: %v", err)
|
||||
}
|
||||
stores := store.NewStores(gdb)
|
||||
srv, stores := newTestServer(t)
|
||||
|
||||
domain := &db.Domain{Name: "example.com"}
|
||||
if err := stores.Domains.Create(domain); err != nil {
|
||||
t.Fatalf("create domain: %v", err)
|
||||
}
|
||||
user := &db.User{Username: "alice", DomainID: domain.ID, IsActive: true}
|
||||
if err := stores.Users.Create(user); err != nil {
|
||||
t.Fatalf("create user: %v", err)
|
||||
}
|
||||
|
||||
msg := &db.Message{UserID: user.ID, Folder: "INBOX", FromAddr: "x@y", Subject: "s", Date: time.Now()}
|
||||
msg := &db.Message{UserID: 1, Folder: "INBOX", FromAddr: "x@y", Subject: "s", Date: time.Now()}
|
||||
if err := stores.Mails.Create(msg); err != nil {
|
||||
t.Fatalf("create message: %v", err)
|
||||
}
|
||||
msg.IsRead = true
|
||||
msg.IsFlagged = true
|
||||
|
||||
hub := connhub.New()
|
||||
srv := NewIMAPServer(config.IMAPConfig{}, stores, nil, config.BanConfig{}, hub)
|
||||
srv.newServer("127.0.0.1:143", nil)
|
||||
hub := srv.hubForOrCreate("alice@example.com", "INBOX")
|
||||
sess := fakeSession()
|
||||
hub.add(sess)
|
||||
|
||||
srv.PushFlagsChanged("alice@example.com", "INBOX", msg)
|
||||
|
||||
srv.beMu.Lock()
|
||||
b := srv.bes[0]
|
||||
srv.beMu.Unlock()
|
||||
|
||||
select {
|
||||
case upd := <-b.updates:
|
||||
mu, ok := upd.(*backend.MessageUpdate)
|
||||
if !ok {
|
||||
t.Fatalf("update type = %T, want *MessageUpdate", upd)
|
||||
}
|
||||
if mu.Username() != "alice@example.com" || mu.Mailbox() != "INBOX" {
|
||||
t.Fatalf("update targeting = %s/%s", mu.Username(), mu.Mailbox())
|
||||
}
|
||||
if mu.Message.Uid != uint32(msg.ID) {
|
||||
t.Fatalf("uid = %d, want %d", mu.Message.Uid, msg.ID)
|
||||
}
|
||||
got := make(map[string]bool)
|
||||
for _, f := range mu.Message.Flags {
|
||||
got[f] = true
|
||||
}
|
||||
if !got["\\Seen"] || !got["\\Flagged"] {
|
||||
t.Fatalf("flags = %v, want \\Seen and \\Flagged", mu.Message.Flags)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("no flags update received")
|
||||
updates := sess.takeUpdates(true)
|
||||
if len(updates) != 1 || updates[0].fetch == nil {
|
||||
t.Fatalf("updates = %+v, want 1 条 FETCH", updates)
|
||||
}
|
||||
f := updates[0].fetch
|
||||
if f.uid != imap.UID(msg.ID) {
|
||||
t.Fatalf("uid = %d, want %d", f.uid, msg.ID)
|
||||
}
|
||||
got := make(map[imap.Flag]bool)
|
||||
for _, fl := range f.flags {
|
||||
got[fl] = true
|
||||
}
|
||||
if !got[imap.FlagSeen] || !got[imap.FlagFlagged] {
|
||||
t.Fatalf("flags = %v, want \\Seen and \\Flagged", f.flags)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPushExpunged 验证删除推送:每条序号一个 ExpungeUpdate。
|
||||
// TestPushExpunged 验证删除推送:每条序号一个 EXPUNGE 更新。
|
||||
func TestPushExpunged(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)
|
||||
srv, _ := newTestServer(t)
|
||||
|
||||
hub := srv.hubForOrCreate("alice@example.com", "INBOX")
|
||||
sess := fakeSession()
|
||||
hub.add(sess)
|
||||
|
||||
hub := connhub.New()
|
||||
srv := NewIMAPServer(config.IMAPConfig{}, stores, nil, config.BanConfig{}, hub)
|
||||
srv.newServer("127.0.0.1:143", nil)
|
||||
srv.PushExpunged("alice@example.com", "INBOX", []uint32{2, 5})
|
||||
|
||||
srv.beMu.Lock()
|
||||
b := srv.bes[0]
|
||||
srv.beMu.Unlock()
|
||||
|
||||
var seqs []uint32
|
||||
for i := 0; i < 2; i++ {
|
||||
select {
|
||||
case upd := <-b.updates:
|
||||
eu, ok := upd.(*backend.ExpungeUpdate)
|
||||
if !ok {
|
||||
t.Fatalf("update type = %T, want *ExpungeUpdate", upd)
|
||||
}
|
||||
if eu.Username() != "alice@example.com" || eu.Mailbox() != "INBOX" {
|
||||
t.Fatalf("update targeting = %s/%s", eu.Username(), eu.Mailbox())
|
||||
}
|
||||
seqs = append(seqs, eu.SeqNum)
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("no expunge update received")
|
||||
}
|
||||
}
|
||||
if seqs[0] != 2 || seqs[1] != 5 {
|
||||
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)
|
||||
}
|
||||
}
|
||||
updates := sess.takeUpdates(true)
|
||||
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")
|
||||
if updates[0].expunge == nil || updates[1].expunge == nil {
|
||||
t.Fatalf("updates = %+v, want expunge updates", updates)
|
||||
}
|
||||
|
||||
// 模拟两个 listenUpdates 各自执行 close(update.Done()):修复前必 panic
|
||||
for _, upd := range updates {
|
||||
close(upd.Done())
|
||||
if *updates[0].expunge != 2 || *updates[1].expunge != 5 {
|
||||
t.Fatalf("seqs = %d,%d, want 2,5", *updates[0].expunge, *updates[1].expunge)
|
||||
}
|
||||
}
|
||||
|
||||
// TestHubExcludesSource 验证来源会话不会收到自己动作的回声推送:
|
||||
// 会话 A 的 STORE/EXPUNGE 只分发给同邮箱的其他会话(本会话的响应已由
|
||||
// 命令本身写回)。回归:v2 库自带 MailboxTracker 对 EXPUNGE/EXISTS 无法
|
||||
// 排除来源,曾导致本会话在后续 Poll 时收到重复 EXPUNGE。
|
||||
func TestHubExcludesSource(t *testing.T) {
|
||||
srv, _ := newTestServer(t)
|
||||
|
||||
hub := srv.hubForOrCreate("alice@example.com", "INBOX")
|
||||
source := fakeSession()
|
||||
other := fakeSession()
|
||||
hub.add(source)
|
||||
hub.add(other)
|
||||
|
||||
hub.enqueue(sessionUpdate{expunge: ptrU32(3)}, source)
|
||||
|
||||
if got := source.takeUpdates(true); len(got) != 0 {
|
||||
t.Fatalf("source 收到 %d 条回声更新,want 0", len(got))
|
||||
}
|
||||
got := other.takeUpdates(true)
|
||||
if len(got) != 1 || got[0].expunge == nil || *got[0].expunge != 3 {
|
||||
t.Fatalf("other updates = %+v, want 1 条 expunge(3)", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPollAllowExpunge 验证 FETCH/STORE/SEARCH 期间不下发 EXPUNGE
|
||||
// (allowExpunge=false 时遇到 EXPUNGE 停止),与 RFC 一致。
|
||||
func TestPollAllowExpunge(t *testing.T) {
|
||||
srv, _ := newTestServer(t)
|
||||
|
||||
hub := srv.hubForOrCreate("alice@example.com", "INBOX")
|
||||
sess := fakeSession()
|
||||
hub.add(sess)
|
||||
|
||||
hub.enqueue(sessionUpdate{fetch: &sessionFetchUpdate{seq: 1, uid: 1, flags: []imap.Flag{imap.FlagSeen}}}, nil)
|
||||
hub.enqueue(sessionUpdate{expunge: ptrU32(2)}, nil)
|
||||
hub.enqueue(sessionUpdate{fetch: &sessionFetchUpdate{seq: 3, uid: 3, flags: []imap.Flag{imap.FlagSeen}}}, nil)
|
||||
|
||||
// allowExpunge=false:只取到第一条 FETCH,EXPUNGE 及其后的留在队列
|
||||
updates := sess.takeUpdates(false)
|
||||
if len(updates) != 1 || updates[0].fetch == nil {
|
||||
t.Fatalf("updates = %+v, want 仅 1 条 FETCH", updates)
|
||||
}
|
||||
// allowExpunge=true:剩余全部下发
|
||||
rest := sess.takeUpdates(true)
|
||||
if len(rest) != 2 {
|
||||
t.Fatalf("rest = %d, want 2", len(rest))
|
||||
}
|
||||
if rest[0].expunge == nil || *rest[0].expunge != 2 {
|
||||
t.Fatalf("rest[0] = %+v, want expunge(2)", rest[0])
|
||||
}
|
||||
}
|
||||
+116
-152
@@ -14,23 +14,22 @@ import (
|
||||
"mail_go/internal/store"
|
||||
"mail_go/internal/tlsutil"
|
||||
|
||||
"github.com/emersion/go-imap"
|
||||
"github.com/emersion/go-imap/backend"
|
||||
imapserver "github.com/emersion/go-imap/server"
|
||||
"github.com/emersion/go-imap/v2"
|
||||
imapserver "github.com/emersion/go-imap/v2/imapserver"
|
||||
)
|
||||
|
||||
// Pusher 是 IMAP 实时推送接口:SMTP/POP3/Web 在邮件状态变化后调用,
|
||||
// 由 go-imap 广播给相关客户端(按用户名+邮箱过滤,IDLE 时即时送达)。
|
||||
// 更新经 mailboxHub 分发给已选中对应邮箱的其他会话(IDLE 时即时送达)。
|
||||
type Pusher interface {
|
||||
// PushNewMessage 推送新邮件(本地投递成功)。
|
||||
PushNewMessage(userEmail string, msg *db.Message)
|
||||
// PushFlagsChanged 推送已读/星标等标志变化(MessageUpdate)。
|
||||
// PushFlagsChanged 推送已读/星标等标志变化。
|
||||
PushFlagsChanged(userEmail, mailbox string, msg *db.Message)
|
||||
// PushExpunged 推送邮件被删除(ExpungeUpdate,seqNums 为删除前序号)。
|
||||
// PushExpunged 推送邮件被删除(seqNums 为删除前序号)。
|
||||
PushExpunged(userEmail, mailbox string, seqNums []uint32)
|
||||
}
|
||||
|
||||
// IMAPServer wraps a go-imap Server and provides mailbox access capability.
|
||||
// IMAPServer 管理 IMAP/IMAPS 监听器、会话注册与跨会话推送。
|
||||
type IMAPServer struct {
|
||||
stores *store.Stores
|
||||
cfg config.IMAPConfig
|
||||
@@ -38,9 +37,11 @@ type IMAPServer struct {
|
||||
tlsLoader *tlsutil.Loader
|
||||
hub *connhub.Hub
|
||||
|
||||
beMu sync.Mutex
|
||||
bes []*imapBackend // 各监听器(明文/TLS)的 backend,用于新邮件推送
|
||||
srvs []*imapserver.Server // 各监听器实例,用于强制断开连接
|
||||
// hubs 按「用户邮箱 + 文件夹」索引的推送中心,会话 SELECT 时加入。
|
||||
mu sync.Mutex
|
||||
hubs map[string]*mailboxHub
|
||||
// sessions 全部活跃会话(DisconnectByAddr 用)。
|
||||
sessions map[*imapSession]struct{}
|
||||
}
|
||||
|
||||
// NewIMAPServer creates a new IMAP server instance. tlsLoader may be nil
|
||||
@@ -52,158 +53,107 @@ func NewIMAPServer(cfg config.IMAPConfig, stores *store.Stores, tlsLoader *tlsut
|
||||
banCfg: banCfg,
|
||||
tlsLoader: tlsLoader,
|
||||
hub: hub,
|
||||
hubs: make(map[string]*mailboxHub),
|
||||
sessions: make(map[*imapSession]struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
// NotifyNewMessage 向所有 IMAP 监听器推送新邮件通知(go-imap 广播时按
|
||||
// 用户名+邮箱过滤,只送达已选中 INBOX 的客户端,IDLE 挂起时实时收到
|
||||
// FETCH 响应)。由 SMTP/Web 本地投递成功时调用;channel 满时非阻塞丢弃。
|
||||
// hubKey 生成推送中心索引键。
|
||||
func hubKey(userEmail, mailbox string) string {
|
||||
return userEmail + "\x00" + mailbox
|
||||
}
|
||||
|
||||
// hubFor 返回已存在的推送中心,不存在返回 nil。
|
||||
func (s *IMAPServer) hubFor(userEmail, mailbox string) *mailboxHub {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.hubs[hubKey(userEmail, mailbox)]
|
||||
}
|
||||
|
||||
// hubForOrCreate 返回推送中心,不存在则创建。
|
||||
func (s *IMAPServer) hubForOrCreate(userEmail, mailbox string) *mailboxHub {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
key := hubKey(userEmail, mailbox)
|
||||
h := s.hubs[key]
|
||||
if h == nil {
|
||||
h = newMailboxHub()
|
||||
s.hubs[key] = h
|
||||
}
|
||||
return h
|
||||
}
|
||||
|
||||
func (s *IMAPServer) unregisterSession(sess *imapSession) {
|
||||
s.mu.Lock()
|
||||
delete(s.sessions, sess)
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// NotifyNewMessage 推送新邮件到达通知:EXISTS 计数 + 标记更新。
|
||||
// 由 SMTP/Web 本地投递成功时调用;无会话选中该邮箱时为 no-op。
|
||||
func (s *IMAPServer) PushNewMessage(userEmail string, msg *db.Message) {
|
||||
if s == nil || userEmail == "" || msg == nil {
|
||||
return
|
||||
}
|
||||
update := buildNewMessageUpdate(s.stores, userEmail, "INBOX", msg)
|
||||
if update == nil {
|
||||
hub := s.hubFor(userEmail, "INBOX")
|
||||
if hub == nil {
|
||||
return
|
||||
}
|
||||
// 先发 EXISTS 通知再发 FETCH 更新:RFC 2177(IDLE)要求新邮件到达
|
||||
// 时服务器发送 EXISTS,不少客户端(如 Apple Mail)只认 EXISTS 才会
|
||||
// 唤醒并主动拉取,仅裸 FETCH 更新会被忽略(表现为"必须手动同步")。
|
||||
if count, err := s.stores.Mails.CountByUserAndFolder(msg.UserID, "INBOX"); err == nil {
|
||||
s.pushExists(userEmail, "INBOX", uint32(count))
|
||||
hub.enqueue(sessionUpdate{exists: ptrU32(uint32(count))}, nil)
|
||||
}
|
||||
s.broadcastUpdate(update, userEmail, msg.ID)
|
||||
}
|
||||
|
||||
// PushFlagsChanged 推送邮件标志(已读/星标等)变化给同用户的其他客户端。
|
||||
// PushFlagsChanged 推送邮件标志(已读/星标/删除标记)变化给同用户其他会话。
|
||||
func (s *IMAPServer) PushFlagsChanged(userEmail, mailbox string, msg *db.Message) {
|
||||
if s == nil || userEmail == "" || mailbox == "" || msg == nil {
|
||||
return
|
||||
}
|
||||
update := buildFlagsUpdate(s.stores, userEmail, mailbox, msg, false)
|
||||
if update == nil {
|
||||
hub := s.hubFor(userEmail, mailbox)
|
||||
if hub == nil {
|
||||
return
|
||||
}
|
||||
s.broadcastUpdate(update, userEmail, msg.ID)
|
||||
seq := seqOf(s.stores, msg.UserID, mailbox, msg.ID)
|
||||
if seq == 0 {
|
||||
return
|
||||
}
|
||||
hub.enqueue(sessionUpdate{fetch: &sessionFetchUpdate{
|
||||
seq: seq,
|
||||
uid: imap.UID(msg.ID),
|
||||
flags: flagsOf(msg.IsRead, msg.IsFlagged, msg.IsDeleted),
|
||||
}}, nil)
|
||||
}
|
||||
|
||||
// PushExpunged 推送邮件被删除(每条序号一个 ExpungeUpdate)。
|
||||
// PushExpunged 推送邮件被删除(每条序号一个 EXPUNGE 更新)。
|
||||
func (s *IMAPServer) PushExpunged(userEmail, mailbox string, seqNums []uint32) {
|
||||
if s == nil || userEmail == "" || mailbox == "" || len(seqNums) == 0 {
|
||||
return
|
||||
}
|
||||
hub := s.hubFor(userEmail, mailbox)
|
||||
if hub == nil {
|
||||
return
|
||||
}
|
||||
for _, seq := range seqNums {
|
||||
update := &backend.ExpungeUpdate{
|
||||
Update: backend.NewUpdate(userEmail, mailbox),
|
||||
SeqNum: seq,
|
||||
}
|
||||
s.broadcastUpdate(update, userEmail, 0)
|
||||
hub.enqueue(sessionUpdate{expunge: ptrU32(seq)}, nil)
|
||||
}
|
||||
}
|
||||
|
||||
// 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...)
|
||||
s.beMu.Unlock()
|
||||
|
||||
for _, b := range bes {
|
||||
select {
|
||||
case b.updates <- cloneUpdate(update):
|
||||
default:
|
||||
log.Printf("IMAP: 推送通道已满,丢弃 %s 的更新 (msg=%d)", userEmail, msgID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// pushExists 向所有监听器中「已登录该用户且已选中该邮箱」的连接直接写入
|
||||
// 未请求的 "* N EXISTS" 响应(绕过 go-imap 更新通道——其仅支持 FETCH/
|
||||
// EXPUNGE 类更新,无法表达 EXISTS)。通道满时非阻塞丢弃,与广播一致。
|
||||
func (s *IMAPServer) pushExists(userEmail, mailbox string, exists uint32) {
|
||||
s.beMu.Lock()
|
||||
srvs := append([]*imapserver.Server(nil), s.srvs...)
|
||||
s.beMu.Unlock()
|
||||
|
||||
for _, srv := range srvs {
|
||||
srv.ForEachConn(func(conn imapserver.Conn) {
|
||||
ctx := conn.Context()
|
||||
if ctx == nil || ctx.User == nil || ctx.Mailbox == nil {
|
||||
return
|
||||
}
|
||||
if ctx.User.Username() != userEmail || ctx.Mailbox.Name() != mailbox {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case ctx.Responses <- existsResponse(exists):
|
||||
default:
|
||||
log.Printf("IMAP: EXISTS 推送通道已满,丢弃 (user=%s mailbox=%s)", userEmail, mailbox)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// existsResponse 序列化为 "* N EXISTS\r\n"。
|
||||
type existsResponse uint32
|
||||
|
||||
func (n existsResponse) WriteTo(w *imap.Writer) error {
|
||||
_, err := fmt.Fprintf(w, "* %d EXISTS\r\n", uint32(n))
|
||||
return err
|
||||
}
|
||||
|
||||
// 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()
|
||||
s.bes = append(s.bes, be)
|
||||
s.beMu.Unlock()
|
||||
}
|
||||
|
||||
// registerServer 记录监听器实例(用于强制断开连接)。
|
||||
func (s *IMAPServer) registerServer(srv *imapserver.Server) {
|
||||
s.beMu.Lock()
|
||||
s.srvs = append(s.srvs, srv)
|
||||
s.beMu.Unlock()
|
||||
}
|
||||
|
||||
// DisconnectByAddr 强制断开指定远端地址的连接(管理后台「断开并封禁」)。
|
||||
// 关闭连接会触发 go-imap 的收尾流程(user.Logout、协议日志回填、hub 注销)。
|
||||
// 发送 BYE 并关闭底层连接,触发库的收尾流程(session.Close、协议日志回填)。
|
||||
func (s *IMAPServer) DisconnectByAddr(remoteAddr string) {
|
||||
if s == nil || remoteAddr == "" {
|
||||
return
|
||||
}
|
||||
s.beMu.Lock()
|
||||
srvs := append([]*imapserver.Server(nil), s.srvs...)
|
||||
s.beMu.Unlock()
|
||||
|
||||
for _, srv := range srvs {
|
||||
srv.ForEachConn(func(conn imapserver.Conn) {
|
||||
info := conn.Info()
|
||||
if info != nil && info.RemoteAddr != nil && info.RemoteAddr.String() == remoteAddr {
|
||||
_ = conn.Close()
|
||||
}
|
||||
})
|
||||
s.mu.Lock()
|
||||
var targets []*imapSession
|
||||
for sess := range s.sessions {
|
||||
if sess.remoteAddr == remoteAddr {
|
||||
targets = append(targets, sess)
|
||||
}
|
||||
}
|
||||
s.mu.Unlock()
|
||||
for _, sess := range targets {
|
||||
_ = sess.conn.Bye("Connection closed by administrator")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -215,25 +165,42 @@ func (s *IMAPServer) tlsConfig() (*tls.Config, error) {
|
||||
return &tls.Config{GetCertificate: s.tlsLoader.GetCertificate}, nil
|
||||
}
|
||||
|
||||
// newServer creates a configured imapserver.Server with the given address.
|
||||
// imapCaps 是服务器支持的能力集:UIDPLUS(RFC 4315)提供 UID EXPUNGE、
|
||||
// COPYUID/APPENDUID 支持;MOVE 由 SessionMove 实现;IDLE/UNSELECT 在
|
||||
// IMAP4rev1 认证后由库自动广告。
|
||||
var imapCaps = imap.CapSet{
|
||||
imap.CapIMAP4rev1: {},
|
||||
imap.CapUIDPlus: {},
|
||||
imap.CapMove: {},
|
||||
imap.CapLiteralPlus: {},
|
||||
imap.CapChildren: {},
|
||||
imap.CapSpecialUse: {},
|
||||
}
|
||||
|
||||
// newServer 创建指定监听地址的 IMAP 服务(会话工厂捕获监听端口)。
|
||||
func (s *IMAPServer) newServer(addr string, tlsConfig *tls.Config) *imapserver.Server {
|
||||
be := &imapBackend{
|
||||
stores: s.stores,
|
||||
banCfg: s.banCfg,
|
||||
port: portOf(addr),
|
||||
hub: s.hub,
|
||||
updates: make(chan backend.Update, 256),
|
||||
disconnectAddr: s.DisconnectByAddr,
|
||||
}
|
||||
s.registerBackend(be)
|
||||
srv := imapserver.New(be)
|
||||
srv.Addr = addr
|
||||
srv.TLSConfig = tlsConfig
|
||||
srv.AllowInsecureAuth = tlsConfig == nil
|
||||
s.registerServer(srv)
|
||||
port := portOf(addr)
|
||||
srv := imapserver.New(&imapserver.Options{
|
||||
Caps: imapCaps,
|
||||
NewSession: s.newSessionFactory(port),
|
||||
TLSConfig: tlsConfig,
|
||||
InsecureAuth: tlsConfig == nil,
|
||||
})
|
||||
return srv
|
||||
}
|
||||
|
||||
// newSessionFactory 构造会话工厂:每个连接一个 imapSession,注册到
|
||||
// 会话表(DisconnectByAddr 用)。
|
||||
func (s *IMAPServer) newSessionFactory(port int) func(conn *imapserver.Conn) (imapserver.Session, *imapserver.GreetingData, error) {
|
||||
return func(conn *imapserver.Conn) (imapserver.Session, *imapserver.GreetingData, error) {
|
||||
sess := newImapSession(s, conn, port)
|
||||
s.mu.Lock()
|
||||
s.sessions[sess] = struct{}{}
|
||||
s.mu.Unlock()
|
||||
return sess, nil, nil
|
||||
}
|
||||
}
|
||||
|
||||
// portOf 从监听地址解析端口号,失败返回 0。
|
||||
func portOf(addr string) int {
|
||||
_, portStr, err := net.SplitHostPort(addr)
|
||||
@@ -247,29 +214,26 @@ func portOf(addr string) int {
|
||||
return port
|
||||
}
|
||||
|
||||
// Start starts the IMAP server on the plain-text port.
|
||||
// Start starts the IMAP server on the plain-text port (STARTTLS enabled
|
||||
// when a certificate is configured).
|
||||
func (s *IMAPServer) Start() error {
|
||||
tlsConfig, err := s.tlsConfig()
|
||||
if err != nil {
|
||||
log.Printf("IMAP STARTTLS 未启用: %v", err)
|
||||
tlsConfig = nil
|
||||
}
|
||||
srv := s.newServer(s.cfg.Addr, tlsConfig)
|
||||
log.Printf("IMAP server listening on %s", s.cfg.Addr)
|
||||
return srv.ListenAndServe()
|
||||
return srv.ListenAndServe(s.cfg.Addr)
|
||||
}
|
||||
|
||||
// StartTLS starts the IMAP server on the TLS port.
|
||||
// StartTLS starts the IMAP server on the implicit TLS port.
|
||||
func (s *IMAPServer) StartTLS() error {
|
||||
tlsConfig, err := s.tlsConfig()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
srv := s.newServer(s.cfg.TLSAddr, tlsConfig)
|
||||
|
||||
log.Printf("IMAPS server listening on %s", s.cfg.TLSAddr)
|
||||
return srv.ListenAndServeTLS()
|
||||
return srv.ListenAndServeTLS(s.cfg.TLSAddr)
|
||||
}
|
||||
|
||||
// ensure imapBackend satisfies backend.Backend at compile time
|
||||
var _ backend.Backend = (*imapBackend)(nil)
|
||||
File diff suppressed because it is too large
Load Diff
@@ -33,6 +33,13 @@ type MailStore interface {
|
||||
// SetFlaggedStates 批量设置多封邮件的星标状态(单条 UPDATE ... IN)。
|
||||
SetFlaggedStates(ids []uint, flagged bool) error
|
||||
MoveToFolder(id uint, folder string) error
|
||||
// SetDeletedStates 批量设置多封邮件的 \Deleted 标记(单条 UPDATE ... IN)。
|
||||
SetDeletedStates(ids []uint, deleted bool) error
|
||||
// ListDeletedByUserAndFolder 列出某文件夹中所有已标记 \Deleted 的邮件
|
||||
// (按 date DESC, id DESC 排序,与全量列表一致,序号映射全链路相同)。
|
||||
ListDeletedByUserAndFolder(userID uint, folder string) ([]db.Message, error)
|
||||
// DeleteMany 批量硬删除多封邮件(单条 DELETE ... IN)。
|
||||
DeleteMany(ids []uint) error
|
||||
Delete(id uint) error
|
||||
CountUnread(userID uint, folder string) (int64, error)
|
||||
CountByFolder(folder string) (int64, error)
|
||||
@@ -130,6 +137,32 @@ func (s *mailStoreGorm) MoveToFolder(id uint, folder string) error {
|
||||
return s.db.Model(&db.Message{}).Where("id = ?", id).Update("folder", folder).Error
|
||||
}
|
||||
|
||||
// SetDeletedStates 批量设置多封邮件的 \Deleted 标记。
|
||||
func (s *mailStoreGorm) SetDeletedStates(ids []uint, deleted bool) error {
|
||||
if len(ids) == 0 {
|
||||
return nil
|
||||
}
|
||||
return s.db.Model(&db.Message{}).Where("id IN ?", ids).Update("is_deleted", deleted).Error
|
||||
}
|
||||
|
||||
// ListDeletedByUserAndFolder 列出某文件夹中所有已标记 \Deleted 的邮件。
|
||||
func (s *mailStoreGorm) ListDeletedByUserAndFolder(userID uint, folder string) ([]db.Message, error) {
|
||||
var messages []db.Message
|
||||
if err := s.db.Where("user_id = ? AND folder = ? AND is_deleted = ?", userID, folder, true).
|
||||
Order("date DESC, id DESC").Find(&messages).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return messages, nil
|
||||
}
|
||||
|
||||
// DeleteMany 批量硬删除多封邮件。
|
||||
func (s *mailStoreGorm) DeleteMany(ids []uint) error {
|
||||
if len(ids) == 0 {
|
||||
return nil
|
||||
}
|
||||
return s.db.Where("id IN ?", ids).Delete(&db.Message{}).Error
|
||||
}
|
||||
|
||||
// Delete removes a message by ID.
|
||||
func (s *mailStoreGorm) Delete(id uint) error {
|
||||
return s.db.Delete(&db.Message{}, id).Error
|
||||
|
||||
Reference in New Issue
Block a user