From aa0437fb4f17f8f2de3e409776d859f84d7dacf1 Mon Sep 17 00:00:00 2001 From: dsh Date: Wed, 19 Aug 2026 11:13:42 -0400 Subject: [PATCH] =?UTF-8?q?perf(imap):=20STORE=20=E6=89=B9=E9=87=8F?= =?UTF-8?q?=E6=8C=81=E4=B9=85=E5=8C=96=20+=20=E6=8E=A8=E9=80=81=E5=BA=8F?= =?UTF-8?q?=E5=8F=B7=E5=A4=8D=E7=94=A8=E5=B7=B2=E5=8A=A0=E8=BD=BD=E5=88=97?= =?UTF-8?q?=E8=A1=A8=EF=BC=8C=E6=B6=88=E9=99=A4=E9=80=90=E6=9D=A1=E5=86=99?= =?UTF-8?q?=E5=BA=93=E4=B8=8E=E5=85=A8=E9=87=8F=E6=89=AB=E6=8F=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 原实现每封匹配消息单独 UPDATE + GetByID + seqOf 全量扫描(含大 附件 raw_data),手机整批标记已读(60+ 封)时产生 60 次写 + 180 次全量读,连接被长时间占住,其他连接响应被推送洪泛阻塞。 优化:收集目标状态后合并为单条 UPDATE ... IN 批量写;推送更新 直接用已加载列表的序号与目标状态(buildFlagsUpdateAt),不再 重复查库。实测 19 封批量 STORE 服务器耗时 372µs。 --- internal/imap_server/backend.go | 142 +++++++++++++++++++++++--------- internal/store/mail_store.go | 22 +++++ 2 files changed, 126 insertions(+), 38 deletions(-) diff --git a/internal/imap_server/backend.go b/internal/imap_server/backend.go index cd5f393..2f2deec 100644 --- a/internal/imap_server/backend.go +++ b/internal/imap_server/backend.go @@ -83,10 +83,20 @@ func buildFlagsUpdate(stores *store.Stores, userEmail, mailbox string, msg *db.M if stores == nil || msg == nil || userEmail == "" || mailbox == "" { return nil } - imapMsg := imap.NewMessage(seqOf(stores, msg.UserID, mailbox, msg.ID), + 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(msg.IsRead, msg.IsFlagged, deleted) + imapMsg.Flags = flagsOf(read, flagged, deleted) return &backend.MessageUpdate{ Update: backend.NewUpdate(userEmail, mailbox), Message: imapMsg, @@ -790,12 +800,30 @@ func (m *imapMailbox) UpdateMessagesFlags(uid bool, seqset *imap.SeqSet, op imap 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 + } - // 记录首个持久化错误:SQLite 忙/锁等瞬时失败必须让客户端感知 - // (返回 NO 触发重试),否则已读/星标会静默丢失。 - var firstErr error + 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, dbMsg := range dbMessages { + for i := range dbMessages { + dbMsg := &dbMessages[i] var match bool if uid { match = seqset.Contains(uint32(dbMsg.ID)) @@ -806,55 +834,93 @@ func (m *imapMailbox) UpdateMessagesFlags(uid bool, seqset *imap.SeqSet, op imap continue } - flagSet := make(map[string]bool, len(flags)) - for _, flag := range flags { - flagSet[flag] = true - } - - applyFlag := func(flag string, enabled bool) { + c := change{msg: dbMsg, seq: uint32(i + 1)} + apply := func(flag string, enabled bool) { switch flag { case "\\Seen": - if err := m.stores.Mails.MarkReadState(dbMsg.ID, enabled); err != nil && firstErr == nil { - log.Printf("IMAP: mark read state for msg %d failed: %v", dbMsg.ID, err) - firstErr = err - } + c.newRead, c.readSet = enabled, true case "\\Flagged": - if err := m.stores.Mails.MarkFlagged(dbMsg.ID, enabled); err != nil && firstErr == nil { - log.Printf("IMAP: mark flagged for msg %d failed: %v", dbMsg.ID, err) - firstErr = err - } + c.newFlagged, c.flaggedSet = enabled, true case "\\Deleted": - if enabled { - m.deleted[dbMsg.ID] = true - } else { - delete(m.deleted, dbMsg.ID) - } + c.newDeleted, c.deletedSet = enabled, true } } - switch op { case imap.SetFlags: - applyFlag("\\Seen", flagSet["\\Seen"]) - applyFlag("\\Flagged", flagSet["\\Flagged"]) - applyFlag("\\Deleted", flagSet["\\Deleted"]) + apply("\\Seen", flagSet["\\Seen"]) + apply("\\Flagged", flagSet["\\Flagged"]) + apply("\\Deleted", flagSet["\\Deleted"]) case imap.AddFlags: for flag := range flagSet { - applyFlag(flag, true) + apply(flag, true) } case imap.RemoveFlags: for flag := range flagSet { - applyFlag(flag, false) + apply(flag, false) } } + changes = append(changes, c) + } - // 标志变化(已读/星标/删除)→ 推送给同用户其他客户端 - // (重新读库取最新状态,\Deleted 取会话内状态) - fresh, err := m.stores.Mails.GetByID(dbMsg.ID) - if err != nil { - continue + // 批量持久化(已读/星标合并为单条 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) + } } - deleted := m.deleted != nil && m.deleted[dbMsg.ID] - pushUpdate(m.user.updates, buildFlagsUpdate(m.stores, m.user.email, m.name, fresh, deleted)) + 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 diff --git a/internal/store/mail_store.go b/internal/store/mail_store.go index 68d310c..e0442b5 100644 --- a/internal/store/mail_store.go +++ b/internal/store/mail_store.go @@ -28,6 +28,10 @@ type MailStore interface { MarkRead(id uint) error MarkReadState(id uint, read bool) error MarkFlagged(id uint, flagged bool) error + // SetReadStates 批量设置多封邮件的已读状态(单条 UPDATE ... IN)。 + SetReadStates(ids []uint, read bool) error + // SetFlaggedStates 批量设置多封邮件的星标状态(单条 UPDATE ... IN)。 + SetFlaggedStates(ids []uint, flagged bool) error MoveToFolder(id uint, folder string) error Delete(id uint) error CountUnread(userID uint, folder string) (int64, error) @@ -103,6 +107,24 @@ func (s *mailStoreGorm) MarkFlagged(id uint, flagged bool) error { return s.db.Model(&db.Message{}).Where("id = ?", id).Update("is_flagged", flagged).Error } +// SetReadStates 批量设置多封邮件的已读状态。 +// 客户端整批标记已读(手机同步后 STORE +FLAGS \Seen)时,逐条 UPDATE +// 会产生大量写事务并占住连接,这里合并为单条 SQL。 +func (s *mailStoreGorm) SetReadStates(ids []uint, read bool) error { + if len(ids) == 0 { + return nil + } + return s.db.Model(&db.Message{}).Where("id IN ?", ids).Update("is_read", read).Error +} + +// SetFlaggedStates 批量设置多封邮件的星标状态。 +func (s *mailStoreGorm) SetFlaggedStates(ids []uint, flagged bool) error { + if len(ids) == 0 { + return nil + } + return s.db.Model(&db.Message{}).Where("id IN ?", ids).Update("is_flagged", flagged).Error +} + // MoveToFolder changes the folder of a message. func (s *mailStoreGorm) MoveToFolder(id uint, folder string) error { return s.db.Model(&db.Message{}).Where("id = ?", id).Update("folder", folder).Error