fix(send_queue): 修复已投递标记恢复逻辑
- 新增投递标记文件持久化,避免重启后重复发送已投递的队列项。 - 在恢复卡住的发送任务时,识别已存在的投递标记并直接修复数据库状态。 - 同步补充回归测试,覆盖陈旧投递标记不应触发重发的场景。
This commit is contained in:
@@ -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) {
|
func TestSubmissionAuthRequiresMailboxPasswordAndSendPermission(t *testing.T) {
|
||||||
a := newTestApp(t)
|
a := newTestApp(t)
|
||||||
user, mailbox, err := a.authenticateSubmission(context.Background(), "admin@lanqin.local", "ChangeMe123!")
|
user, mailbox, err := a.authenticateSubmission(context.Background(), "admin@lanqin.local", "ChangeMe123!")
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ import (
|
|||||||
"database/sql"
|
"database/sql"
|
||||||
"encoding/base64"
|
"encoding/base64"
|
||||||
"errors"
|
"errors"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -26,6 +28,8 @@ const (
|
|||||||
|
|
||||||
sendQueueStaleAfter = 15 * time.Minute
|
sendQueueStaleAfter = 15 * time.Minute
|
||||||
sendQueueConcurrency = 4
|
sendQueueConcurrency = 4
|
||||||
|
|
||||||
|
sendQueueDeliveredMarkerDir = "send_queue_delivered"
|
||||||
)
|
)
|
||||||
|
|
||||||
type sendQueueInput struct {
|
type sendQueueInput struct {
|
||||||
@@ -87,6 +91,7 @@ func (a *App) enqueueSend(ctx context.Context, in sendQueueInput) (string, error
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
a.deleteSendQueueDeliveredMarker(existingID)
|
||||||
a.recordSendAudit(ctx, sendAuditQueued, sendQueueStatusQueued, sendAuditInput{
|
a.recordSendAudit(ctx, sendAuditQueued, sendQueueStatusQueued, sendAuditInput{
|
||||||
QueueID: existingID,
|
QueueID: existingID,
|
||||||
UserID: in.UserID,
|
UserID: in.UserID,
|
||||||
@@ -195,6 +200,14 @@ func (a *App) recoverStaleSendQueueItems(ctx context.Context) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
item.Recipients = jsonDecodeSlice(recipientsJSON)
|
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)
|
mimeBytes, err := base64.StdEncoding.DecodeString(mimeBase64)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -207,6 +220,16 @@ func (a *App) recoverStaleSendQueueItems(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
now := a.now().UTC().Format(time.RFC3339Nano)
|
now := a.now().UTC().Format(time.RFC3339Nano)
|
||||||
for _, item := range items {
|
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)
|
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 {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -230,12 +253,23 @@ func (a *App) processSendQueueItem(ctx context.Context, id string) {
|
|||||||
a.markSendQueueFailed(ctx, item, err)
|
a.markSendQueueFailed(ctx, item, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
now := a.now().UTC().Format(time.RFC3339Nano)
|
if err := a.writeSendQueueDeliveredMarker(item.ID); err != nil {
|
||||||
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 {
|
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)
|
a.log.Warn("failed to mark send queue delivered", "id", item.ID, "error", err)
|
||||||
return
|
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, ""))
|
a.recordSendAudit(ctx, sendAuditDelivered, sendQueueStatusDelivered, sendAuditInputFromQueue(item, ""))
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (a *App) claimSendQueueItem(ctx context.Context, id string) (sendQueueItem, error) {
|
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) {
|
func (a *App) authorizedSender(ctx context.Context, mb *Mailbox, from string) (string, string, error) {
|
||||||
from = normalizeEmail(from)
|
from = normalizeEmail(from)
|
||||||
if from == "" {
|
if from == "" {
|
||||||
|
|||||||
Reference in New Issue
Block a user