diff --git a/README.md b/README.md index 6214b95..392de59 100644 --- a/README.md +++ b/README.md @@ -403,8 +403,8 @@ mailgo/ │ │ ├── manager.go # 外发队列、重试、限速、退信 │ │ └── sign.go # DKIM 签名 │ ├── imap_server/ -│ │ ├── server.go # IMAP 服务 -│ │ └── backend.go # IMAP 后端 +│ │ ├── server.go # IMAP 服务(监听器/能力/跨会话推送) +│ │ └── session.go # IMAP Session 实现(SELECT/FETCH/STORE/EXPUNGE) │ ├── pop3_server/server.go # POP3 服务 │ ├── connhub/hub.go # 协议连接注册中心(当前连接监控) │ ├── storage/attachment.go # 附件文件存储 @@ -521,7 +521,7 @@ sudo journalctl -u mailgo -f | 数据库 | SQLite(默认)/ MySQL | | 配置格式 | TOML | | SMTP | github.com/emersion/go-smtp | -| IMAP | github.com/emersion/go-imap v1 | +| IMAP | github.com/emersion/go-imap/v2 | | POP3 | 手工实现 TCP 协议 | | 密码哈希 | golang.org/x/crypto/bcrypt | | 富文本 | Quill.js (CDN) | diff --git a/docs/architecture.md b/docs/architecture.md index 21eb309..05a17bb 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -83,8 +83,8 @@ mail_go/ │ ├── smtp_server/ │ │ └── server.go # SMTP 服务端(go-smtp Backend 实现) │ ├── imap_server/ -│ │ ├── server.go # IMAP 服务端启动 -│ │ └── backend.go # go-imap Backend/User/Mailbox/Message 实现 +│ │ ├── server.go # IMAP 服务端启动(监听器/能力/跨会话推送) +│ │ └── session.go # go-imap/v2 imapserver.Session 实现 │ ├── pop3_server/ │ │ └── server.go # POP3 服务端(TCP 监听 + 文本协议) │ ├── web/ @@ -373,11 +373,11 @@ func (s *smtpSession) Logout() error #### IMAP Backend 接口(go-imap/v2 要求) ```go -// internal/imap_server/backend.go +// internal/imap_server/session.go -// 实现 go-imap/v2 的 backend.Backend 接口 -type imapBackend struct { - userStore store.UserStore +// 实现 go-imap/v2 的 imapserver.Session 接口 +type imapSession struct { + stores *store.Stores mailStore store.MailStore domainStore store.DomainStore attStore store.AttachmentStore @@ -557,7 +557,7 @@ sequenceDiagram | Task ID | 任务名称 | 依赖 | 涉及文件 | 优先级 | |---------|---------|------|---------|--------| | T01 | 项目基础设施:go.mod + 入口 + 配置系统 + 数据库层 | — | go.mod, main.go, config/config.go, config/defaults.go, internal/db/db.go, internal/db/models.go, internal/store/*.go | P0 | -| T02 | 邮件协议服务端(SMTP + IMAP + POP3) | T01 | internal/smtp_server/server.go, internal/imap_server/server.go, internal/imap_server/backend.go, internal/pop3_server/server.go | P0 | +| T02 | 邮件协议服务端(SMTP + IMAP + POP3) | T01 | internal/smtp_server/server.go, internal/imap_server/server.go, internal/imap_server/session.go, internal/pop3_server/server.go | P0 | | T03 | Web 服务核心:路由 + 中间件 + 认证 + 邮件页面 | T01 | internal/web/server.go, internal/web/middleware/auth.go, internal/web/middleware/admin.go, internal/web/handlers/auth.go, internal/web/handlers/mail.go | P0 | | T04 | 管理后台 + 附件存储 + 模板 | T01, T03 | internal/web/handlers/admin.go, internal/storage/attachment.go, internal/web/templates/*.html | P0 | | T05 | 集成调试 + 安装脚本 | T02, T03, T04 | main.go(更新), scripts/install.sh | P1 | @@ -614,7 +614,7 @@ go 1.22 require ( github.com/emersion/go-smtp v0.21.0 - github.com/emersion/go-imap/v2 v2.0.0-beta.5 + github.com/emersion/go-imap/v2 v2.0.0-beta.8 github.com/emersion/go-message v0.18.0 github.com/gin-gonic/gin v1.10.0 github.com/gin-contrib/sessions v0.0.5 diff --git a/go.mod b/go.mod index 0b4ee92..11bd737 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.25.0 require ( github.com/BurntSushi/toml v1.4.0 - github.com/emersion/go-imap v1.2.1 + github.com/emersion/go-imap/v2 v2.0.0-beta.8 github.com/emersion/go-message v0.18.2 github.com/emersion/go-msgauth v0.7.0 github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6 diff --git a/go.sum b/go.sum index 9593529..601bee3 100644 --- a/go.sum +++ b/go.sum @@ -17,19 +17,16 @@ github.com/cloudwego/base64x v0.1.6/go.mod h1:OFcloc187FXDaYHvrNIjxSe8ncn0OOM8gE github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/emersion/go-imap v1.2.1 h1:+s9ZjMEjOB8NzZMVTM3cCenz2JrQIGGo5j1df19WjTA= -github.com/emersion/go-imap v1.2.1/go.mod h1:Qlx1FSx2FTxjnjWpIlVNEuX+ylerZQNFE5NsmKFSejY= -github.com/emersion/go-message v0.15.0/go.mod h1:wQUEfE+38+7EW8p8aZ96ptg6bAb1iwdgej19uXASlE4= +github.com/emersion/go-imap/v2 v2.0.0-beta.8 h1:5IXZK1E33DyeP526320J3RS7eFlCYGFgtbrfapqDPug= +github.com/emersion/go-imap/v2 v2.0.0-beta.8/go.mod h1:dhoFe2Q0PwLrMD7oZw8ODuaD0vLYPe5uj2wcOMnvh48= github.com/emersion/go-message v0.18.2 h1:rl55SQdjd9oJcIoQNhubD2Acs1E6IzlZISRTK7x/Lpg= github.com/emersion/go-message v0.18.2/go.mod h1:XpJyL70LwRvq2a8rVbHXikPgKj8+aI0kGdHlg16ibYA= github.com/emersion/go-msgauth v0.7.0 h1:vj2hMn6KhFtW41kshIBTXvp6KgYSqpA/ZN9Pv4g1INc= github.com/emersion/go-msgauth v0.7.0/go.mod h1:mmS9I6HkSovrNgq0HNXTeu8l3sRAAuQ9RMvbM4KU7Ck= -github.com/emersion/go-sasl v0.0.0-20200509203442-7bfe0ed36a21/go.mod h1:iL2twTeMvZnrg54ZoPDNfJaJaqy0xIQFuBdrLsmspwQ= github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6 h1:oP4q0fw+fOSWn3DfFi4EXdT+B+gTtzx8GC9xsc26Znk= github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6/go.mod h1:iL2twTeMvZnrg54ZoPDNfJaJaqy0xIQFuBdrLsmspwQ= github.com/emersion/go-smtp v0.24.0 h1:g6AfoF140mvW0vLNPD/LuCBLEAdlxOjIXqbIkJIS6Wk= github.com/emersion/go-smtp v0.24.0/go.mod h1:ZtRRkbTyp2XTHCA+BmyTFTrj8xY4I+b4McvHxCU2gsQ= -github.com/emersion/go-textwrapper v0.0.0-20200911093747-65d896831594/go.mod h1:aqO8z8wPrjkscevZJFVE1wXJrLpC5LtJG7fqLOsPb2U= github.com/gabriel-vasile/mimetype v1.4.12 h1:e9hWvmLYvtp846tLHam2o++qitpguFiYCKbn0w9jyqw= github.com/gabriel-vasile/mimetype v1.4.12/go.mod h1:d+9Oxyo1wTzWdyVUPMmXFvp4F9tea18J8ufA774AB3s= github.com/gin-contrib/sessions v1.1.0 h1:00mhHfNEGF5sP2fwxa98aRqj1FOJdL6IkR86n2hOiBo= @@ -163,7 +160,6 @@ golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuX golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= diff --git a/internal/db/models.go b/internal/db/models.go index f554029..9f2c6f9 100644 --- a/internal/db/models.go +++ b/internal/db/models.go @@ -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"` } diff --git a/internal/imap_server/backend.go b/internal/imap_server/backend.go deleted file mode 100644 index 2f2deec..0000000 --- a/internal/imap_server/backend.go +++ /dev/null @@ -1,1152 +0,0 @@ -package imap_server - -import ( - "bufio" - "bytes" - "fmt" - "io" - "log" - "net/mail" - "strings" - "time" - - "mail_go/config" - "mail_go/internal/connhub" - "mail_go/internal/db" - "mail_go/internal/mailutil" - "mail_go/internal/store" - - "github.com/emersion/go-imap" - "github.com/emersion/go-imap/backend" - "github.com/emersion/go-imap/backend/backendutil" - asgomail "github.com/emersion/go-message/mail" - "github.com/emersion/go-message/textproto" -) - -// ---------- imapBackend ---------- - -// imapBackend implements backend.Backend and backend.BackendUpdater. -type imapBackend struct { - stores *store.Stores - banCfg config.BanConfig - port int - hub *connhub.Hub - - // updates 承载新邮件等后端更新,由 go-imap 服务器广播给相关客户端。 - updates chan backend.Update - // disconnectAddr 强制断开指定远端地址的连接(管理后台断开封禁用)。 - disconnectAddr func(addr string) -} - -// Updates 实现 backend.BackendUpdater:新邮件推送通道(广播按用户名与 -// 邮箱过滤,只送达已选中对应邮箱的客户端)。 -func (b *imapBackend) Updates() <-chan backend.Update { - return b.updates -} - -// buildNewMessageUpdate 为一条新投递到 mailbox 的邮件构造 IMAP 更新。 -// seq 取该邮件在邮箱中的实际序号(最新在前,通常为 1)。 -func buildNewMessageUpdate(stores *store.Stores, userEmail, mailbox string, msg *db.Message) *backend.MessageUpdate { - if stores == nil || msg == nil || userEmail == "" || mailbox == "" { - return nil - } - seq := seqOf(stores, msg.UserID, mailbox, msg.ID) - if seq == 0 { - seq = 1 - } - - imapMsg := imap.NewMessage(seq, []imap.FetchItem{imap.FetchUid, imap.FetchFlags, imap.FetchInternalDate, imap.FetchRFC822Size, imap.FetchEnvelope}) - imapMsg.Uid = uint32(msg.ID) - imapMsg.Flags = flagsOf(msg.IsRead, msg.IsFlagged, false) - imapMsg.InternalDate = msg.Date - imapMsg.Size = uint32(len(msg.RawData)) - imapMsg.Envelope = &imap.Envelope{ - Date: msg.Date, - Subject: msg.Subject, - From: parseAddressList(msg.FromAddr), - Sender: parseAddressList(msg.FromAddr), - ReplyTo: parseAddressList(msg.FromAddr), - To: parseAddressList(msg.ToAddr), - Cc: parseAddressList(msg.CcAddr), - MessageId: msg.MessageID, - } - - return &backend.MessageUpdate{ - Update: backend.NewUpdate(userEmail, mailbox), - Message: imapMsg, - } -} - -// buildFlagsUpdate 为一条消息的标志变化构造 IMAP 更新(已读/星标/删除标记)。 -// deleted 为会话内 \Deleted 标记(IMAP STORE 会话状态,不入库)。 -func buildFlagsUpdate(stores *store.Stores, userEmail, mailbox string, msg *db.Message, deleted bool) *backend.MessageUpdate { - if stores == nil || msg == nil || userEmail == "" || mailbox == "" { - return nil - } - return buildFlagsUpdateAt(seqOf(stores, msg.UserID, mailbox, msg.ID), userEmail, mailbox, msg, msg.IsRead, msg.IsFlagged, deleted) -} - -// buildFlagsUpdateAt 与 buildFlagsUpdate 相同,但序号与已读/星标状态由 -// 调用方直接提供(批量 STORE 路径已加载全量列表,避免逐条 GetByID + -// seqOf 全量扫描)。 -func buildFlagsUpdateAt(seq uint32, userEmail, mailbox string, msg *db.Message, read, flagged, deleted bool) *backend.MessageUpdate { - if msg == nil || userEmail == "" || mailbox == "" { - return nil - } - imapMsg := imap.NewMessage(seq, - []imap.FetchItem{imap.FetchUid, imap.FetchFlags}) - imapMsg.Uid = uint32(msg.ID) - imapMsg.Flags = flagsOf(read, flagged, deleted) - return &backend.MessageUpdate{ - Update: backend.NewUpdate(userEmail, mailbox), - Message: imapMsg, - } -} - -// seqOf 返回消息在文件夹中的序号(1 基),未找到返回 0。 -func seqOf(stores *store.Stores, userID uint, mailbox string, msgID uint) uint32 { - msgs, err := stores.Mails.ListAllByUserAndFolder(userID, mailbox) - if err != nil { - return 0 - } - for i := range msgs { - if msgs[i].ID == msgID { - return uint32(i + 1) - } - } - return 0 -} - -// flagsOf 按数据库状态生成 IMAP 标志列表(deleted 为会话内 \Deleted 标记)。 -func flagsOf(read, flagged, deleted bool) []string { - flags := make([]string, 0, 3) - if read { - flags = append(flags, "\\Seen") - } - if flagged { - flags = append(flags, "\\Flagged") - } - if deleted { - flags = append(flags, "\\Deleted") - } - return flags -} - -// pushUpdate 非阻塞地把一条后端更新送入推送通道(满则丢弃并记日志)。 -func pushUpdate(ch chan backend.Update, u backend.Update) { - if ch == nil { - return - } - select { - case ch <- u: - default: - log.Printf("IMAP: 推送通道已满,丢弃更新") - } -} - -// Login authenticates a user by email and password. -func (b *imapBackend) Login(connInfo *imap.ConnInfo, username, password string) (backend.User, error) { - clientIP := store.ClientIPFromAddr(connInfo.RemoteAddr) - now := time.Now() - - // 已封禁 IP 一律拒绝认证(防协议层暴力破解) - if banned, _ := b.stores.Bans.IsBanned(clientIP); banned { - b.recordLogin(clientIP, username, false, "IP已被封禁", "认证被拒绝(IP 已封禁)", 0, now) - return nil, backend.ErrInvalidCredentials - } - - user, err := b.stores.Users.AuthenticateLogin(username, password) - if err != nil { - // 认证失败计数,达到阈值按档位封禁(与 Web 登录共用 ban_entries) - b.stores.RecordAuthFailure(clientIP, b.banCfg.MaxFailAttempts, b.banCfg.BanDurationMin, "邮件协议认证失败次数过多") - b.recordLogin(clientIP, username, false, "用户名或密码错误", "LOGIN 失败", 0, now) - return nil, fmt.Errorf("invalid credentials: %w", err) - } - - // 登录成功清零失败计数(与 Web 登录一致):否则协议客户端的失败计数 - // 只增不减(如配置探测、输错密码、APP 用裸用户名重试等),累计触发 - // 档位封禁,合法用户 IP 被反复误封。 - b.stores.Bans.ResetFail(clientIP) - - email := user.Username + "@" - domain, err := b.stores.Domains.GetByID(user.DomainID) - if err == nil { - email = user.Username + "@" + domain.Name - } - - logID := b.recordLogin(clientIP, username, true, "", "LOGIN 成功", 0, now) - - // 连接追踪:注册到当前连接中心,Logout 时注销; - // 注册强制断开回调(按远端地址匹配,断开时由 go-imap 走正常收尾)。 - conn := b.hub.Register("imap", clientIP, b.port, connInfo.TLS != nil) - if conn != nil { - conn.SetUser(email) - remoteAddr := "" - if connInfo.RemoteAddr != nil { - remoteAddr = connInfo.RemoteAddr.String() - } - if b.disconnectAddr != nil && remoteAddr != "" { - addr := remoteAddr - conn.SetDisconnect(func() { b.disconnectAddr(addr) }) - } - } - - return &imapUser{ - stores: b.stores, - id: user.ID, - email: email, - logID: logID, - clientIP: clientIP, - startedAt: now, - conn: conn, - updates: b.updates, - }, nil -} - -// recordLogin 写入一条 IMAP 登录日志,返回新记录的 ID(失败时为 0)。 -func (b *imapBackend) recordLogin(ip, username string, success bool, failReason, detail string, durationMs int64, at time.Time) uint { - entry := &db.ProtocolLog{ - Protocol: db.ProtocolIMAP, - Port: b.port, - ClientIP: ip, - Username: username, - Success: success, - FailReason: failReason, - Detail: detail, - DurationMs: durationMs, - CreatedAt: at, - } - if err := b.stores.ProtocolLogs.Create(entry); err != nil { - log.Printf("IMAP: 写入协议日志失败: %v", err) - return 0 - } - return entry.ID -} - -// ---------- imapUser ---------- - -// imapUser implements backend.User. -type imapUser struct { - stores *store.Stores - id uint - email string - logID uint - clientIP string - startedAt time.Time - conn *connhub.Conn - // updates 所在 backend 的推送通道(STORE/EXPUNGE 等实时同步用)。 - updates chan backend.Update -} - -// Username returns the user's email address. -func (u *imapUser) Username() string { - return u.email -} - -// ListMailboxes returns the standard mailbox list. -func (u *imapUser) ListMailboxes(subscribed bool) ([]backend.Mailbox, error) { - folders := []struct { - name string - delimiter string - attributes []string - }{ - {"INBOX", "/", nil}, - {"Sent", "/", []string{"\\Sent"}}, - {"Drafts", "/", []string{"\\Drafts"}}, - {"Trash", "/", []string{"\\Trash"}}, - } - - mailboxes := make([]backend.Mailbox, 0, len(folders)) - for _, f := range folders { - mailboxes = append(mailboxes, &imapMailbox{ - stores: u.stores, - user: u, - name: f.name, - delimiter: f.delimiter, - attributes: f.attributes, - }) - } - return mailboxes, nil -} - -// GetMailbox returns a mailbox by name. -func (u *imapUser) GetMailbox(name string) (backend.Mailbox, error) { - normalized, ok := canonicalMailboxName(name) - if !ok { - return nil, backend.ErrNoSuchMailbox - } - - return &imapMailbox{ - stores: u.stores, - user: u, - name: normalized, - delimiter: "/", - }, nil -} - -func canonicalMailboxName(name string) (string, bool) { - switch strings.ToUpper(strings.TrimSpace(name)) { - case "INBOX": - return "INBOX", true - case "SENT": - return "Sent", true - case "DRAFTS": - return "Drafts", true - case "TRASH": - return "Trash", true - default: - return "", false - } -} - -// CreateMailbox creates a new mailbox (not supported in this version). -func (u *imapUser) CreateMailbox(name string) error { - return fmt.Errorf("mailbox creation not supported") -} - -// DeleteMailbox deletes a mailbox (not supported in this version). -func (u *imapUser) DeleteMailbox(name string) error { - return fmt.Errorf("mailbox deletion not supported") -} - -// RenameMailbox renames a mailbox (not supported in this version). -func (u *imapUser) RenameMailbox(existingName, newName string) error { - return fmt.Errorf("mailbox rename not supported") -} - -// Logout is called when the user session ends. -func (u *imapUser) Logout() error { - // 回填会话时长,登录记录在 Login 时已写入 - if u.logID == 0 { - return nil - } - durationMs := time.Since(u.startedAt).Milliseconds() - if err := u.stores.ProtocolLogs.UpdateDuration(u.logID, durationMs); err != nil { - log.Printf("IMAP: 更新协议日志失败: %v", err) - } - u.conn.Close() - return nil -} - -// ---------- imapMailbox ---------- - -// imapMailbox implements backend.Mailbox. -type imapMailbox struct { - stores *store.Stores - user *imapUser - name string - delimiter string - attributes []string - // deleted tracks messages marked as \Deleted in this session - deleted map[uint]bool -} - -// Name returns the mailbox name. -func (m *imapMailbox) Name() string { - return m.name -} - -// Info returns mailbox metadata. -func (m *imapMailbox) Info() (*imap.MailboxInfo, error) { - attrs := m.attributes - if attrs == nil { - attrs = []string{} - } - return &imap.MailboxInfo{ - Name: m.name, - Delimiter: m.delimiter, - Attributes: attrs, - }, nil -} - -// Status returns mailbox status information. -func (m *imapMailbox) Status(items []imap.StatusItem) (*imap.MailboxStatus, error) { - status := imap.NewMailboxStatus(m.name, items) - status.Flags = []string{"\\Answered", "\\Flagged", "\\Deleted", "\\Seen", "\\Draft"} - status.PermanentFlags = []string{"\\Answered", "\\Flagged", "\\Deleted", "\\Seen", "\\Draft", "\\*"} - - messages, err := m.stores.Mails.ListAllByUserAndFolder(m.user.id, m.name) - if err != nil { - return nil, err - } - status.Messages = uint32(len(messages)) - - var unseenCount uint32 - for _, msg := range messages { - if !msg.IsRead { - unseenCount++ - } - } - status.Unseen = unseenCount - status.Recent = 0 - maxID, err := m.stores.Mails.MaxIDByUserAndFolder(m.user.id, m.name) - if err != nil { - return nil, err - } - status.UidNext = uint32(maxID + 1) - // UIDVALIDITY 持久化随机值(RFC 3501):数据库重建导致消息 ID 空间 - // 变化时该值随之改变,客户端才会丢弃旧缓存全量重同步。此前硬编码 1, - // 数据库重建后 Thunderbird 等客户端缓存永不失效(只下载"新增"的 - // UID),表现为列表只剩少量邮件。 - uidValidity, err := m.stores.MailboxState.UidValidity(m.user.id, m.name) - if err != nil { - log.Printf("IMAP: 获取 UIDVALIDITY 失败 folder=%s: %v", m.name, err) - uidValidity = 1 - } - status.UidValidity = uidValidity - - return status, nil -} - -// SetSubscribed sets the subscribed status (no-op for now). -func (m *imapMailbox) SetSubscribed(subscribed bool) error { - return nil -} - -// Check is a no-op checkpoint. -func (m *imapMailbox) Check() error { - return nil -} - -// ListMessages returns messages matching the sequence set and fetch items. -func (m *imapMailbox) ListMessages(uid bool, seqset *imap.SeqSet, items []imap.FetchItem, ch chan<- *imap.Message) error { - defer close(ch) - - // Fetch all messages in this mailbox - dbMessages, err := m.stores.Mails.ListAllByUserAndFolder(m.user.id, m.name) - if err != nil { - return err - } - if len(dbMessages) == 0 { - return nil - } - - // Build a mapping of sequence number (1-based) to db.Message - type seqEntry struct { - seqNum uint32 - msg *db.Message - } - entries := make([]seqEntry, len(dbMessages)) - for i := range dbMessages { - entries[i] = seqEntry{ - seqNum: uint32(i + 1), - msg: &dbMessages[i], - } - } - - for _, entry := range entries { - var match bool - if uid { - match = seqset.Contains(uint32(entry.msg.ID)) - } else { - match = seqset.Contains(entry.seqNum) - } - if !match { - continue - } - - imapMsg, err := m.buildIMAPMessage(entry.msg, entry.seqNum, items) - if err != nil { - log.Printf("IMAP: error building message %d: %v", entry.msg.ID, err) - continue - } - ch <- imapMsg - } - - return nil -} - -// buildIMAPMessage constructs an imap.Message from a db.Message with the requested items. -func (m *imapMailbox) buildIMAPMessage(dbMsg *db.Message, seqNum uint32, items []imap.FetchItem) (*imap.Message, error) { - imapMsg := imap.NewMessage(seqNum, items) - imapMsg.Uid = uint32(dbMsg.ID) - imapMsg.Flags = m.getMessageFlags(dbMsg) - imapMsg.InternalDate = dbMsg.Date - - rawMsg := messageRawData(dbMsg) - imapMsg.Size = uint32(len(rawMsg)) - - for _, item := range items { - switch item { - case imap.FetchUid, imap.FetchFlags, imap.FetchInternalDate, imap.FetchRFC822Size: - continue - case imap.FetchEnvelope: - hdr, _, err := headerAndBody(rawMsg) - if err == nil { - imapMsg.Envelope, _ = backendutil.FetchEnvelope(hdr) - } - if imapMsg.Envelope == nil { - imapMsg.Envelope = m.buildEnvelope(dbMsg) - } - case imap.FetchBody, imap.FetchBodyStructure: - hdr, body, err := headerAndBody(rawMsg) - if err == nil { - imapMsg.BodyStructure, _ = backendutil.FetchBodyStructure(hdr, body, item == imap.FetchBodyStructure) - } - // 防御:FetchBodyStructure 对部分合法/畸形 MIME 会失败并返回 - // nil(典型:message/rfc822 附件为 base64 编码时库内不解码 - // 直接按嵌套消息解析头;或 multipart 边界截断)。BodyStructure - // 为 nil 时 go-imap 格式化 FETCH 响应会在 send() 协程 panic - // (nil 指针解引用),连接中断导致客户端只收到部分邮件甚至 - // 一直卡在同步。解析失败时降级为 text/plain 单段结构。 - if imapMsg.BodyStructure == nil { - imapMsg.BodyStructure = fallbackBodyStructure(rawMsg) - } - default: - section, err := imap.ParseBodySectionName(item) - if err != nil { - continue - } - hdr, body, err := headerAndBody(rawMsg) - if err != nil { - return nil, err - } - literal, _ := backendutil.FetchBodySection(hdr, body, section) - if literal != nil { - imapMsg.Body[section] = literal - } - } - } - - return imapMsg, nil -} - -// fallbackBodyStructure 构造一个 text/plain 单段 BodyStructure,用于 -// MIME 解析失败的消息(保证 FETCH BODY/BODYSTRUCTURE 不因 nil 崩溃)。 -func fallbackBodyStructure(raw []byte) *imap.BodyStructure { - size := uint32(len(raw)) - lines := uint32(bytes.Count(raw, []byte{'\n'})) - if len(raw) > 0 && raw[len(raw)-1] != '\n' { - lines++ - } - return &imap.BodyStructure{ - MIMEType: "text", - MIMESubType: "plain", - Params: map[string]string{"charset": "utf-8"}, - Encoding: "8bit", - Size: size, - Lines: lines, - } -} - -func messageRawData(msg *db.Message) []byte { - if msg.RawData != "" { - return []byte(msg.RawData) - } - return buildRawMessage(msg) -} - -func headerAndBody(raw []byte) (textproto.Header, io.Reader, error) { - body := bufio.NewReader(bytes.NewReader(raw)) - hdr, err := textproto.ReadHeader(body) - return hdr, body, err -} - -// getMessageFlags returns IMAP flags for a database message. -func (m *imapMailbox) getMessageFlags(dbMsg *db.Message) []string { - flags := make([]string, 0) - if dbMsg.IsRead { - flags = append(flags, "\\Seen") - } - if dbMsg.IsFlagged { - flags = append(flags, "\\Flagged") - } - if m.deleted != nil && m.deleted[dbMsg.ID] { - flags = append(flags, "\\Deleted") - } - return flags -} - -// buildEnvelope constructs an imap.Envelope from a db.Message. -func (m *imapMailbox) buildEnvelope(dbMsg *db.Message) *imap.Envelope { - env := &imap.Envelope{ - Date: dbMsg.Date, - Subject: dbMsg.Subject, - From: parseAddressList(dbMsg.FromAddr), - Sender: parseAddressList(dbMsg.FromAddr), - ReplyTo: parseAddressList(dbMsg.FromAddr), - To: parseAddressList(dbMsg.ToAddr), - Cc: parseAddressList(dbMsg.CcAddr), - MessageId: dbMsg.MessageID, - } - return env -} - -// SearchMessages returns sequence numbers or UIDs of messages matching the criteria. -func (m *imapMailbox) SearchMessages(uid bool, criteria *imap.SearchCriteria) ([]uint32, error) { - dbMessages, err := m.stores.Mails.ListAllByUserAndFolder(m.user.id, m.name) - if err != nil { - return nil, err - } - - var results []uint32 - for i, dbMsg := range dbMessages { - if m.matchesCriteria(&dbMsg, criteria) { - if uid { - results = append(results, uint32(dbMsg.ID)) - } else { - results = append(results, uint32(i+1)) - } - } - } - return results, nil -} - -// matchesCriteria checks if a message matches the given search criteria. -func (m *imapMailbox) matchesCriteria(msg *db.Message, criteria *imap.SearchCriteria) bool { - // Check WithFlags (messages must have all these flags) - for _, flag := range criteria.WithFlags { - switch flag { - case "\\Seen": - if !msg.IsRead { - return false - } - case "\\Flagged": - if !msg.IsFlagged { - return false - } - case "\\Deleted": - if m.deleted == nil || !m.deleted[msg.ID] { - return false - } - } - } - - // Check WithoutFlags (messages must NOT have any of these flags) - for _, flag := range criteria.WithoutFlags { - switch flag { - case "\\Seen": - if msg.IsRead { - return false - } - case "\\Flagged": - if msg.IsFlagged { - return false - } - case "\\Deleted": - if m.deleted != nil && m.deleted[msg.ID] { - return false - } - } - } - - // Check date range - if !criteria.Since.IsZero() && msg.Date.Before(criteria.Since) { - return false - } - if !criteria.Before.IsZero() && !msg.Date.Before(criteria.Before) { - return false - } - - // Check header fields - if criteria.Header != nil { - if subject := criteria.Header.Get("Subject"); subject != "" { - if !strings.Contains(strings.ToLower(msg.Subject), strings.ToLower(subject)) { - return false - } - } - if from := criteria.Header.Get("From"); from != "" { - if !strings.Contains(strings.ToLower(msg.FromAddr), strings.ToLower(from)) { - return false - } - } - if to := criteria.Header.Get("To"); to != "" { - if !strings.Contains(strings.ToLower(msg.ToAddr), strings.ToLower(to)) { - return false - } - } - } - - // Check body text - for _, text := range criteria.Body { - bodyText := strings.ToLower(msg.TextBody + " " + msg.HtmlBody) - if !strings.Contains(bodyText, strings.ToLower(text)) { - return false - } - } - - // Check generic text (searches headers + body) - for _, text := range criteria.Text { - allText := strings.ToLower(msg.Subject + " " + msg.FromAddr + " " + msg.ToAddr + " " + msg.TextBody + " " + msg.HtmlBody) - if !strings.Contains(allText, strings.ToLower(text)) { - return false - } - } - - // Check NOT criteria - for _, notCrit := range criteria.Not { - if m.matchesCriteria(msg, notCrit) { - return false - } - } - - // Check OR criteria (at least one must match) - for _, orPair := range criteria.Or { - if !m.matchesCriteria(msg, orPair[0]) && !m.matchesCriteria(msg, orPair[1]) { - return false - } - } - - return true -} - -// CreateMessage appends a new message to the mailbox (IMAP APPEND command). -func (m *imapMailbox) CreateMessage(flags []string, date time.Time, body imap.Literal) error { - // Read the message literal - data, err := io.ReadAll(body) - if err != nil { - return fmt.Errorf("failed to read message body: %w", err) - } - - // Parse as MIME message - mr, err := asgomail.CreateReader(bytes.NewReader(data)) - if err != nil { - return fmt.Errorf("failed to parse MIME message: %w", err) - } - - header := mr.Header - fromAddr := mailutil.FormatAddressList(&header, "From") - toAddr := mailutil.FormatAddressList(&header, "To") - ccAddr := mailutil.FormatAddressList(&header, "Cc") - subject, _ := header.Subject() - messageID, _ := header.MessageID() - msgDate, _ := header.Date() - if msgDate.IsZero() { - msgDate = date - } - if msgDate.IsZero() { - msgDate = time.Now() - } - - var textBody, htmlBody string - for { - p, err := mr.NextPart() - if err == io.EOF { - break - } - if err != nil { - break - } - - switch h := p.Header.(type) { - case *asgomail.InlineHeader: - contentType, params, _ := h.ContentType() - buf, _ := io.ReadAll(p.Body) - // 检测并转换字符集 - charset := "" - if cs, ok := params["charset"]; ok { - charset = cs - } - decoded := mailutil.DecodeCharset(buf, charset) - if strings.HasPrefix(contentType, "text/plain") { - textBody = decoded - } else if strings.HasPrefix(contentType, "text/html") { - htmlBody = decoded - } - case *asgomail.AttachmentHeader: - // Attachments from APPEND are not saved in this simple implementation - } - } - - if textBody == "" && htmlBody == "" { - textBody = string(data) - } - - // Determine initial flag state - isRead := false - isFlagged := false - for _, flag := range flags { - switch flag { - case "\\Seen": - isRead = true - case "\\Flagged": - isFlagged = true - } - } - - msg := &db.Message{ - UserID: m.user.id, - MessageID: messageID, - Folder: m.name, - FromAddr: fromAddr, - ToAddr: toAddr, - CcAddr: ccAddr, - Subject: subject, - TextBody: textBody, - HtmlBody: htmlBody, - RawData: string(data), - IsRead: isRead, - IsFlagged: isFlagged, - Date: msgDate, - } - - if err := m.stores.Mails.Create(msg); err != nil { - return fmt.Errorf("failed to create message: %w", err) - } - - // 新邮件(IMAP APPEND)→ 推送给同用户其他已选中该邮箱的客户端 - pushUpdate(m.user.updates, buildNewMessageUpdate(m.stores, m.user.email, m.name, msg)) - - return nil -} - -// UpdateMessagesFlags modifies flags on messages. -func (m *imapMailbox) UpdateMessagesFlags(uid bool, seqset *imap.SeqSet, op imap.FlagsOp, flags []string) error { - if m.deleted == nil { - m.deleted = make(map[uint]bool) - } - - dbMessages, err := m.stores.Mails.ListAllByUserAndFolder(m.user.id, m.name) - if err != nil { - return err - } - // 批量模式:先收集每封匹配消息的目标标志状态,再合并成少量 SQL - // 写入。此前逐条 UPDATE + 逐条 GetByID + 逐条 seqOf 全量扫描, - // 手机整批标记已读(60+ 封)时会产生 60 次写 + 180 次全量读 - // (含大附件 raw_data),该连接命令循环被长时间占住,其他连接 - // 的响应通道也被推送洪泛阻塞——表现为"卡在接收邮件"。 - flagSet := make(map[string]bool, len(flags)) - for _, flag := range flags { - flagSet[flag] = true - } - - type change struct { - msg *db.Message - seq uint32 - newRead bool - readSet bool - newFlagged bool - flaggedSet bool - newDeleted bool - deletedSet bool - } - var changes []change - - for i := range dbMessages { - dbMsg := &dbMessages[i] - var match bool - if uid { - match = seqset.Contains(uint32(dbMsg.ID)) - } else { - match = seqset.Contains(uint32(i + 1)) - } - if !match { - continue - } - - c := change{msg: dbMsg, seq: uint32(i + 1)} - apply := func(flag string, enabled bool) { - switch flag { - case "\\Seen": - c.newRead, c.readSet = enabled, true - case "\\Flagged": - c.newFlagged, c.flaggedSet = enabled, true - case "\\Deleted": - c.newDeleted, c.deletedSet = enabled, true - } - } - switch op { - case imap.SetFlags: - apply("\\Seen", flagSet["\\Seen"]) - apply("\\Flagged", flagSet["\\Flagged"]) - apply("\\Deleted", flagSet["\\Deleted"]) - case imap.AddFlags: - for flag := range flagSet { - apply(flag, true) - } - case imap.RemoveFlags: - for flag := range flagSet { - apply(flag, false) - } - } - changes = append(changes, c) - } - - // 批量持久化(已读/星标合并为单条 UPDATE ... IN) - var readTrue, readFalse, flagTrue, flagFalse []uint - for _, c := range changes { - if c.readSet && c.newRead != c.msg.IsRead { - if c.newRead { - readTrue = append(readTrue, c.msg.ID) - } else { - readFalse = append(readFalse, c.msg.ID) - } - } - if c.flaggedSet && c.newFlagged != c.msg.IsFlagged { - if c.newFlagged { - flagTrue = append(flagTrue, c.msg.ID) - } else { - flagFalse = append(flagFalse, c.msg.ID) - } - } - } - // 记录首个持久化错误:SQLite 忙/锁等瞬时失败必须让客户端感知 - // (返回 NO 触发重试),否则已读/星标会静默丢失。 - var firstErr error - mark := func(err error) { - if err != nil && firstErr == nil { - firstErr = err - } - } - if len(readTrue) > 0 { - mark(m.stores.Mails.SetReadStates(readTrue, true)) - } - if len(readFalse) > 0 { - mark(m.stores.Mails.SetReadStates(readFalse, false)) - } - if len(flagTrue) > 0 { - mark(m.stores.Mails.SetFlaggedStates(flagTrue, true)) - } - if len(flagFalse) > 0 { - mark(m.stores.Mails.SetFlaggedStates(flagFalse, false)) - } - - // 标志变化 → 推送给同用户其他客户端(状态取目标值,序号用已加载 - // 列表的下标——与 ListMessages/Status 全链路一致,不再重复全量扫描) - for _, c := range changes { - if c.deletedSet { - if c.newDeleted { - m.deleted[c.msg.ID] = true - } else { - delete(m.deleted, c.msg.ID) - } - } - read := c.msg.IsRead - if c.readSet { - read = c.newRead - } - flagged := c.msg.IsFlagged - if c.flaggedSet { - flagged = c.newFlagged - } - deleted := m.deleted != nil && m.deleted[c.msg.ID] - pushUpdate(m.user.updates, buildFlagsUpdateAt(c.seq, m.user.email, m.name, c.msg, read, flagged, deleted)) - } - - return firstErr -} - -// CopyMessages copies messages to another mailbox. -func (m *imapMailbox) CopyMessages(uid bool, seqset *imap.SeqSet, dest string) error { - dest, ok := canonicalMailboxName(dest) - if !ok { - return backend.ErrNoSuchMailbox - } - - dbMessages, err := m.stores.Mails.ListAllByUserAndFolder(m.user.id, m.name) - if err != nil { - return err - } - - for i, dbMsg := range dbMessages { - var match bool - if uid { - match = seqset.Contains(uint32(dbMsg.ID)) - } else { - match = seqset.Contains(uint32(i + 1)) - } - if !match { - continue - } - - // Create a copy in the destination mailbox - copyMsg := &db.Message{ - UserID: m.user.id, - MessageID: dbMsg.MessageID, - Folder: dest, - FromAddr: dbMsg.FromAddr, - ToAddr: dbMsg.ToAddr, - CcAddr: dbMsg.CcAddr, - Subject: dbMsg.Subject, - TextBody: dbMsg.TextBody, - HtmlBody: dbMsg.HtmlBody, - RawData: dbMsg.RawData, - IsRead: dbMsg.IsRead, - IsFlagged: dbMsg.IsFlagged, - Date: dbMsg.Date, - } - if err := m.stores.Mails.Create(copyMsg); err != nil { - log.Printf("IMAP: failed to copy message %d to %s: %v", dbMsg.ID, dest, err) - continue - } - // 目标邮箱新增 → 推送给同用户其他客户端 - pushUpdate(m.user.updates, buildNewMessageUpdate(m.stores, m.user.email, dest, copyMsg)) - } - - return nil -} - -// MoveMessages moves messages to another mailbox. -func (m *imapMailbox) MoveMessages(uid bool, seqset *imap.SeqSet, dest string) error { - dest, ok := canonicalMailboxName(dest) - if !ok { - return backend.ErrNoSuchMailbox - } - - dbMessages, err := m.stores.Mails.ListAllByUserAndFolder(m.user.id, m.name) - if err != nil { - return err - } - - for i, dbMsg := range dbMessages { - var match bool - if uid { - match = seqset.Contains(uint32(dbMsg.ID)) - } else { - match = seqset.Contains(uint32(i + 1)) - } - if !match { - continue - } - if err := m.stores.Mails.MoveToFolder(dbMsg.ID, dest); err != nil { - log.Printf("IMAP: failed to move message %d to %s: %v", dbMsg.ID, dest, err) - continue - } - // 目标邮箱新增(移动后消息在 dest)→ 推送给同用户其他客户端 - if moved, err := m.stores.Mails.GetByID(dbMsg.ID); err == nil { - pushUpdate(m.user.updates, buildNewMessageUpdate(m.stores, m.user.email, dest, moved)) - } - } - return nil -} - -// Expunge permanently removes messages marked as \Deleted. -func (m *imapMailbox) Expunge() error { - if m.deleted == nil { - return nil - } - - // 删除前计算各消息的序号(Expunge 响应序号为删除前状态下的序号) - var seqs []uint32 - for msgID := range m.deleted { - if seq := seqOf(m.stores, m.user.id, m.name, msgID); seq > 0 { - seqs = append(seqs, seq) - } - } - - for msgID := range m.deleted { - if err := m.stores.Mails.Delete(msgID); err != nil { - log.Printf("IMAP: failed to expunge message %d: %v", msgID, err) - } - } - m.deleted = make(map[uint]bool) - - // 删除 → 推送给同用户其他客户端(每条序号一个 ExpungeUpdate) - for _, seq := range seqs { - pushUpdate(m.user.updates, &backend.ExpungeUpdate{ - Update: backend.NewUpdate(m.user.email, m.name), - SeqNum: seq, - }) - } - return nil -} - -// ---------- Helper functions ---------- - -// parseAddressList parses a comma-separated address string into imap.Address slice. -func parseAddressList(addrStr string) []*imap.Address { - if addrStr == "" { - return nil - } - - addresses, err := mail.ParseAddressList(addrStr) - if err != nil { - // Fallback: treat the whole string as a single address - return []*imap.Address{{ - MailboxName: addrStr, - HostName: "", - }} - } - - result := make([]*imap.Address, 0, len(addresses)) - for _, addr := range addresses { - parts := strings.SplitN(addr.Address, "@", 2) - mailbox := parts[0] - host := "" - if len(parts) > 1 { - host = parts[1] - } - result = append(result, &imap.Address{ - PersonalName: addr.Name, - MailboxName: mailbox, - HostName: host, - }) - } - return result -} - -// buildRawMessage reconstructs a raw RFC822 message from a db.Message. -// If RawData is available, it uses the original raw data directly. -func buildRawMessage(msg *db.Message) []byte { - // 优先使用原始邮件数据 - if msg.RawData != "" { - return []byte(msg.RawData) - } - - // 降级:从字段重建 - var buf bytes.Buffer - - // Write headers - buf.WriteString(fmt.Sprintf("From: %s\r\n", msg.FromAddr)) - buf.WriteString(fmt.Sprintf("To: %s\r\n", msg.ToAddr)) - if msg.CcAddr != "" { - buf.WriteString(fmt.Sprintf("Cc: %s\r\n", msg.CcAddr)) - } - buf.WriteString(fmt.Sprintf("Subject: %s\r\n", msg.Subject)) - buf.WriteString(fmt.Sprintf("Date: %s\r\n", msg.Date.Format(time.RFC1123Z))) - if msg.MessageID != "" { - buf.WriteString(fmt.Sprintf("Message-ID: %s\r\n", msg.MessageID)) - } - buf.WriteString("MIME-Version: 1.0\r\n") - - // Write body - if msg.HtmlBody != "" && msg.TextBody != "" { - boundary := fmt.Sprintf("mailgo_%d", msg.ID) - buf.WriteString(fmt.Sprintf("Content-Type: multipart/alternative; boundary=\"%s\"\r\n", boundary)) - buf.WriteString("\r\n") - buf.WriteString(fmt.Sprintf("--%s\r\n", boundary)) - buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n\r\n") - buf.WriteString(msg.TextBody) - buf.WriteString("\r\n") - buf.WriteString(fmt.Sprintf("--%s\r\n", boundary)) - buf.WriteString("Content-Type: text/html; charset=utf-8\r\n\r\n") - buf.WriteString(msg.HtmlBody) - buf.WriteString("\r\n") - buf.WriteString(fmt.Sprintf("--%s--\r\n", boundary)) - } else if msg.TextBody != "" { - buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n\r\n") - buf.WriteString(msg.TextBody) - } else if msg.HtmlBody != "" { - buf.WriteString("Content-Type: text/html; charset=utf-8\r\n\r\n") - buf.WriteString(msg.HtmlBody) - } else { - buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n\r\n") - } - - return buf.Bytes() -} - -// buildRawHeader reconstructs just the header portion of a raw RFC822 message. -func buildRawHeader(msg *db.Message) []byte { - var buf bytes.Buffer - buf.WriteString(fmt.Sprintf("From: %s\r\n", msg.FromAddr)) - buf.WriteString(fmt.Sprintf("To: %s\r\n", msg.ToAddr)) - if msg.CcAddr != "" { - buf.WriteString(fmt.Sprintf("Cc: %s\r\n", msg.CcAddr)) - } - buf.WriteString(fmt.Sprintf("Subject: %s\r\n", msg.Subject)) - buf.WriteString(fmt.Sprintf("Date: %s\r\n", msg.Date.Format(time.RFC1123Z))) - if msg.MessageID != "" { - buf.WriteString(fmt.Sprintf("Message-ID: %s\r\n", msg.MessageID)) - } - buf.WriteString("MIME-Version: 1.0\r\n") - if msg.HtmlBody != "" && msg.TextBody != "" { - boundary := fmt.Sprintf("mailgo_%d", msg.ID) - buf.WriteString(fmt.Sprintf("Content-Type: multipart/alternative; boundary=\"%s\"\r\n", boundary)) - } else if msg.TextBody != "" { - buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n") - } else if msg.HtmlBody != "" { - buf.WriteString("Content-Type: text/html; charset=utf-8\r\n") - } - return buf.Bytes() -} diff --git a/internal/imap_server/integration_test.go b/internal/imap_server/integration_test.go index c9bde59..122b47a 100644 --- a/internal/imap_server/integration_test.go +++ b/internal/imap_server/integration_test.go @@ -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) } } diff --git a/internal/imap_server/notify_test.go b/internal/imap_server/notify_test.go index e9aa374..8ecfa5b 100644 --- a/internal/imap_server/notify_test.go +++ b/internal/imap_server/notify_test.go @@ -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: "", - 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]) } } diff --git a/internal/imap_server/server.go b/internal/imap_server/server.go index daed02b..f70e180 100644 --- a/internal/imap_server/server.go +++ b/internal/imap_server/server.go @@ -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) diff --git a/internal/imap_server/session.go b/internal/imap_server/session.go new file mode 100644 index 0000000..2e493e9 --- /dev/null +++ b/internal/imap_server/session.go @@ -0,0 +1,1404 @@ +package imap_server + +import ( + "bufio" + "bytes" + "crypto/tls" + "fmt" + "io" + "log" + "net/mail" + "regexp" + "strings" + "sync" + "time" + + "mail_go/internal/connhub" + "mail_go/internal/db" + "mail_go/internal/mailutil" + "mail_go/internal/store" + + "github.com/emersion/go-imap/v2" + "github.com/emersion/go-imap/v2/imapserver" + asgomail "github.com/emersion/go-message/mail" + "github.com/emersion/go-message/textproto" +) + +// ---------- 跨会话推送 ---------- + +// sessionUpdate 是一条待下发给客户端的邮箱状态更新。 +type sessionUpdate struct { + expunge *uint32 // EXPUNGE 序号 + exists *uint32 // EXISTS 计数 + fetch *sessionFetchUpdate +} + +type sessionFetchUpdate struct { + seq uint32 + uid imap.UID + flags []imap.Flag +} + +// mailboxHub 跟踪一个(用户+文件夹)邮箱的所有已选中会话,把其他会话 +// 产生的变化(EXISTS/EXPUNGE/FLAGS)推送给其余会话(排除来源会话)。 +// 自行实现而非复用 imapserver.MailboxTracker:库版本对 EXPUNGE/EXISTS +// 无法排除来源会话,会导致本会话 EXPUNGE 时给自己排队重复响应。 +type mailboxHub struct { + mu sync.Mutex + sessions map[*imapSession]struct{} +} + +func newMailboxHub() *mailboxHub { + return &mailboxHub{sessions: make(map[*imapSession]struct{})} +} + +func (h *mailboxHub) add(s *imapSession) { + h.mu.Lock() + h.sessions[s] = struct{}{} + h.mu.Unlock() +} + +func (h *mailboxHub) remove(s *imapSession) { + h.mu.Lock() + delete(h.sessions, s) + h.mu.Unlock() +} + +// enqueue 把更新分发给除 source 外的所有会话。 +func (h *mailboxHub) enqueue(u sessionUpdate, source *imapSession) { + h.mu.Lock() + targets := make([]*imapSession, 0, len(h.sessions)) + for s := range h.sessions { + if s != source { + targets = append(targets, s) + } + } + h.mu.Unlock() + for _, s := range targets { + s.pushUpdate(u) + } +} + +// ---------- imapSession ---------- + +// imapSession 实现 imapserver.Session 与 imapserver.SessionMove。 +// 每个连接一个会话实例;所有命令串行执行,跨会话更新经 mailboxHub 推送。 +type imapSession struct { + srv *IMAPServer + conn *imapserver.Conn + + port int // 监听端口(协议日志) + remoteAddr string // 远端地址字符串(封禁/断开用) + + mu sync.Mutex + userID uint + email string + logID uint + startedAt time.Time + clientIP string + hubConn *connhub.Conn + + // selected state + selected string + readOnly bool + hub *mailboxHub + + // 待下发的推送队列(Idle/Poll 时写出) + queue []sessionUpdate + notify chan struct{} +} + +func newImapSession(srv *IMAPServer, conn *imapserver.Conn, port int) *imapSession { + s := &imapSession{ + srv: srv, + conn: conn, + port: port, + startedAt: time.Now(), + notify: make(chan struct{}, 1), + } + if nc := conn.NetConn(); nc != nil && nc.RemoteAddr() != nil { + s.remoteAddr = nc.RemoteAddr().String() + } + return s +} + +// pushUpdate 入队一条推送并唤醒 Idle(非阻塞,队列共享锁保护)。 +func (s *imapSession) pushUpdate(u sessionUpdate) { + s.mu.Lock() + s.queue = append(s.queue, u) + s.mu.Unlock() + select { + case s.notify <- struct{}{}: + default: + } +} + +// takeUpdates 取出待下发更新;allowExpunge=false 时遇到 EXPUNGE 停止 +// (与 RFC 要求一致:FETCH/STORE/SEARCH 期间不能下发 EXPUNGE)。 +func (s *imapSession) takeUpdates(allowExpunge bool) []sessionUpdate { + s.mu.Lock() + defer s.mu.Unlock() + if allowExpunge { + updates := s.queue + s.queue = nil + return updates + } + stop := -1 + for i, u := range s.queue { + if u.expunge != nil { + stop = i + break + } + } + if stop < 0 { + updates := s.queue + s.queue = nil + return updates + } + updates := s.queue[:stop] + s.queue = s.queue[stop:] + return updates +} + +// writeUpdates 把更新按类型写入连接。 +func writeUpdates(w *imapserver.UpdateWriter, updates []sessionUpdate) error { + for _, u := range updates { + switch { + case u.expunge != nil: + if err := w.WriteExpunge(*u.expunge); err != nil { + return err + } + case u.exists != nil: + if err := w.WriteNumMessages(*u.exists); err != nil { + return err + } + case u.fetch != nil: + if err := w.WriteMessageFlags(u.fetch.seq, u.fetch.uid, u.fetch.flags); err != nil { + return err + } + } + } + return nil +} + +// ---------- Session 基础 ---------- + +// Close 在连接结束时调用:注销推送、关闭连接中心记录、回填协议日志。 +func (s *imapSession) Close() error { + s.mu.Lock() + hub := s.hub + s.hub = nil + hubConn := s.hubConn + s.hubConn = nil + logID := s.logID + startedAt := s.startedAt + s.mu.Unlock() + + if hub != nil { + hub.remove(s) + } + if hubConn != nil { + hubConn.Close() + } + if logID != 0 { + s.srv.stores.ProtocolLogs.UpdateDuration(logID, time.Since(startedAt).Milliseconds()) + } + s.srv.unregisterSession(s) + return nil +} + +// Login 认证:封禁检查 → 认证 → 失败计数 → 协议日志 → 连接中心注册。 +// 与 v1 imapBackend.Login 完全一致的策略(含裸用户名支持、成功后清零失败计数)。 +func (s *imapSession) Login(username, password string) error { + clientIP := store.ClientIPFromAddr(s.conn.NetConn().RemoteAddr()) + now := time.Now() + + if banned, _ := s.srv.stores.Bans.IsBanned(clientIP); banned { + s.recordLogin(clientIP, username, false, "IP已被封禁", "认证被拒绝(IP 已封禁)", now) + return imapserver.ErrAuthFailed + } + + user, err := s.srv.stores.Users.AuthenticateLogin(username, password) + if err != nil { + s.srv.stores.RecordAuthFailure(clientIP, s.srv.banCfg.MaxFailAttempts, s.srv.banCfg.BanDurationMin, "邮件协议认证失败次数过多") + s.recordLogin(clientIP, username, false, "用户名或密码错误", "LOGIN 失败", now) + return imapserver.ErrAuthFailed + } + + // 登录成功清零失败计数(与 Web 登录一致),避免协议客户端失败计数 + // 只增不减(配置探测、APP 裸用户名重试等)累计触发误封。 + s.srv.stores.Bans.ResetFail(clientIP) + + email := user.Username + "@" + if domain, err := s.srv.stores.Domains.GetByID(user.DomainID); err == nil { + email = user.Username + "@" + domain.Name + } + + logID := s.recordLogin(clientIP, username, true, "", "LOGIN 成功", now) + + // STARTTLS 后连接对象会被库替换为 tls.Conn,此处按当前实际加密状态标记。 + tlsOn := false + if _, ok := s.conn.NetConn().(*tls.Conn); ok { + tlsOn = true + } + conn := s.srv.hub.Register("imap", clientIP, s.port, tlsOn) + if conn != nil { + conn.SetUser(email) + if s.remoteAddr != "" { + addr := s.remoteAddr + conn.SetDisconnect(func() { s.srv.DisconnectByAddr(addr) }) + } + } + + s.mu.Lock() + s.userID = user.ID + s.email = email + s.logID = logID + s.clientIP = clientIP + s.hubConn = conn + s.mu.Unlock() + return nil +} + +// recordLogin 写入一条 IMAP 协议日志,返回记录 ID(失败为 0)。 +func (s *imapSession) recordLogin(ip, username string, success bool, failReason, detail string, at time.Time) uint { + entry := &db.ProtocolLog{ + Protocol: db.ProtocolIMAP, + Port: s.port, + ClientIP: ip, + Username: username, + Success: success, + FailReason: failReason, + Detail: detail, + CreatedAt: at, + } + if err := s.srv.stores.ProtocolLogs.Create(entry); err != nil { + log.Printf("IMAP: 写入协议日志失败: %v", err) + return 0 + } + return entry.ID +} + +// systemMailboxes 是系统支持的全部文件夹。 +var systemMailboxes = []struct { + name string + attrs []imap.MailboxAttr +}{ + {"INBOX", nil}, + {"Sent", []imap.MailboxAttr{imap.MailboxAttrSent}}, + {"Drafts", []imap.MailboxAttr{imap.MailboxAttrDrafts}}, + {"Trash", []imap.MailboxAttr{imap.MailboxAttrTrash}}, +} + +// matchPattern 按 RFC 3501 LIST wildcard 语义匹配(* 任意、% 不跨分隔符)。 +// INBOX 大小写不敏感,其余文件夹按系统定义大小写精确匹配。 +func matchPattern(name, pattern string) bool { + var sb strings.Builder + sb.WriteString("^") + for _, r := range pattern { + switch r { + case '*': + sb.WriteString(".*") + case '%': + sb.WriteString("[^/]*") + default: + sb.WriteString(regexp.QuoteMeta(string(r))) + } + } + sb.WriteString("$") + re, err := regexp.Compile(sb.String()) + if err != nil { + return false + } + if strings.EqualFold(name, "INBOX") { + return re.MatchString(strings.ToUpper(name)) + } + return re.MatchString(name) +} + +// matchAnyPattern 判断名字是否命中任意一个 pattern。 +func matchAnyPattern(name string, patterns []string) bool { + for _, p := range patterns { + if matchPattern(name, p) { + return true + } + } + return false +} + +// List 返回系统文件夹列表(LSUB 同样返回全部——订阅状态恒为已订阅)。 +func (s *imapSession) List(w *imapserver.ListWriter, ref string, patterns []string, options *imap.ListOptions) error { + for _, mb := range systemMailboxes { + if !matchAnyPattern(mb.name, patterns) { + continue + } + data := &imap.ListData{ + Attrs: mb.attrs, + Delim: '/', + Mailbox: mb.name, + } + if err := w.WriteList(data); err != nil { + return err + } + } + return nil +} + +// Create 不支持。 +func (s *imapSession) Create(mailbox string, options *imap.CreateOptions) error { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "mailbox creation not supported"} +} + +// Delete 不支持。 +func (s *imapSession) Delete(mailbox string) error { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "mailbox deletion not supported"} +} + +// Rename 不支持。 +func (s *imapSession) Rename(mailbox, newName string, options *imap.RenameOptions) error { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "mailbox rename not supported"} +} + +// Subscribe 恒为已订阅(no-op)。 +func (s *imapSession) Subscribe(mailbox string) error { + return nil +} + +// Unsubscribe no-op。 +func (s *imapSession) Unsubscribe(mailbox string) error { + return nil +} + +// ---------- 选中状态 ---------- + +// Select 选中邮箱:登记 mailboxHub 会话,返回状态数据。 +func (s *imapSession) Select(mailbox string, options *imap.SelectOptions) (*imap.SelectData, error) { + name, ok := canonicalMailboxName(mailbox) + if !ok { + return nil, &imap.Error{Type: imap.StatusResponseTypeNo, Text: "No such mailbox"} + } + + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, name) + if err != nil { + return nil, err + } + var unseen uint32 + for i := range msgs { + if !msgs[i].IsRead { + unseen++ + } + } + maxID, err := s.srv.stores.Mails.MaxIDByUserAndFolder(userID, name) + if err != nil { + return nil, err + } + uidValidity, err := s.srv.stores.MailboxState.UidValidity(userID, name) + if err != nil { + log.Printf("IMAP: 获取 UIDVALIDITY 失败 folder=%s: %v", name, err) + uidValidity = 1 + } + + // 换绑推送(若此前已选中,先注销旧邮箱) + s.mu.Lock() + oldHub := s.hub + s.hub = nil + s.mu.Unlock() + if oldHub != nil { + oldHub.remove(s) + } + hub := s.srv.hubForOrCreate(s.currentEmail(), name) + hub.add(s) + + s.mu.Lock() + s.selected = name + s.readOnly = options.ReadOnly + s.hub = hub + s.mu.Unlock() + + flags := []imap.Flag{imap.FlagAnswered, imap.FlagFlagged, imap.FlagDeleted, imap.FlagSeen, imap.FlagDraft} + return &imap.SelectData{ + Flags: flags, + PermanentFlags: append(flags, imap.FlagWildcard), + NumMessages: uint32(len(msgs)), + NumRecent: 0, + UIDNext: imap.UID(maxID + 1), + UIDValidity: uidValidity, + }, nil +} + +// Unselect 取消选中:注销 mailboxHub 会话。 +func (s *imapSession) Unselect() error { + s.mu.Lock() + hub := s.hub + s.hub = nil + s.selected = "" + s.readOnly = false + s.queue = nil + s.mu.Unlock() + if hub != nil { + hub.remove(s) + } + return nil +} + +// Status 返回邮箱状态(计数/未读/UIDNEXT/UIDVALIDITY/已删除标记数/大小)。 +func (s *imapSession) Status(mailbox string, options *imap.StatusOptions) (*imap.StatusData, error) { + name, ok := canonicalMailboxName(mailbox) + if !ok { + return nil, &imap.Error{Type: imap.StatusResponseTypeNo, Text: "No such mailbox"} + } + userID := s.currentUserID() + + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, name) + if err != nil { + return nil, err + } + data := &imap.StatusData{Mailbox: name} + + if options.NumMessages || options.NumUnseen || options.NumDeleted || options.Size { + var unseen, deleted uint32 + var size int64 + for i := range msgs { + if !msgs[i].IsRead { + unseen++ + } + if msgs[i].IsDeleted { + deleted++ + } + size += int64(len(messageRawData(&msgs[i]))) + } + if options.NumMessages { + n := uint32(len(msgs)) + data.NumMessages = &n + } + if options.NumUnseen { + data.NumUnseen = &unseen + } + if options.NumDeleted { + data.NumDeleted = &deleted + } + if options.Size { + data.Size = &size + } + } + if options.NumRecent { + zero := uint32(0) + data.NumRecent = &zero + } + if options.UIDNext { + maxID, err := s.srv.stores.Mails.MaxIDByUserAndFolder(userID, name) + if err != nil { + return nil, err + } + data.UIDNext = imap.UID(maxID + 1) + } + if options.UIDValidity { + uidValidity, err := s.srv.stores.MailboxState.UidValidity(userID, name) + if err != nil { + return nil, err + } + data.UIDValidity = uidValidity + } + return data, nil +} + +// ---------- 消息操作 ---------- + +// numSetContains 判断消息是否命中序号/UID 集合(含 "*" 动态上界)。 +func numSetContains(numSet imap.NumSet, seq uint32, uid imap.UID, maxSeq, maxUID uint32) bool { + switch set := numSet.(type) { + case imap.SeqSet: + for _, r := range set { + stop := r.Stop + if stop == 0 { + stop = maxSeq + } + if seq >= r.Start && seq <= stop { + return true + } + } + case imap.UIDSet: + for _, r := range set { + stop := r.Stop + if stop == 0 { + stop = imap.UID(maxUID) + } + if uid >= r.Start && uid <= stop { + return true + } + } + } + return false +} + +// Fetch 下发匹配消息的 FETCH 响应。 +func (s *imapSession) Fetch(w *imapserver.FetchWriter, numSet imap.NumSet, options *imap.FetchOptions) error { + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, s.currentMailbox()) + if err != nil { + return err + } + maxSeq := uint32(len(msgs)) + var maxUID uint32 + for i := range msgs { + if uint32(msgs[i].ID) > maxUID { + maxUID = uint32(msgs[i].ID) + } + } + + for i := range msgs { + msg := &msgs[i] + seq := uint32(i + 1) + uid := imap.UID(msg.ID) + if !numSetContains(numSet, seq, uid, maxSeq, maxUID) { + continue + } + + raw := messageRawData(msg) + fw := w.CreateMessage(seq) + if options.UID { + fw.WriteUID(uid) + } + if options.Flags { + fw.WriteFlags(flagsOf(msg.IsRead, msg.IsFlagged, msg.IsDeleted)) + } + if options.RFC822Size { + fw.WriteRFC822Size(int64(len(raw))) + } + if options.InternalDate { + fw.WriteInternalDate(msg.Date) + } + if options.Envelope { + var env *imap.Envelope + if hdr, _, err := headerAndBody(raw); err == nil { + env = imapserver.ExtractEnvelope(hdr) + } + if env == nil { + env = envelopeFromDB(msg) + } + fw.WriteEnvelope(env) + } + if options.BodyStructure != nil { + bs := imapserver.ExtractBodyStructure(bytes.NewReader(raw)) + if bs == nil { + // 防御:畸形 MIME 解析失败时降级为 text/plain 单段 + bs = fallbackBodyStructure(raw) + } + fw.WriteBodyStructure(bs) + } + for _, section := range options.BodySection { + b := imapserver.ExtractBodySection(bytes.NewReader(raw), section) + lit := fw.WriteBodySection(section, int64(len(b))) + lit.Write(b) + lit.Close() + } + if err := fw.Close(); err != nil { + return err + } + } + return nil +} + +// Search 按条件返回匹配消息的序号/UID 集合。 +func (s *imapSession) Search(kind imapserver.NumKind, criteria *imap.SearchCriteria, options *imap.SearchOptions) (*imap.SearchData, error) { + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, s.currentMailbox()) + if err != nil { + return nil, err + } + + data := &imap.SearchData{} + var seqSet imap.SeqSet + var uidSet imap.UIDSet + for i := range msgs { + msg := &msgs[i] + if !matchesCriteria(msg, criteria) { + continue + } + if kind == imapserver.NumKindSeq { + seqSet.AddNum(uint32(i + 1)) + } else { + uidSet.AddNum(imap.UID(msg.ID)) + } + } + if kind == imapserver.NumKindSeq { + data.All = seqSet + } else { + data.All = uidSet + } + return data, nil +} + +// matchesCriteria 判断消息是否满足搜索条件(\Deleted 读数据库持久化标记)。 +func matchesCriteria(msg *db.Message, criteria *imap.SearchCriteria) bool { + for _, flag := range criteria.Flag { + switch flag { + case imap.FlagSeen: + if !msg.IsRead { + return false + } + case imap.FlagFlagged: + if !msg.IsFlagged { + return false + } + case imap.FlagDeleted: + if !msg.IsDeleted { + return false + } + case imap.FlagAnswered, imap.FlagDraft: + return false + } + } + for _, flag := range criteria.NotFlag { + switch flag { + case imap.FlagSeen: + if msg.IsRead { + return false + } + case imap.FlagFlagged: + if msg.IsFlagged { + return false + } + case imap.FlagDeleted: + if msg.IsDeleted { + return false + } + } + } + + if !criteria.Since.IsZero() && msg.Date.Before(criteria.Since) { + return false + } + if !criteria.Before.IsZero() && !msg.Date.Before(criteria.Before) { + return false + } + + for _, f := range criteria.Header { + switch strings.ToLower(f.Key) { + case "subject": + if !strings.Contains(strings.ToLower(msg.Subject), strings.ToLower(f.Value)) { + return false + } + case "from": + if !strings.Contains(strings.ToLower(msg.FromAddr), strings.ToLower(f.Value)) { + return false + } + case "to": + if !strings.Contains(strings.ToLower(msg.ToAddr), strings.ToLower(f.Value)) { + return false + } + case "cc": + if !strings.Contains(strings.ToLower(msg.CcAddr), strings.ToLower(f.Value)) { + return false + } + } + } + + rawSize := int64(len(messageRawData(msg))) + if criteria.Larger > 0 && rawSize <= criteria.Larger { + return false + } + if criteria.Smaller > 0 && rawSize >= criteria.Smaller { + return false + } + + for _, text := range criteria.Body { + bodyText := strings.ToLower(msg.TextBody + " " + msg.HtmlBody) + if !strings.Contains(bodyText, strings.ToLower(text)) { + return false + } + } + for _, text := range criteria.Text { + allText := strings.ToLower(msg.Subject + " " + msg.FromAddr + " " + msg.ToAddr + " " + msg.TextBody + " " + msg.HtmlBody) + if !strings.Contains(allText, strings.ToLower(text)) { + return false + } + } + + for _, notCrit := range criteria.Not { + if matchesCriteria(msg, ¬Crit) { + return false + } + } + for _, orPair := range criteria.Or { + if !matchesCriteria(msg, &orPair[0]) && !matchesCriteria(msg, &orPair[1]) { + return false + } + } + + return true +} + +// Store 修改消息标志:批量持久化(含 \Deleted,落库不丢失)并推送其他会话。 +func (s *imapSession) Store(w *imapserver.FetchWriter, numSet imap.NumSet, flags *imap.StoreFlags, options *imap.StoreOptions) error { + if s.isReadOnly() { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "Mailbox is read-only"} + } + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, s.currentMailbox()) + if err != nil { + return err + } + maxSeq := uint32(len(msgs)) + var maxUID uint32 + for i := range msgs { + if uint32(msgs[i].ID) > maxUID { + maxUID = uint32(msgs[i].ID) + } + } + + flagSet := make(map[imap.Flag]bool, len(flags.Flags)) + for _, flag := range flags.Flags { + flagSet[flag] = true + } + + type change struct { + msg *db.Message + seq uint32 + newRead bool + readSet bool + newFlagged bool + flaggedSet bool + newDeleted bool + deletedSet bool + } + var changes []change + + for i := range msgs { + msg := &msgs[i] + seq := uint32(i + 1) + uid := imap.UID(msg.ID) + if !numSetContains(numSet, seq, uid, maxSeq, maxUID) { + continue + } + c := change{msg: msg, seq: seq} + apply := func(flag imap.Flag, enabled bool) { + switch flag { + case imap.FlagSeen: + c.newRead, c.readSet = enabled, true + case imap.FlagFlagged: + c.newFlagged, c.flaggedSet = enabled, true + case imap.FlagDeleted: + c.newDeleted, c.deletedSet = enabled, true + } + } + switch flags.Op { + case imap.StoreFlagsSet: + apply(imap.FlagSeen, flagSet[imap.FlagSeen]) + apply(imap.FlagFlagged, flagSet[imap.FlagFlagged]) + apply(imap.FlagDeleted, flagSet[imap.FlagDeleted]) + case imap.StoreFlagsAdd: + for flag := range flagSet { + apply(flag, true) + } + case imap.StoreFlagsDel: + for flag := range flagSet { + apply(flag, false) + } + } + changes = append(changes, c) + } + + // 批量持久化(已读/星标/删除各合并为单条 UPDATE ... IN) + var readTrue, readFalse, flagTrue, flagFalse, delTrue, delFalse []uint + for _, c := range changes { + if c.readSet && c.newRead != c.msg.IsRead { + if c.newRead { + readTrue = append(readTrue, c.msg.ID) + } else { + readFalse = append(readFalse, c.msg.ID) + } + } + if c.flaggedSet && c.newFlagged != c.msg.IsFlagged { + if c.newFlagged { + flagTrue = append(flagTrue, c.msg.ID) + } else { + flagFalse = append(flagFalse, c.msg.ID) + } + } + if c.deletedSet && c.newDeleted != c.msg.IsDeleted { + if c.newDeleted { + delTrue = append(delTrue, c.msg.ID) + } else { + delFalse = append(delFalse, c.msg.ID) + } + } + } + // 记录首个持久化错误:SQLite 忙/锁等瞬时失败必须让客户端感知 + // (返回 NO 触发重试),否则标志会静默丢失。 + var firstErr error + mark := func(err error) { + if err != nil && firstErr == nil { + firstErr = err + } + } + mark(s.srv.stores.Mails.SetReadStates(readTrue, true)) + mark(s.srv.stores.Mails.SetReadStates(readFalse, false)) + mark(s.srv.stores.Mails.SetFlaggedStates(flagTrue, true)) + mark(s.srv.stores.Mails.SetFlaggedStates(flagFalse, false)) + mark(s.srv.stores.Mails.SetDeletedStates(delTrue, true)) + mark(s.srv.stores.Mails.SetDeletedStates(delFalse, false)) + + // 标志变化 → 回写本连接(非 SILENT)并推送同用户其他会话 + hub := s.hubForSelected() + for _, c := range changes { + read := c.msg.IsRead + if c.readSet { + read = c.newRead + } + flagged := c.msg.IsFlagged + if c.flaggedSet { + flagged = c.newFlagged + } + deleted := c.msg.IsDeleted + if c.deletedSet { + deleted = c.newDeleted + } + fl := flagsOf(read, flagged, deleted) + if !flags.Silent { + fw := w.CreateMessage(c.seq) + fw.WriteUID(imap.UID(c.msg.ID)) + fw.WriteFlags(fl) + if err := fw.Close(); err != nil && firstErr == nil { + firstErr = err + } + } + if hub != nil { + hub.enqueue(sessionUpdate{fetch: &sessionFetchUpdate{seq: c.seq, uid: imap.UID(c.msg.ID), flags: fl}}, s) + } + } + return firstErr +} + +// Copy 把匹配消息复制到目标文件夹。 +func (s *imapSession) Copy(numSet imap.NumSet, dest string) (*imap.CopyData, error) { + destName, ok := canonicalMailboxName(dest) + if !ok { + return nil, &imap.Error{Type: imap.StatusResponseTypeNo, Text: "No such mailbox"} + } + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, s.currentMailbox()) + if err != nil { + return nil, err + } + maxSeq := uint32(len(msgs)) + var maxUID uint32 + for i := range msgs { + if uint32(msgs[i].ID) > maxUID { + maxUID = uint32(msgs[i].ID) + } + } + + var sourceUIDs, destUIDs imap.UIDSet + for i := range msgs { + dbMsg := &msgs[i] + seq := uint32(i + 1) + uid := imap.UID(dbMsg.ID) + if !numSetContains(numSet, seq, uid, maxSeq, maxUID) { + continue + } + copyMsg := &db.Message{ + UserID: userID, + MessageID: dbMsg.MessageID, + Folder: destName, + FromAddr: dbMsg.FromAddr, + ToAddr: dbMsg.ToAddr, + CcAddr: dbMsg.CcAddr, + Subject: dbMsg.Subject, + TextBody: dbMsg.TextBody, + HtmlBody: dbMsg.HtmlBody, + RawData: dbMsg.RawData, + IsRead: dbMsg.IsRead, + IsFlagged: dbMsg.IsFlagged, + Date: dbMsg.Date, + } + if err := s.srv.stores.Mails.Create(copyMsg); err != nil { + log.Printf("IMAP: failed to copy message %d to %s: %v", dbMsg.ID, destName, err) + continue + } + sourceUIDs.AddNum(uid) + destUIDs.AddNum(imap.UID(copyMsg.ID)) + } + + // 目标文件夹会话收到 EXISTS 通知 + if hub := s.srv.hubFor(s.currentEmail(), destName); hub != nil { + if count, err := s.srv.stores.Mails.CountByUserAndFolder(userID, destName); err == nil { + hub.enqueue(sessionUpdate{exists: ptrU32(uint32(count))}, nil) + } + } + + uidValidity, err := s.srv.stores.MailboxState.UidValidity(userID, destName) + if err != nil { + uidValidity = 0 + } + return &imap.CopyData{ + UIDValidity: uidValidity, + SourceUIDs: sourceUIDs, + DestUIDs: destUIDs, + }, nil +} + +// Move 把匹配消息移动到目标文件夹(COPYUID + EXPUNGE 响应)。 +func (s *imapSession) Move(w *imapserver.MoveWriter, numSet imap.NumSet, dest string) error { + if s.isReadOnly() { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "Mailbox is read-only"} + } + destName, ok := canonicalMailboxName(dest) + if !ok { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "No such mailbox"} + } + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, s.currentMailbox()) + if err != nil { + return err + } + maxSeq := uint32(len(msgs)) + var maxUID uint32 + for i := range msgs { + if uint32(msgs[i].ID) > maxUID { + maxUID = uint32(msgs[i].ID) + } + } + + var sourceUIDs, destUIDs imap.UIDSet + type moved struct { + seq uint32 + uid imap.UID + } + var movedMsgs []moved + for i := range msgs { + dbMsg := &msgs[i] + seq := uint32(i + 1) + uid := imap.UID(dbMsg.ID) + if !numSetContains(numSet, seq, uid, maxSeq, maxUID) { + continue + } + if err := s.srv.stores.Mails.MoveToFolder(dbMsg.ID, destName); err != nil { + log.Printf("IMAP: failed to move message %d to %s: %v", dbMsg.ID, destName, err) + continue + } + sourceUIDs.AddNum(uid) + destUIDs.AddNum(uid) + movedMsgs = append(movedMsgs, moved{seq: seq, uid: uid}) + } + + uidValidity, err := s.srv.stores.MailboxState.UidValidity(userID, destName) + if err != nil { + uidValidity = 0 + } + if err := w.WriteCopyData(&imap.CopyData{ + UIDValidity: uidValidity, + SourceUIDs: sourceUIDs, + DestUIDs: destUIDs, + }); err != nil { + return err + } + for _, m := range movedMsgs { + if err := w.WriteExpunge(m.seq); err != nil { + return err + } + } + + // 源文件夹其余会话收到 EXPUNGE;目标文件夹会话收到 EXISTS + if hub := s.hubForSelected(); hub != nil { + for _, m := range movedMsgs { + hub.enqueue(sessionUpdate{expunge: ptrU32(m.seq)}, s) + } + } + if hub := s.srv.hubFor(s.currentEmail(), destName); hub != nil { + if count, err := s.srv.stores.Mails.CountByUserAndFolder(userID, destName); err == nil { + hub.enqueue(sessionUpdate{exists: ptrU32(uint32(count))}, nil) + } + } + return nil +} + +// Expunge 永久删除 \Deleted 消息。 +// uids == nil:删除该文件夹全部 \Deleted(EXPUNGE/CLOSE); +// uids != nil:只删除集合内且带 \Deleted 的消息(UID EXPUNGE,RFC 4315)。 +func (s *imapSession) Expunge(w *imapserver.ExpungeWriter, uids *imap.UIDSet) error { + if s.isReadOnly() { + return &imap.Error{Type: imap.StatusResponseTypeNo, Text: "Mailbox is read-only"} + } + userID := s.currentUserID() + msgs, err := s.srv.stores.Mails.ListAllByUserAndFolder(userID, s.currentMailbox()) + if err != nil { + return err + } + + // 删除前计算序号(EXPUNGE 响应序号为删除前状态下的序号) + type target struct { + id uint + seq uint32 + } + var targets []target + for i := range msgs { + msg := &msgs[i] + if !msg.IsDeleted { + continue + } + if uids != nil && !uids.Contains(imap.UID(msg.ID)) { + continue + } + targets = append(targets, target{id: msg.ID, seq: uint32(i + 1)}) + } + if len(targets) == 0 { + return nil + } + + ids := make([]uint, len(targets)) + for i, t := range targets { + ids[i] = t.id + } + if err := s.srv.stores.Mails.DeleteMany(ids); err != nil { + log.Printf("IMAP: failed to expunge %d messages: %v", len(ids), err) + return err + } + + // 本连接按序号升序写 EXPUNGE 响应,其余会话经 hub 推送 + for _, t := range targets { + if err := w.WriteExpunge(t.seq); err != nil { + return err + } + } + if hub := s.hubForSelected(); hub != nil { + for _, t := range targets { + hub.enqueue(sessionUpdate{expunge: ptrU32(t.seq)}, s) + } + } + return nil +} + +// Append 追加一封新邮件(IMAP APPEND)。 +func (s *imapSession) Append(mailbox string, r imap.LiteralReader, options *imap.AppendOptions) (*imap.AppendData, error) { + name, ok := canonicalMailboxName(mailbox) + if !ok { + return nil, &imap.Error{Type: imap.StatusResponseTypeNo, Text: "No such mailbox"} + } + userID := s.currentUserID() + + data, err := io.ReadAll(r) + if err != nil { + return nil, fmt.Errorf("failed to read message body: %w", err) + } + + mr, err := asgomail.CreateReader(bytes.NewReader(data)) + if err != nil { + return nil, fmt.Errorf("failed to parse MIME message: %w", err) + } + header := mr.Header + fromAddr := mailutil.FormatAddressList(&header, "From") + toAddr := mailutil.FormatAddressList(&header, "To") + ccAddr := mailutil.FormatAddressList(&header, "Cc") + subject, _ := header.Subject() + messageID, _ := header.MessageID() + msgDate, _ := header.Date() + if msgDate.IsZero() { + msgDate = options.Time + } + if msgDate.IsZero() { + msgDate = time.Now() + } + + var textBody, htmlBody string + for { + p, err := mr.NextPart() + if err == io.EOF { + break + } + if err != nil { + break + } + switch h := p.Header.(type) { + case *asgomail.InlineHeader: + contentType, params, _ := h.ContentType() + buf, _ := io.ReadAll(p.Body) + charset := "" + if cs, ok := params["charset"]; ok { + charset = cs + } + decoded := mailutil.DecodeCharset(buf, charset) + if strings.HasPrefix(contentType, "text/plain") { + textBody = decoded + } else if strings.HasPrefix(contentType, "text/html") { + htmlBody = decoded + } + case *asgomail.AttachmentHeader: + // APPEND 附件不落盘(与 v1 行为一致) + } + } + if textBody == "" && htmlBody == "" { + textBody = string(data) + } + + var isRead, isFlagged, isDeleted bool + for _, flag := range options.Flags { + switch flag { + case imap.FlagSeen: + isRead = true + case imap.FlagFlagged: + isFlagged = true + case imap.FlagDeleted: + isDeleted = true + } + } + + msg := &db.Message{ + UserID: userID, + MessageID: messageID, + Folder: name, + FromAddr: fromAddr, + ToAddr: toAddr, + CcAddr: ccAddr, + Subject: subject, + TextBody: textBody, + HtmlBody: htmlBody, + RawData: string(data), + IsRead: isRead, + IsFlagged: isFlagged, + IsDeleted: isDeleted, + Date: msgDate, + } + if err := s.srv.stores.Mails.Create(msg); err != nil { + return nil, fmt.Errorf("failed to create message: %w", err) + } + + // 该文件夹其余会话收到 EXISTS 通知 + if hub := s.srv.hubFor(s.currentEmail(), name); hub != nil { + if count, err := s.srv.stores.Mails.CountByUserAndFolder(userID, name); err == nil { + hub.enqueue(sessionUpdate{exists: ptrU32(uint32(count))}, nil) + } + } + + uidValidity, err := s.srv.stores.MailboxState.UidValidity(userID, name) + if err != nil { + uidValidity = 0 + } + return &imap.AppendData{UID: imap.UID(msg.ID), UIDValidity: uidValidity}, nil +} + +// Poll 下发积压的推送更新(NOOP/COPY/APPEND 等命令后由库调用)。 +func (s *imapSession) Poll(w *imapserver.UpdateWriter, allowExpunge bool) error { + return writeUpdates(w, s.takeUpdates(allowExpunge)) +} + +// Idle 阻塞直到收到推送或客户端发 DONE(库负责读取并关闭 stop)。 +func (s *imapSession) Idle(w *imapserver.UpdateWriter, stop <-chan struct{}) error { + for { + if err := writeUpdates(w, s.takeUpdates(true)); err != nil { + return err + } + select { + case <-stop: + return nil + case <-s.notify: + } + } +} + +// ---------- 会话内部辅助 ---------- + +func (s *imapSession) currentUserID() uint { + s.mu.Lock() + defer s.mu.Unlock() + return s.userID +} + +func (s *imapSession) currentEmail() string { + s.mu.Lock() + defer s.mu.Unlock() + return s.email +} + +func (s *imapSession) currentMailbox() string { + s.mu.Lock() + defer s.mu.Unlock() + return s.selected +} + +func (s *imapSession) isReadOnly() bool { + s.mu.Lock() + defer s.mu.Unlock() + return s.readOnly +} + +func (s *imapSession) hubForSelected() *mailboxHub { + s.mu.Lock() + defer s.mu.Unlock() + return s.hub +} + +func ptrU32(v uint32) *uint32 { + return &v +} + +// ---------- Helper functions ---------- + +// canonicalMailboxName 规范化邮箱名。 +func canonicalMailboxName(name string) (string, bool) { + switch strings.ToUpper(strings.TrimSpace(name)) { + case "INBOX": + return "INBOX", true + case "SENT": + return "Sent", true + case "DRAFTS": + return "Drafts", true + case "TRASH": + return "Trash", true + default: + return "", false + } +} + +// flagsOf 按数据库状态生成 IMAP 标志列表。 +func flagsOf(read, flagged, deleted bool) []imap.Flag { + flags := make([]imap.Flag, 0, 3) + if read { + flags = append(flags, imap.FlagSeen) + } + if flagged { + flags = append(flags, imap.FlagFlagged) + } + if deleted { + flags = append(flags, imap.FlagDeleted) + } + return flags +} + +// seqOf 返回消息在文件夹中的序号(1 基),未找到返回 0。 +func seqOf(stores *store.Stores, userID uint, mailbox string, msgID uint) uint32 { + msgs, err := stores.Mails.ListAllByUserAndFolder(userID, mailbox) + if err != nil { + return 0 + } + for i := range msgs { + if msgs[i].ID == msgID { + return uint32(i + 1) + } + } + return 0 +} + +// messageRawData 返回消息原始字节(优先 RawData,降级按字段重建)。 +func messageRawData(msg *db.Message) []byte { + if msg.RawData != "" { + return []byte(msg.RawData) + } + return buildRawMessage(msg) +} + +func headerAndBody(raw []byte) (textproto.Header, io.Reader, error) { + body := bufio.NewReader(bytes.NewReader(raw)) + hdr, err := textproto.ReadHeader(body) + return hdr, body, err +} + +// fallbackBodyStructure 构造一个 text/plain 单段 BodyStructure(含 +// Extended,保证 BODYSTRUCTURE 扩展请求不因 nil Extended 触发库 panic)。 +func fallbackBodyStructure(raw []byte) imap.BodyStructure { + size := uint32(len(raw)) + lines := int64(bytes.Count(raw, []byte{'\n'})) + if len(raw) > 0 && raw[len(raw)-1] != '\n' { + lines++ + } + return &imap.BodyStructureSinglePart{ + Type: "text", + Subtype: "plain", + Params: map[string]string{"charset": "utf-8"}, + Encoding: "8bit", + Size: size, + Text: &imap.BodyStructureText{NumLines: lines}, + Extended: &imap.BodyStructureSinglePartExt{}, + } +} + +// envelopeFromDB 从数据库字段构造信封(原始邮件头解析失败时的降级)。 +func envelopeFromDB(msg *db.Message) *imap.Envelope { + return &imap.Envelope{ + Date: msg.Date, + Subject: msg.Subject, + From: parseAddressList(msg.FromAddr), + Sender: parseAddressList(msg.FromAddr), + ReplyTo: parseAddressList(msg.FromAddr), + To: parseAddressList(msg.ToAddr), + Cc: parseAddressList(msg.CcAddr), + MessageID: msg.MessageID, + } +} + +// parseAddressList 解析逗号分隔的地址字符串为 imap.Address 切片。 +func parseAddressList(addrStr string) []imap.Address { + if addrStr == "" { + return nil + } + addresses, err := mail.ParseAddressList(addrStr) + if err != nil { + return []imap.Address{{Mailbox: addrStr}} + } + result := make([]imap.Address, 0, len(addresses)) + for _, addr := range addresses { + parts := strings.SplitN(addr.Address, "@", 2) + mailbox := parts[0] + host := "" + if len(parts) > 1 { + host = parts[1] + } + result = append(result, imap.Address{ + Name: addr.Name, + Mailbox: mailbox, + Host: host, + }) + } + return result +} + +// buildRawMessage reconstructs a raw RFC822 message from a db.Message. +// If RawData is available, it uses the original raw data directly. +func buildRawMessage(msg *db.Message) []byte { + if msg.RawData != "" { + return []byte(msg.RawData) + } + + var buf bytes.Buffer + buf.WriteString(fmt.Sprintf("From: %s\r\n", msg.FromAddr)) + buf.WriteString(fmt.Sprintf("To: %s\r\n", msg.ToAddr)) + if msg.CcAddr != "" { + buf.WriteString(fmt.Sprintf("Cc: %s\r\n", msg.CcAddr)) + } + buf.WriteString(fmt.Sprintf("Subject: %s\r\n", msg.Subject)) + buf.WriteString(fmt.Sprintf("Date: %s\r\n", msg.Date.Format(time.RFC1123Z))) + if msg.MessageID != "" { + buf.WriteString(fmt.Sprintf("Message-ID: %s\r\n", msg.MessageID)) + } + buf.WriteString("MIME-Version: 1.0\r\n") + + if msg.HtmlBody != "" && msg.TextBody != "" { + boundary := fmt.Sprintf("mailgo_%d", msg.ID) + buf.WriteString(fmt.Sprintf("Content-Type: multipart/alternative; boundary=\"%s\"\r\n", boundary)) + buf.WriteString("\r\n") + buf.WriteString(fmt.Sprintf("--%s\r\n", boundary)) + buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n\r\n") + buf.WriteString(msg.TextBody) + buf.WriteString("\r\n") + buf.WriteString(fmt.Sprintf("--%s\r\n", boundary)) + buf.WriteString("Content-Type: text/html; charset=utf-8\r\n\r\n") + buf.WriteString(msg.HtmlBody) + buf.WriteString("\r\n") + buf.WriteString(fmt.Sprintf("--%s--\r\n", boundary)) + } else if msg.TextBody != "" { + buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n\r\n") + buf.WriteString(msg.TextBody) + } else if msg.HtmlBody != "" { + buf.WriteString("Content-Type: text/html; charset=utf-8\r\n\r\n") + buf.WriteString(msg.HtmlBody) + } else { + buf.WriteString("Content-Type: text/plain; charset=utf-8\r\n\r\n") + } + + return buf.Bytes() +} + +// 编译期断言:imapSession 满足 Session 与 SessionMove 接口。 +var _ imapserver.Session = (*imapSession)(nil) +var _ imapserver.SessionMove = (*imapSession)(nil) diff --git a/internal/store/mail_store.go b/internal/store/mail_store.go index e0442b5..fd5f411 100644 --- a/internal/store/mail_store.go +++ b/internal/store/mail_store.go @@ -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