feat(app): 引入发送队列与可授权发件人

- Webmail/API 与第三方提交统一先写入已发送,再入发送队列异步投递。
- 新增队列重试、失败审计、陈旧任务恢复与发送状态追踪。
- 支持本人邮箱、启用别名和 `send_as_grants` 显式授权的 send-as。
- 更新 README 与部署文档,说明新的外发流程和排障方式。
This commit is contained in:
LanQin_
2026-06-24 11:27:32 +08:00
parent 5c9bea075a
commit 7b0717412c
7 changed files with 646 additions and 46 deletions
+194 -11
View File
@@ -46,6 +46,20 @@ func newTestApp(t *testing.T) *App {
return a
}
func defaultAdminUserAndMailbox(t *testing.T, a *App) (*User, *Mailbox) {
t.Helper()
ctx := context.Background()
user, _, err := a.userByEmail(ctx, "admin@lanqin.local")
if err != nil {
t.Fatal(err)
}
mb, err := a.mailboxByAddress(ctx, "admin@lanqin.local")
if err != nil {
t.Fatal(err)
}
return user, mb
}
func startFakeSMTP(t *testing.T) (string, string, <-chan string) {
t.Helper()
ln, err := net.Listen("tcp", "127.0.0.1:0")
@@ -775,7 +789,7 @@ func TestCatchAllStoresUnregisteredMailForAdminOnly(t *testing.T) {
}
}
func TestMailSendReturnsSMTPFailure(t *testing.T) {
func TestMailSendQueuesSMTPFailureForRetry(t *testing.T) {
a := newTestApp(t)
a.cfg.SMTPHost = "127.0.0.1"
a.cfg.SMTPPort = "1"
@@ -792,12 +806,93 @@ func TestMailSendReturnsSMTPFailure(t *testing.T) {
"subject": "smtp failure should surface",
"text": "hello",
}
var errBody map[string]any
if code := admin.do("POST", "/api/mail/send", payload, &errBody); code != http.StatusBadGateway {
t.Fatalf("smtp failure code=%d body=%v", code, errBody)
var sent MailMessage
if code := admin.do("POST", "/api/mail/send", payload, &sent); code != http.StatusCreated {
t.Fatalf("smtp queued send code=%d body=%+v", code, sent)
}
if got, _ := errBody["error"].(string); !strings.Contains(got, "smtp delivery failed") {
t.Fatalf("error=%q", got)
if err := a.processDueSendQueue(context.Background()); err != nil {
t.Fatal(err)
}
var status, lastError string
if err := a.db.QueryRow(`SELECT status,last_error FROM send_queue WHERE sent_message_id=?`, sent.ID).Scan(&status, &lastError); err != nil {
t.Fatal(err)
}
if status != sendQueueStatusFailed || lastError == "" {
t.Fatalf("queue status=%q lastError=%q", status, lastError)
}
var auditCount int
if err := a.db.QueryRow(`SELECT COUNT(1) FROM send_audit_events WHERE sent_message_id=? AND event=?`, sent.ID, sendAuditRetry).Scan(&auditCount); err != nil {
t.Fatal(err)
}
if auditCount != 1 {
t.Fatalf("retry audit count=%d, want 1", auditCount)
}
}
func TestMailSendRollsBackSentCopyWhenQueueInsertFails(t *testing.T) {
a := newTestApp(t)
a.cfg.SMTPHost = "postfix"
a.cfg.SMTPPort = "25"
user, mb := defaultAdminUserAndMailbox(t, a)
if _, err := a.db.ExecContext(context.Background(), `DROP TABLE send_queue`); err != nil {
t.Fatal(err)
}
_, err := a.sendMailNow(context.Background(), user, mb, mailComposeInput{
To: []string{"person@example.com"},
Subject: "queue insert failure",
Text: "hello",
})
if err == nil || !strings.Contains(err.Error(), "failed to enqueue delivery") {
t.Fatalf("sendMailNow error=%v, want enqueue failure", err)
}
var count int
if err := a.db.QueryRow(`SELECT COUNT(1) FROM messages WHERE mailbox_id=? AND subject=?`, mb.ID, "queue insert failure").Scan(&count); err != nil {
t.Fatal(err)
}
if count != 0 {
t.Fatalf("sent copy should be removed after enqueue failure, count=%d", count)
}
}
func TestSendQueueRecoversStaleSendingItems(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: stale\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.processDueSendQueue(context.Background()); err != nil {
t.Fatal(err)
}
select {
case <-received:
case <-time.After(2 * time.Second):
t.Fatal("recovered queue item was not relayed")
}
var status string
if err := a.db.QueryRow(`SELECT status FROM send_queue WHERE id=?`, queueID).Scan(&status); err != nil {
t.Fatal(err)
}
if status != sendQueueStatusDelivered {
t.Fatalf("queue status=%q, want delivered", status)
}
}
@@ -866,6 +961,9 @@ func TestSubmissionSendsRelayAndStoresSentCopy(t *testing.T) {
if err := a.submitSMTPMessage(context.Background(), user, mb, mb.Address, []string{"person@example.com", "hidden@example.com"}, strings.NewReader(raw)); err != nil {
t.Fatalf("submit smtp message: %v", err)
}
if err := a.processDueSendQueue(context.Background()); err != nil {
t.Fatal(err)
}
select {
case body := <-received:
if strings.Contains(strings.ToLower(body), "\r\nbcc:") || strings.Contains(body, "hidden@example.com") {
@@ -889,6 +987,13 @@ func TestSubmissionSendsRelayAndStoresSentCopy(t *testing.T) {
if got := jsonDecodeSlice(bccJSON); len(got) != 1 || got[0] != "hidden@example.com" {
t.Fatalf("bcc json=%s", bccJSON)
}
var deliveredAudits int
if err := a.db.QueryRow(`SELECT COUNT(1) FROM send_audit_events WHERE event=? AND status=?`, sendAuditDelivered, sendQueueStatusDelivered).Scan(&deliveredAudits); err != nil {
t.Fatal(err)
}
if deliveredAudits != 1 {
t.Fatalf("delivered audit count=%d, want 1", deliveredAudits)
}
}
func TestSubmissionRejectsMismatchedSender(t *testing.T) {
@@ -911,7 +1016,7 @@ func TestSubmissionRejectsMismatchedSender(t *testing.T) {
}
}
func TestSubmissionRelayFailureRemovesSentCopy(t *testing.T) {
func TestSubmissionRelayFailureKeepsSentCopyAndRetries(t *testing.T) {
a := newTestApp(t)
a.cfg.SMTPHost = "127.0.0.1"
a.cfg.SMTPPort = "1"
@@ -920,15 +1025,25 @@ func TestSubmissionRelayFailureRemovesSentCopy(t *testing.T) {
t.Fatal(err)
}
raw := "From: admin@lanqin.local\r\nTo: person@example.com\r\nSubject: relay fail\r\nMessage-ID: <relay-fail@example.test>\r\n\r\nbody"
if err := a.submitSMTPMessage(context.Background(), user, mb, mb.Address, []string{"person@example.com"}, strings.NewReader(raw)); err == nil {
t.Fatal("relay failure should fail")
if err := a.submitSMTPMessage(context.Background(), user, mb, mb.Address, []string{"person@example.com"}, strings.NewReader(raw)); err != nil {
t.Fatalf("submission should queue relay failure for retry: %v", err)
}
if err := a.processDueSendQueue(context.Background()); err != nil {
t.Fatal(err)
}
var count int
if err := a.db.QueryRow(`SELECT COUNT(1) FROM messages WHERE mailbox_id=? AND message_id=?`, mb.ID, "<relay-fail@example.test>").Scan(&count); err != nil {
t.Fatal(err)
}
if count != 0 {
t.Fatalf("sent copy should be removed after relay failure, count=%d", count)
if count != 1 {
t.Fatalf("sent copy should remain after queued relay failure, count=%d", count)
}
var status, lastError string
if err := a.db.QueryRow(`SELECT status,last_error FROM send_queue WHERE mailbox_id=? AND sent_message_id <> ''`, mb.ID).Scan(&status, &lastError); err != nil {
t.Fatal(err)
}
if status != sendQueueStatusFailed || lastError == "" {
t.Fatalf("queue status=%q lastError=%q", status, lastError)
}
}
@@ -958,6 +1073,68 @@ func TestSubmissionSentCopyDedupesByMessageID(t *testing.T) {
if count != 1 {
t.Fatalf("sent copy count=%d, want 1", count)
}
var queueCount int
if err := a.db.QueryRow(`SELECT COUNT(1) FROM send_queue WHERE mailbox_id=? AND message_id=?`, mb.ID, "<dedupe@example.test>").Scan(&queueCount); err != nil {
t.Fatal(err)
}
if queueCount != 1 {
t.Fatalf("send queue count=%d, want 1", queueCount)
}
}
func TestSubmissionAllowsAuthorizedAliasSendAs(t *testing.T) {
a := newTestApp(t)
ctx := context.Background()
if _, err := a.db.ExecContext(ctx, `INSERT INTO aliases(id,domain_id,source,destination,enabled,created_at,updated_at) VALUES(?,?,?,?,?,?,?)`, newID("als"), mustDefaultDomainID(t, a), "team@lanqin.local", "admin@lanqin.local", 1, a.now().UTC().Format(time.RFC3339Nano), a.now().UTC().Format(time.RFC3339Nano)); err != nil {
t.Fatal(err)
}
user, mb, err := a.authenticateSubmission(ctx, "admin@lanqin.local", "ChangeMe123!")
if err != nil {
t.Fatal(err)
}
raw := "From: Team <team@lanqin.local>\r\nTo: person@example.com\r\nSubject: alias send-as\r\nMessage-ID: <alias-send-as@example.test>\r\n\r\nbody"
if err := a.submitSMTPMessage(ctx, user, mb, "team@lanqin.local", []string{"person@example.com"}, strings.NewReader(raw)); err != nil {
t.Fatalf("authorized alias send-as should submit: %v", err)
}
sentFolderID, err := a.ensureFolder(ctx, mb.ID, "Sent")
if err != nil {
t.Fatal(err)
}
var fromAddr string
if err := a.db.QueryRow(`SELECT from_addr FROM messages WHERE mailbox_id=? AND folder_id=? AND message_id=?`, mb.ID, sentFolderID, "<alias-send-as@example.test>").Scan(&fromAddr); err != nil {
t.Fatal(err)
}
if fromAddr != "team@lanqin.local" {
t.Fatalf("from_addr=%q, want alias", fromAddr)
}
}
func TestSubmissionAllowsExplicitSendAsGrant(t *testing.T) {
a := newTestApp(t)
ctx := context.Background()
user, mb, err := a.authenticateSubmission(ctx, "admin@lanqin.local", "ChangeMe123!")
if err != nil {
t.Fatal(err)
}
now := a.now().UTC().Format(time.RFC3339Nano)
if _, err := a.db.ExecContext(ctx, `INSERT INTO send_as_grants(id,mailbox_id,address,display_name,enabled,created_at,updated_at) VALUES(?,?,?,?,?,?,?)`, newID("sag"), mb.ID, "support@example.com", "Support", 1, now, now); err != nil {
t.Fatal(err)
}
raw := "From: Support <support@example.com>\r\nTo: person@example.com\r\nSubject: explicit send-as\r\nMessage-ID: <explicit-send-as@example.test>\r\n\r\nbody"
if err := a.submitSMTPMessage(ctx, user, mb, "support@example.com", []string{"person@example.com"}, strings.NewReader(raw)); err != nil {
t.Fatalf("explicit send-as grant should submit: %v", err)
}
sentFolderID, err := a.ensureFolder(ctx, mb.ID, "Sent")
if err != nil {
t.Fatal(err)
}
var fromAddr, fromName string
if err := a.db.QueryRow(`SELECT from_addr,from_name FROM messages WHERE mailbox_id=? AND folder_id=? AND message_id=?`, mb.ID, sentFolderID, "<explicit-send-as@example.test>").Scan(&fromAddr, &fromName); err != nil {
t.Fatal(err)
}
if fromAddr != "support@example.com" || fromName != "Support" {
t.Fatalf("from=%q name=%q, want explicit grant", fromAddr, fromName)
}
}
func TestSentMessageDedupeTableExists(t *testing.T) {
@@ -1012,6 +1189,9 @@ func TestSubmissionServersAcceptStartTLSAndImplicitTLS(t *testing.T) {
t.Fatal(err)
}
_ = client.Close()
if err := a.processDueSendQueue(context.Background()); err != nil {
t.Fatal(err)
}
select {
case <-received:
case <-time.After(2 * time.Second):
@@ -1031,6 +1211,9 @@ func TestSubmissionServersAcceptStartTLSAndImplicitTLS(t *testing.T) {
t.Fatal(err)
}
_ = client.Close()
if err := a.processDueSendQueue(context.Background()); err != nil {
t.Fatal(err)
}
select {
case <-received:
case <-time.After(2 * time.Second):