fix(imap): 修复 Thunderbird 未读状态丢失(SQLite 并发写 + 客户端 seq 错位)
A. SQLite 并发写失败被吞(高概率根因): - DSN 追加 _busy_timeout=5000&_journal_mode=WAL&_synchronous=NORMAL, WAL 下读不阻塞写,消除瞬时 SQLITE_BUSY - UpdateMessagesFlags 不再吞 MarkReadState/MarkFlagged 错误: 出错记日志并返回 error(客户端收到 NO 会重试) - Web MarkRead/Delete 同款静默丢错补日志 B. 客户端 seq 视图错位(集成测试复现后修复): - 真实服务器+脚本客户端测试复现:客户端按日期倒序自编号发 seq 式 STORE 时,旧排序(id ASC)把最旧邮件标为已读 - 规范排序改为 date DESC, id DESC(最新在前),与主流客户端 默认视图一致;buildNewMessageUpdate 改用 seqOf 取真实序号 - 新增 UID STORE / 服务器下发 seq STORE / 自编号 seq STORE 集成测试 注:go-imap v1.2.1 存在库内 *conn.silent() 数据竞争(启用推送后 必然触发),集成测试以 //go:build !race 排除,-race 下由单元测试 覆盖推送逻辑
This commit is contained in:
@@ -4,6 +4,7 @@ import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"mail_go/config"
|
||||
|
||||
@@ -31,6 +32,15 @@ func InitDB(cfg config.DatabaseConfig, storageCfg config.StorageConfig) (*gorm.D
|
||||
if err := os.MkdirAll(dir, 0755); err != nil {
|
||||
return nil, fmt.Errorf("创建数据库目录失败 %s: %w", dir, err)
|
||||
}
|
||||
// 多连接并发(SMTP/IMAP/POP3/Web/外发 worker/推送)下:
|
||||
// - WAL 模式:读不阻塞写,消除瞬时 SQLITE_BUSY 导致写失败被吞
|
||||
// - busy_timeout=5000ms:写竞争时等待而非立刻失败
|
||||
// - synchronous=NORMAL:WAL 下安全且写入更快
|
||||
sep := "?"
|
||||
if strings.Contains(dsn, "?") {
|
||||
sep = "&"
|
||||
}
|
||||
dsn = dsn + sep + "_busy_timeout=5000&_journal_mode=WAL&_synchronous=NORMAL"
|
||||
dialector = sqlite.Open(dsn)
|
||||
case "mysql":
|
||||
dialector = mysql.Open(cfg.DSN)
|
||||
|
||||
@@ -45,14 +45,14 @@ func (b *imapBackend) Updates() <-chan backend.Update {
|
||||
}
|
||||
|
||||
// buildNewMessageUpdate 为一条新投递到 mailbox 的邮件构造 IMAP 更新。
|
||||
// seq 取该邮件在邮箱中的序号(通知时机在入库之后,取当前列表长度)。
|
||||
// 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 := uint32(1)
|
||||
if msgs, err := stores.Mails.ListAllByUserAndFolder(msg.UserID, mailbox); err == nil {
|
||||
seq = uint32(len(msgs))
|
||||
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})
|
||||
@@ -751,6 +751,10 @@ func (m *imapMailbox) UpdateMessagesFlags(uid bool, seqset *imap.SeqSet, op imap
|
||||
return err
|
||||
}
|
||||
|
||||
// 记录首个持久化错误:SQLite 忙/锁等瞬时失败必须让客户端感知
|
||||
// (返回 NO 触发重试),否则已读/星标会静默丢失。
|
||||
var firstErr error
|
||||
|
||||
for i, dbMsg := range dbMessages {
|
||||
var match bool
|
||||
if uid {
|
||||
@@ -770,9 +774,15 @@ func (m *imapMailbox) UpdateMessagesFlags(uid bool, seqset *imap.SeqSet, op imap
|
||||
applyFlag := func(flag string, enabled bool) {
|
||||
switch flag {
|
||||
case "\\Seen":
|
||||
_ = m.stores.Mails.MarkReadState(dbMsg.ID, enabled)
|
||||
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
|
||||
}
|
||||
case "\\Flagged":
|
||||
_ = m.stores.Mails.MarkFlagged(dbMsg.ID, enabled)
|
||||
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
|
||||
}
|
||||
case "\\Deleted":
|
||||
if enabled {
|
||||
m.deleted[dbMsg.ID] = true
|
||||
@@ -807,7 +817,7 @@ func (m *imapMailbox) UpdateMessagesFlags(uid bool, seqset *imap.SeqSet, op imap
|
||||
pushUpdate(m.user.updates, buildFlagsUpdate(m.stores, m.user.email, m.name, fresh, deleted))
|
||||
}
|
||||
|
||||
return nil
|
||||
return firstErr
|
||||
}
|
||||
|
||||
// CopyMessages copies messages to another mailbox.
|
||||
|
||||
@@ -0,0 +1,190 @@
|
||||
//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)承担。
|
||||
package imap_server
|
||||
|
||||
import (
|
||||
"net"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"mail_go/config"
|
||||
"mail_go/internal/connhub"
|
||||
"mail_go/internal/db"
|
||||
"mail_go/internal/store"
|
||||
|
||||
"github.com/emersion/go-imap"
|
||||
"github.com/emersion/go-imap/client"
|
||||
"golang.org/x/crypto/bcrypt"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// startIntegrationServer 启动一个真实的 IMAP 监听(随机端口)供客户端测试。
|
||||
func startIntegrationServer(t *testing.T) (*store.Stores, string) {
|
||||
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{}, &db.ProtocolLog{}, &db.BanEntry{}); err != nil {
|
||||
t.Fatalf("migrate: %v", err)
|
||||
}
|
||||
stores := store.NewStores(gdb)
|
||||
|
||||
domain := &db.Domain{Name: "example.com"}
|
||||
if err := stores.Domains.Create(domain); err != nil {
|
||||
t.Fatalf("create domain: %v", err)
|
||||
}
|
||||
hashed, _ := bcrypt.GenerateFromPassword([]byte("secret123"), bcrypt.DefaultCost)
|
||||
user := &db.User{Username: "alice", DomainID: domain.ID, PasswordHash: string(hashed), IsActive: true}
|
||||
if err := stores.Users.Create(user); err != nil {
|
||||
t.Fatalf("create user: %v", err)
|
||||
}
|
||||
|
||||
srv := NewIMAPServer(config.IMAPConfig{}, stores, nil, config.BanConfig{}, connhub.New())
|
||||
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatalf("listen: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { ln.Close() })
|
||||
|
||||
imapSrv := srv.newServer(ln.Addr().String(), nil)
|
||||
go imapSrv.Serve(ln)
|
||||
|
||||
return stores, ln.Addr().String()
|
||||
}
|
||||
|
||||
// seedMailbox 创建 n 封按时间递增的邮件(id 与 date 顺序一致时 id ASC == date ASC)。
|
||||
func seedMailbox(t *testing.T, stores *store.Stores, userID uint, n int) []uint {
|
||||
t.Helper()
|
||||
ids := make([]uint, 0, n)
|
||||
base := time.Now().Add(-time.Duration(n) * time.Hour)
|
||||
for i := 0; i < n; i++ {
|
||||
msg := &db.Message{
|
||||
UserID: userID,
|
||||
Folder: "INBOX",
|
||||
FromAddr: "x@y",
|
||||
ToAddr: "alice@example.com",
|
||||
Subject: "m",
|
||||
Date: base.Add(time.Duration(i) * time.Hour), // 时间递增:id 越大日期越新
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
if err := stores.Mails.Create(msg); err != nil {
|
||||
t.Fatalf("create message: %v", err)
|
||||
}
|
||||
ids = append(ids, msg.ID)
|
||||
}
|
||||
return ids
|
||||
}
|
||||
|
||||
func loginAndSelect(t *testing.T, addr string) *client.Client {
|
||||
t.Helper()
|
||||
c, err := client.Dial(addr)
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { c.Logout() })
|
||||
if err := c.Login("alice@example.com", "secret123"); err != nil {
|
||||
t.Fatalf("login: %v", err)
|
||||
}
|
||||
if _, err := c.Select("INBOX", false); 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)
|
||||
if err != nil {
|
||||
t.Fatalf("get msg: %v", err)
|
||||
}
|
||||
if msg.IsRead != want {
|
||||
t.Fatalf("msg %d IsRead = %v, want %v", msgID, msg.IsRead, want)
|
||||
}
|
||||
}
|
||||
|
||||
// TestUidStorePersists 验证 UID STORE +FLAGS(\Seen) 持久化(RFC 标准流程)。
|
||||
func TestUidStorePersists(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(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 {
|
||||
t.Fatalf("uid store: %v", err)
|
||||
}
|
||||
<-ch
|
||||
|
||||
assertReadState(t, stores, ids[1], true)
|
||||
assertReadState(t, stores, ids[0], false)
|
||||
assertReadState(t, stores, ids[2], false)
|
||||
}
|
||||
|
||||
// TestSeqStoreServerIssued 验证客户端用服务器下发的序号(FETCH 结果)做
|
||||
// seq 式 STORE:任何排序下都应正确持久化。
|
||||
func TestSeqStoreServerIssued(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(t)
|
||||
ids := seedMailbox(t, stores, 1, 3)
|
||||
|
||||
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 {
|
||||
t.Fatalf("fetch: %v", err)
|
||||
}
|
||||
var targetSeq uint32
|
||||
for m := range messages {
|
||||
if m.Uid == uint32(ids[2]) {
|
||||
targetSeq = m.SeqNum
|
||||
}
|
||||
}
|
||||
if targetSeq == 0 {
|
||||
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 {
|
||||
t.Fatalf("store: %v", err)
|
||||
}
|
||||
<-ch
|
||||
|
||||
assertReadState(t, stores, ids[2], true)
|
||||
}
|
||||
|
||||
// TestSeqStoreClientSelfNumbered 复现风险场景:客户端不信任服务器序号,
|
||||
// 按自己的视图(日期倒序,最新在前)自行编号后发 seq 式 STORE。
|
||||
// 服务器规范排序必须与常见客户端视图一致(date DESC, id DESC),
|
||||
// 否则会把另一封邮件标为已读、目标邮件永远未读。
|
||||
func TestSeqStoreClientSelfNumbered(t *testing.T) {
|
||||
stores, addr := startIntegrationServer(t)
|
||||
ids := seedMailbox(t, stores, 1, 3) // 3 封,日期递增,最新的是 ids[2]
|
||||
|
||||
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 {
|
||||
t.Fatalf("store: %v", err)
|
||||
}
|
||||
<-ch
|
||||
|
||||
// 客户端意图是标记最新一封(ids[2])为已读
|
||||
assertReadState(t, stores, ids[2], true)
|
||||
}
|
||||
@@ -36,8 +36,8 @@ func TestPushNewMessage(t *testing.T) {
|
||||
}
|
||||
email := "alice@example.com"
|
||||
|
||||
// 已有一封旧邮件,新邮件应为 INBOX 第 2 封
|
||||
old := &db.Message{UserID: user.ID, Folder: "INBOX", FromAddr: "x@y", Subject: "old", Date: time.Now()}
|
||||
// 已有一封旧邮件(日期更早);规范排序最新在前,新邮件应为 INBOX 第 1 封
|
||||
old := &db.Message{UserID: user.ID, 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)
|
||||
}
|
||||
@@ -87,8 +87,8 @@ func TestPushNewMessage(t *testing.T) {
|
||||
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 != 2 {
|
||||
t.Fatalf("backend %d: seq = %d, want 2", i, mu.Message.SeqNum)
|
||||
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)
|
||||
|
||||
@@ -118,11 +118,14 @@ func (s *mailStoreGorm) CountUnread(userID uint, folder string) (int64, error) {
|
||||
}
|
||||
|
||||
// ListAllByUserAndFolder retrieves all messages for a user in a folder without pagination.
|
||||
// Messages are ordered by ID ascending so that sequence numbers are stable.
|
||||
// 按 date DESC, id DESC 排序(最新在前):与主流邮件客户端(Thunderbird、
|
||||
// 手机客户端等)默认视图一致,客户端自行按日期编号的 seq 式 STORE 不会
|
||||
// 错位标错邮件。所有 IMAP 序号相关路径(Status/ListMessages/推送/seqOf)
|
||||
// 共用本排序,保证序号全链路一致。
|
||||
func (s *mailStoreGorm) ListAllByUserAndFolder(userID uint, folder string) ([]db.Message, error) {
|
||||
var messages []db.Message
|
||||
if err := s.db.Where("user_id = ? AND folder = ?", userID, folder).
|
||||
Order("id ASC").Find(&messages).Error; err != nil {
|
||||
Order("date DESC, id DESC").Find(&messages).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return messages, nil
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"encoding/base64"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"mime"
|
||||
"net/http"
|
||||
"path/filepath"
|
||||
@@ -132,7 +133,9 @@ func (h *MailHandler) View(c *gin.Context) {
|
||||
|
||||
// Auto mark as read
|
||||
if !msg.IsRead {
|
||||
_ = h.stores.Mails.MarkRead(uint(id))
|
||||
if err := h.stores.Mails.MarkRead(uint(id)); err != nil {
|
||||
log.Printf("web: 标记已读失败 msg=%d: %v", id, err)
|
||||
}
|
||||
msg.IsRead = true
|
||||
}
|
||||
|
||||
@@ -637,7 +640,9 @@ func (h *MailHandler) Delete(c *gin.Context) {
|
||||
_ = h.storage.Delete(att.FilePath)
|
||||
_ = h.stores.Users.UpdateUsedBytes(userID, -att.FileSize)
|
||||
}
|
||||
_ = h.stores.Attachments.DeleteByMessage(uint(id))
|
||||
if err := h.stores.Attachments.DeleteByMessage(uint(id)); err != nil {
|
||||
log.Printf("web: 删除附件记录失败 msg=%d: %v", id, err)
|
||||
}
|
||||
|
||||
// 删除前计算消息在所属文件夹中的序号(用于 Expunge 推送)
|
||||
var seq uint32
|
||||
@@ -649,7 +654,9 @@ func (h *MailHandler) Delete(c *gin.Context) {
|
||||
}
|
||||
}
|
||||
}
|
||||
_ = h.stores.Mails.Delete(uint(id))
|
||||
if err := h.stores.Mails.Delete(uint(id)); err != nil {
|
||||
log.Printf("web: 删除邮件失败 msg=%d: %v", id, err)
|
||||
}
|
||||
|
||||
// 删除 → 推送给该用户的其他 IMAP 客户端
|
||||
if h.pusher != nil && seq > 0 {
|
||||
|
||||
Reference in New Issue
Block a user