From 65e4c4f6b5011d238ac5f523f36eaabba88f49ef Mon Sep 17 00:00:00 2001 From: LanQin_ Date: Wed, 24 Jun 2026 13:35:37 +0800 Subject: [PATCH] =?UTF-8?q?fix(send=5Fqueue):=20=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?=E5=B7=B2=E6=8A=95=E9=80=92=E6=A0=87=E8=AE=B0=E6=81=A2=E5=A4=8D?= =?UTF-8?q?=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增投递标记文件持久化,避免重启后重复发送已投递的队列项。 - 在恢复卡住的发送任务时,识别已存在的投递标记并直接修复数据库状态。 - 同步补充回归测试,覆盖陈旧投递标记不应触发重发的场景。 --- apps/api/internal/app/app_test.go | 55 ++++++++++++++++++++ apps/api/internal/app/send_queue.go | 81 ++++++++++++++++++++++++++++- 2 files changed, 134 insertions(+), 2 deletions(-) diff --git a/apps/api/internal/app/app_test.go b/apps/api/internal/app/app_test.go index 78587c7..df9eba1 100644 --- a/apps/api/internal/app/app_test.go +++ b/apps/api/internal/app/app_test.go @@ -975,6 +975,61 @@ func TestSendQueueRecoversStaleSendingItems(t *testing.T) { } } +func TestSendQueueStaleDeliveredMarkerDoesNotRedeliver(t *testing.T) { + a := newTestApp(t) + host, port, received := startCapturingSMTP(t, 1) + a.cfg.SMTPHost = host + a.cfg.SMTPPort = port + user, mb := defaultAdminUserAndMailbox(t, a) + now := a.now().UTC() + mimeBytes := []byte("From: admin@lanqin.local\r\nTo: person@example.com\r\nSubject: marker\r\n\r\nbody") + queueID, err := a.enqueueSend(context.Background(), sendQueueInput{ + UserID: user.ID, + MailboxID: mb.ID, + Source: sendSourceWebmail, + MailFrom: mb.Address, + HeaderFrom: mb.Address, + Recipients: []string{"person@example.com"}, + MIMEBytes: mimeBytes, + Now: now, + }) + if err != nil { + t.Fatal(err) + } + staleAt := now.Add(-sendQueueStaleAfter - time.Minute).Format(time.RFC3339Nano) + if _, err := a.db.Exec(`UPDATE send_queue SET status=?,attempt_count=1,updated_at=? WHERE id=?`, sendQueueStatusSending, staleAt, queueID); err != nil { + t.Fatal(err) + } + if err := a.writeSendQueueDeliveredMarker(queueID); err != nil { + t.Fatal(err) + } + if err := a.processDueSendQueue(context.Background()); err != nil { + t.Fatal(err) + } + select { + case body := <-received: + t.Fatalf("stale delivered marker should not redeliver, got %q", body) + case <-time.After(200 * time.Millisecond): + } + var status, mimeBase64 string + if err := a.db.QueryRow(`SELECT status,mime_base64 FROM send_queue WHERE id=?`, queueID).Scan(&status, &mimeBase64); err != nil { + t.Fatal(err) + } + if status != sendQueueStatusDelivered { + t.Fatalf("queue status=%q, want delivered", status) + } + if mimeBase64 != "" { + t.Fatal("delivered marker recovery should clear raw MIME") + } + delivered, err := a.hasSendQueueDeliveredMarker(queueID) + if err != nil { + t.Fatal(err) + } + if delivered { + t.Fatal("delivered marker should be removed after database state is repaired") + } +} + func TestSubmissionAuthRequiresMailboxPasswordAndSendPermission(t *testing.T) { a := newTestApp(t) user, mailbox, err := a.authenticateSubmission(context.Background(), "admin@lanqin.local", "ChangeMe123!") diff --git a/apps/api/internal/app/send_queue.go b/apps/api/internal/app/send_queue.go index 91d911a..02a03ec 100644 --- a/apps/api/internal/app/send_queue.go +++ b/apps/api/internal/app/send_queue.go @@ -5,6 +5,8 @@ import ( "database/sql" "encoding/base64" "errors" + "os" + "path/filepath" "strings" "time" ) @@ -26,6 +28,8 @@ const ( sendQueueStaleAfter = 15 * time.Minute sendQueueConcurrency = 4 + + sendQueueDeliveredMarkerDir = "send_queue_delivered" ) type sendQueueInput struct { @@ -87,6 +91,7 @@ func (a *App) enqueueSend(ctx context.Context, in sendQueueInput) (string, error if err != nil { return "", err } + a.deleteSendQueueDeliveredMarker(existingID) a.recordSendAudit(ctx, sendAuditQueued, sendQueueStatusQueued, sendAuditInput{ QueueID: existingID, UserID: in.UserID, @@ -195,6 +200,14 @@ func (a *App) recoverStaleSendQueueItems(ctx context.Context) error { return err } item.Recipients = jsonDecodeSlice(recipientsJSON) + delivered, err := a.hasSendQueueDeliveredMarker(item.ID) + if err != nil { + return err + } + if delivered { + items = append(items, item) + continue + } mimeBytes, err := base64.StdEncoding.DecodeString(mimeBase64) if err != nil { return err @@ -207,6 +220,16 @@ func (a *App) recoverStaleSendQueueItems(ctx context.Context) error { } now := a.now().UTC().Format(time.RFC3339Nano) for _, item := range items { + delivered, err := a.hasSendQueueDeliveredMarker(item.ID) + if err != nil { + return err + } + if delivered { + if err := a.markSendQueueDelivered(ctx, item); err != nil { + return err + } + continue + } res, err := a.db.ExecContext(ctx, `UPDATE send_queue SET status=?,next_attempt_at=?,last_error=?,updated_at=? WHERE id=? AND status=?`, sendQueueStatusFailed, now, "send attempt interrupted", now, item.ID, sendQueueStatusSending) if err != nil { return err @@ -230,12 +253,23 @@ func (a *App) processSendQueueItem(ctx context.Context, id string) { a.markSendQueueFailed(ctx, item, err) return } - now := a.now().UTC().Format(time.RFC3339Nano) - if _, err := a.db.ExecContext(ctx, `UPDATE send_queue SET status=?,delivered_at=?,updated_at=?,last_error='',mime_base64='' WHERE id=?`, sendQueueStatusDelivered, now, now, item.ID); err != nil { + if err := a.writeSendQueueDeliveredMarker(item.ID); err != nil { + a.log.Warn("failed to persist send queue delivered marker", "id", item.ID, "error", err) + } + if err := a.markSendQueueDelivered(ctx, item); err != nil { a.log.Warn("failed to mark send queue delivered", "id", item.ID, "error", err) return } +} + +func (a *App) markSendQueueDelivered(ctx context.Context, item sendQueueItem) error { + now := a.now().UTC().Format(time.RFC3339Nano) + if _, err := a.db.ExecContext(ctx, `UPDATE send_queue SET status=?,delivered_at=?,updated_at=?,last_error='',mime_base64='' WHERE id=?`, sendQueueStatusDelivered, now, now, item.ID); err != nil { + return err + } + a.deleteSendQueueDeliveredMarker(item.ID) a.recordSendAudit(ctx, sendAuditDelivered, sendQueueStatusDelivered, sendAuditInputFromQueue(item, "")) + return nil } func (a *App) claimSendQueueItem(ctx context.Context, id string) (sendQueueItem, error) { @@ -329,6 +363,49 @@ func (a *App) recordSendAudit(ctx context.Context, event, status string, in send } } +func (a *App) sendQueueDeliveredMarkerPath(id string) string { + safeID := filepath.Base(strings.TrimSpace(id)) + if safeID == "" || safeID == "." { + safeID = "unknown" + } + return filepath.Join(a.cfg.DataDir, sendQueueDeliveredMarkerDir, safeID+".marker") +} + +func (a *App) writeSendQueueDeliveredMarker(id string) error { + path := a.sendQueueDeliveredMarkerPath(id) + dir := filepath.Dir(path) + if err := os.MkdirAll(dir, 0o700); err != nil { + return err + } + tmp := filepath.Join(dir, filepath.Base(path)+"."+newID("tmp")) + if err := os.WriteFile(tmp, []byte(a.now().UTC().Format(time.RFC3339Nano)), 0o600); err != nil { + return err + } + if err := os.Rename(tmp, path); err != nil { + _ = os.Remove(tmp) + return err + } + return nil +} + +func (a *App) hasSendQueueDeliveredMarker(id string) (bool, error) { + _, err := os.Stat(a.sendQueueDeliveredMarkerPath(id)) + if err == nil { + return true, nil + } + if errors.Is(err, os.ErrNotExist) { + return false, nil + } + return false, err +} + +func (a *App) deleteSendQueueDeliveredMarker(id string) { + err := os.Remove(a.sendQueueDeliveredMarkerPath(id)) + if err != nil && !errors.Is(err, os.ErrNotExist) { + a.log.Warn("failed to remove send queue delivered marker", "id", id, "error", err) + } +} + func (a *App) authorizedSender(ctx context.Context, mb *Mailbox, from string) (string, string, error) { from = normalizeEmail(from) if from == "" {