diff --git a/README.md b/README.md index 9f56e24..d5a7649 100644 --- a/README.md +++ b/README.md @@ -195,8 +195,9 @@ docker compose -f docker-compose.yml -f docker-compose.build.yml up -d --build - 第三方客户端的 SMTP 提交 `465/587` 由 LanQin API 进程处理。 - Postfix 只保留 `25` 端口,用于公网入站邮件和内部/外部 relay。 -- Webmail/API 发信继续由现有 API 发信流程写入 Sent。 -- 第三方客户端发信会先校验邮箱密码,写入 Sent,再 relay 到 `LANQIN_SMTP_HOST:LANQIN_SMTP_PORT`。 +- Webmail/API 和第三方客户端发信都会先写入 Sent,再进入发送队列。 +- 发送队列由 LanQin API 后台 worker relay 到 `LANQIN_SMTP_HOST:LANQIN_SMTP_PORT`,失败会记录审计并按退避策略重试。 +- v1 支持本人邮箱发信;如需 send-as,可使用启用的别名转发 source 指向本人邮箱,或在数据库中配置 `send_as_grants`。 - 如果客户端随后又通过 IMAP APPEND 写入自己的 Sent 副本,Maildir 同步会按 Sent 文件夹内的 `Message-ID` 去重。 ## License diff --git a/apps/api/internal/app/app.go b/apps/api/internal/app/app.go index 2a3524d..29689f5 100644 --- a/apps/api/internal/app/app.go +++ b/apps/api/internal/app/app.go @@ -74,6 +74,7 @@ func New(cfg Config, logger *slog.Logger) (*App, error) { if strings.TrimSpace(a.cfg.MaildirRoot) != "" { go a.maildirWorker(workerCtx) } + go a.sendQueueWorker(workerCtx) go a.smtpEventsCleanupWorker(workerCtx) return a, nil } @@ -240,6 +241,53 @@ func (a *App) migrate(ctx context.Context) error { created_at TEXT NOT NULL, PRIMARY KEY(mailbox_id, folder_id, message_id) )`, + `CREATE TABLE IF NOT EXISTS send_as_grants ( + id TEXT PRIMARY KEY, + mailbox_id TEXT NOT NULL REFERENCES mailboxes(id) ON DELETE CASCADE, + address TEXT NOT NULL, + display_name TEXT NOT NULL DEFAULT '', + enabled INTEGER NOT NULL DEFAULT 1, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + UNIQUE(mailbox_id, address) + )`, + `CREATE TABLE IF NOT EXISTS send_queue ( + id TEXT PRIMARY KEY, + user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE, + mailbox_id TEXT NOT NULL REFERENCES mailboxes(id) ON DELETE CASCADE, + sent_message_id TEXT NOT NULL DEFAULT '', + message_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL, + mail_from TEXT NOT NULL, + header_from TEXT NOT NULL, + recipients_json TEXT NOT NULL, + mime_base64 TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'queued', + attempt_count INTEGER NOT NULL DEFAULT 0, + max_attempts INTEGER NOT NULL DEFAULT 5, + next_attempt_at TEXT NOT NULL, + last_error TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + delivered_at TEXT + )`, + `CREATE INDEX IF NOT EXISTS idx_send_queue_due ON send_queue(status, next_attempt_at, created_at)`, + `CREATE TABLE IF NOT EXISTS send_audit_events ( + id TEXT PRIMARY KEY, + queue_id TEXT NOT NULL DEFAULT '', + user_id TEXT NOT NULL DEFAULT '', + mailbox_id TEXT NOT NULL DEFAULT '', + sent_message_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL, + event TEXT NOT NULL, + status TEXT NOT NULL, + mail_from TEXT NOT NULL DEFAULT '', + header_from TEXT NOT NULL DEFAULT '', + recipients_json TEXT NOT NULL DEFAULT '[]', + error TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL + )`, + `CREATE INDEX IF NOT EXISTS idx_send_audit_events_created ON send_audit_events(created_at)`, `CREATE TABLE IF NOT EXISTS attachments ( id TEXT PRIMARY KEY, message_id TEXT NOT NULL REFERENCES messages(id) ON DELETE CASCADE, @@ -377,12 +425,47 @@ func (a *App) migrate(ctx context.Context) error { if err := a.migratePermissionGroupLimits(ctx); err != nil { return err } + if err := a.migrateSendQueueMessageID(ctx); err != nil { + return err + } if err := a.ensureDefaultPermissionGroups(ctx); err != nil { return err } return nil } +func (a *App) migrateSendQueueMessageID(ctx context.Context) error { + rows, err := a.db.QueryContext(ctx, `PRAGMA table_info(send_queue)`) + if err != nil { + return err + } + hasMessageID := false + for rows.Next() { + var cid int + var name, typ string + var notnull int + var dflt any + var pk int + if err := rows.Scan(&cid, &name, &typ, ¬null, &dflt, &pk); err != nil { + rows.Close() + return err + } + if name == "message_id" { + hasMessageID = true + } + } + if err := rows.Close(); err != nil { + return err + } + if !hasMessageID { + if _, err := a.db.ExecContext(ctx, `ALTER TABLE send_queue ADD COLUMN message_id TEXT NOT NULL DEFAULT ''`); err != nil { + return err + } + } + _, err = a.db.ExecContext(ctx, `CREATE UNIQUE INDEX IF NOT EXISTS idx_send_queue_mailbox_source_message_id ON send_queue(mailbox_id, source, message_id) WHERE message_id <> ''`) + return err +} + func (a *App) migratePermissionGroupLimits(ctx context.Context) error { rows, err := a.db.QueryContext(ctx, `PRAGMA table_info(permission_groups)`) if err != nil { diff --git a/apps/api/internal/app/app_test.go b/apps/api/internal/app/app_test.go index ca1dee2..98f8e91 100644 --- a/apps/api/internal/app/app_test.go +++ b/apps/api/internal/app/app_test.go @@ -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: \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, "").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, "").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 \r\nTo: person@example.com\r\nSubject: alias send-as\r\nMessage-ID: \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, "").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 \r\nTo: person@example.com\r\nSubject: explicit send-as\r\nMessage-ID: \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, "").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): diff --git a/apps/api/internal/app/mail_handlers.go b/apps/api/internal/app/mail_handlers.go index 596c467..43bb84e 100644 --- a/apps/api/internal/app/mail_handlers.go +++ b/apps/api/internal/app/mail_handlers.go @@ -336,6 +336,8 @@ func (a *App) handleMailMessage(w http.ResponseWriter, r *http.Request) { type mailComposeInput struct { MailboxID string `json:"mailboxId"` + From string `json:"from"` + FromName string `json:"fromName"` To []string `json:"to"` CC []string `json:"cc"` BCC []string `json:"bcc"` @@ -358,6 +360,8 @@ type mailDraftInput struct { type scheduledSendPayload struct { MailboxID string `json:"mailboxId"` + From string `json:"from"` + FromName string `json:"fromName"` To []string `json:"to"` CC []string `json:"cc"` BCC []string `json:"bcc"` @@ -412,10 +416,6 @@ func (a *App) handleMailSend(w http.ResponseWriter, r *http.Request) { respondError(w, http.StatusTooManyRequests, err.Error()) return } - if strings.HasPrefix(err.Error(), "smtp delivery failed:") { - respondError(w, http.StatusBadGateway, err.Error()) - return - } respondError(w, http.StatusInternalServerError, err.Error()) return } @@ -448,9 +448,16 @@ func (a *App) sendMailNow(ctx context.Context, user *User, mb *Mailbox, req mail } now := a.now().UTC() + fromAddress, fromName, err := a.authorizedSender(ctx, mb, req.From) + if err != nil { + return nil, err + } + if strings.TrimSpace(req.FromName) != "" && normalizeEmail(req.From) == normalizeEmail(mb.Address) { + fromName = strings.TrimSpace(req.FromName) + } messageID := fmt.Sprintf("<%s@%s>", newID("msg"), strings.Split(mb.Address, "@")[1]) mimeBytes, err := BuildMIME(MIMEMessage{ - From: mb.Address, FromName: mb.DisplayName, To: req.To, CC: req.CC, BCC: req.BCC, Subject: req.Subject, Text: req.Text, HTML: req.HTML, MessageID: messageID, Date: now, Attachments: req.Attachments, + From: fromAddress, FromName: fromName, To: req.To, CC: req.CC, BCC: req.BCC, Subject: req.Subject, Text: req.Text, HTML: req.HTML, MessageID: messageID, Date: now, Attachments: req.Attachments, }) if err != nil { return nil, fmt.Errorf("%w: %v", errInvalidMIME, err) @@ -458,21 +465,21 @@ func (a *App) sendMailNow(ctx context.Context, user *User, mb *Mailbox, req mail if err := a.recordSMTPRate(ctx, user, mb); err != nil { return nil, err } - if a.cfg.SMTPHost != "" { - if err := a.sendSMTP(mb.Address, allRecipients, mimeBytes); err != nil { - return nil, fmt.Errorf("smtp delivery failed: %w", err) - } - } sentFolderID, err := a.ensureFolder(ctx, mb.ID, "Sent") if err != nil { return nil, fmt.Errorf("failed to load sent folder: %w", err) } - base := storedMessage{MailboxID: mb.ID, FolderID: sentFolderID, MessageUID: newID("uid"), MessageID: messageID, Subject: req.Subject, From: mb.Address, FromName: mb.DisplayName, To: req.To, CC: req.CC, BCC: req.BCC, SentAt: now, ReceivedAt: now, Snippet: snippetFrom(req.Text, req.HTML), BodyText: req.Text, BodyHTML: req.HTML, IsRead: true} + base := storedMessage{MailboxID: mb.ID, FolderID: sentFolderID, MessageUID: newID("uid"), MessageID: messageID, Subject: req.Subject, From: fromAddress, FromName: fromName, To: req.To, CC: req.CC, BCC: req.BCC, SentAt: now, ReceivedAt: now, Snippet: snippetFrom(req.Text, req.HTML), BodyText: req.Text, BodyHTML: req.HTML, IsRead: true} sentID, err := a.insertMessage(ctx, base, req.Attachments) if err != nil { return nil, fmt.Errorf("failed to store sent message: %w", err) } + a.recordSendAudit(ctx, sendAuditAccepted, sendQueueStatusQueued, sendAuditInput{UserID: user.ID, MailboxID: mb.ID, SentMessageID: sentID, Source: sendSourceWebmail, MailFrom: fromAddress, HeaderFrom: fromAddress, Recipients: allRecipients}) + if _, err := a.enqueueSend(ctx, sendQueueInput{UserID: user.ID, MailboxID: mb.ID, SentMessageID: sentID, MessageID: messageID, Source: sendSourceWebmail, MailFrom: fromAddress, HeaderFrom: fromAddress, Recipients: allRecipients, MIMEBytes: mimeBytes, Now: now}); err != nil { + a.deleteMessage(ctx, sentID) + return nil, fmt.Errorf("failed to enqueue delivery: %w", err) + } // Development/local-domain delivery: known local recipients go to their Inbox. // When catch-all is enabled, unknown local recipients are stored as unregistered @@ -925,7 +932,7 @@ func (a *App) handleScheduleSend(w http.ResponseWriter, r *http.Request) { } } now := a.now().UTC().Format(time.RFC3339Nano) - payload := scheduledSendPayload{MailboxID: compose.MailboxID, To: compose.To, CC: compose.CC, BCC: compose.BCC, Subject: compose.Subject, Text: compose.Text, HTML: compose.HTML, Attachments: compose.Attachments, DraftID: draftID} + payload := scheduledSendPayload{MailboxID: compose.MailboxID, From: compose.From, FromName: compose.FromName, To: compose.To, CC: compose.CC, BCC: compose.BCC, Subject: compose.Subject, Text: compose.Text, HTML: compose.HTML, Attachments: compose.Attachments, DraftID: draftID} item := ScheduledSend{ID: newID("sched"), MailboxID: mb.ID, DraftID: draftID, Subject: payload.Subject, To: payload.To, Snippet: snippetFrom(payload.Text, payload.HTML), SendAt: sendAt.UTC(), Status: "pending", CreatedAt: parseTime(now), UpdatedAt: parseTime(now)} if _, err := a.db.ExecContext(r.Context(), `INSERT INTO scheduled_sends(id,user_id,mailbox_id,draft_id,payload_json,send_at,status,created_at,updated_at) VALUES(?,?,?,?,?,?,?,?,?)`, item.ID, currentUser(r).ID, mb.ID, nullableString(draftID), jsonEncode(payload), item.SendAt.Format(time.RFC3339Nano), item.Status, now, now); err != nil { respondError(w, http.StatusInternalServerError, "failed to schedule send") diff --git a/apps/api/internal/app/send_queue.go b/apps/api/internal/app/send_queue.go new file mode 100644 index 0000000..1dab033 --- /dev/null +++ b/apps/api/internal/app/send_queue.go @@ -0,0 +1,318 @@ +package app + +import ( + "context" + "database/sql" + "encoding/base64" + "errors" + "fmt" + "strings" + "time" +) + +const ( + sendQueueStatusQueued = "queued" + sendQueueStatusSending = "sending" + sendQueueStatusDelivered = "delivered" + sendQueueStatusFailed = "failed" + + sendAuditAccepted = "accepted" + sendAuditQueued = "queued" + sendAuditDelivered = "delivered" + sendAuditFailed = "failed" + sendAuditRetry = "retry" + + sendSourceWebmail = "webmail" + sendSourceSubmission = "submission" + + sendQueueStaleAfter = 15 * time.Minute +) + +type sendQueueInput struct { + UserID string + MailboxID string + SentMessageID string + MessageID string + Source string + MailFrom string + HeaderFrom string + Recipients []string + MIMEBytes []byte + Now time.Time +} + +type sendQueueItem struct { + ID string + UserID string + MailboxID string + SentMessageID string + MessageID string + Source string + MailFrom string + HeaderFrom string + Recipients []string + MIMEBytes []byte + AttemptCount int + MaxAttempts int +} + +func (a *App) enqueueSend(ctx context.Context, in sendQueueInput) (string, error) { + if strings.TrimSpace(a.cfg.SMTPHost) == "" { + return "", nil + } + now := in.Now.UTC() + if now.IsZero() { + now = a.now().UTC() + } + id := newID("snd") + messageID := strings.TrimSpace(in.MessageID) + _, err := a.db.ExecContext(ctx, `INSERT OR IGNORE INTO send_queue(id,user_id,mailbox_id,sent_message_id,message_id,source,mail_from,header_from,recipients_json,mime_base64,status,next_attempt_at,created_at,updated_at) + VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, + id, in.UserID, in.MailboxID, in.SentMessageID, messageID, in.Source, normalizeEmail(in.MailFrom), normalizeEmail(in.HeaderFrom), jsonEncode(dedupeEmails(in.Recipients)), base64.StdEncoding.EncodeToString(in.MIMEBytes), sendQueueStatusQueued, now.Format(time.RFC3339Nano), now.Format(time.RFC3339Nano), now.Format(time.RFC3339Nano)) + if err != nil { + return "", err + } + if strings.TrimSpace(messageID) != "" { + var existingID string + if err := a.db.QueryRowContext(ctx, `SELECT id FROM send_queue WHERE mailbox_id=? AND source=? AND message_id=?`, in.MailboxID, in.Source, messageID).Scan(&existingID); err != nil { + return "", err + } + if existingID != id { + return existingID, nil + } + } + a.recordSendAudit(ctx, sendAuditQueued, sendQueueStatusQueued, sendAuditInput{ + QueueID: id, + UserID: in.UserID, + MailboxID: in.MailboxID, + SentMessageID: in.SentMessageID, + Source: in.Source, + MailFrom: in.MailFrom, + HeaderFrom: in.HeaderFrom, + Recipients: in.Recipients, + }) + return id, nil +} + +func (a *App) sendQueueWorker(ctx context.Context) { + a.log.Info("send queue worker started") + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + for { + if err := a.processDueSendQueue(ctx); err != nil { + a.log.Warn("send queue worker failed", "error", err) + } + select { + case <-ctx.Done(): + a.log.Info("send queue worker stopped") + return + case <-ticker.C: + } + } +} + +func (a *App) processDueSendQueue(ctx context.Context) error { + if strings.TrimSpace(a.cfg.SMTPHost) == "" { + return nil + } + if err := a.recoverStaleSendQueueItems(ctx); err != nil { + return err + } + rows, err := a.db.QueryContext(ctx, `SELECT id FROM send_queue WHERE (status=? OR (status=? AND attempt_count 0 { + a.recordSendAudit(ctx, sendAuditRetry, sendQueueStatusFailed, sendAuditInputFromQueue(item, "send attempt interrupted")) + } + } + return nil +} + +func (a *App) processSendQueueItem(ctx context.Context, id string) { + item, err := a.claimSendQueueItem(ctx, id) + if err != nil { + if !errors.Is(err, sql.ErrNoRows) { + a.log.Warn("failed to claim send queue item", "id", id, "error", err) + } + return + } + if err := a.sendSMTP(item.MailFrom, item.Recipients, item.MIMEBytes); err != nil { + 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='' WHERE id=?`, sendQueueStatusDelivered, now, now, item.ID); err != nil { + a.log.Warn("failed to mark send queue delivered", "id", item.ID, "error", err) + return + } + a.recordSendAudit(ctx, sendAuditDelivered, sendQueueStatusDelivered, sendAuditInputFromQueue(item, "")) +} + +func (a *App) claimSendQueueItem(ctx context.Context, id string) (sendQueueItem, error) { + now := a.now().UTC().Format(time.RFC3339Nano) + res, err := a.db.ExecContext(ctx, `UPDATE send_queue SET status=?,attempt_count=attempt_count+1,updated_at=? WHERE id=? AND (status=? OR (status=? AND attempt_count= item.MaxAttempts { + nextAttempt = now.Add(365 * 24 * time.Hour) + } + _, err := a.db.ExecContext(ctx, `UPDATE send_queue SET status=?,next_attempt_at=?,last_error=?,updated_at=? WHERE id=?`, status, nextAttempt.Format(time.RFC3339Nano), sendErr.Error(), now.Format(time.RFC3339Nano), item.ID) + if err != nil { + a.log.Warn("failed to mark send queue failed", "id", item.ID, "error", err) + } + event := sendAuditRetry + if item.AttemptCount >= item.MaxAttempts { + event = sendAuditFailed + } + a.recordSendAudit(ctx, event, status, sendAuditInputFromQueue(item, sendErr.Error())) +} + +func sendRetryDelay(attempt int) time.Duration { + if attempt < 1 { + attempt = 1 + } + delays := []time.Duration{30 * time.Second, 2 * time.Minute, 10 * time.Minute, time.Hour, 6 * time.Hour} + if attempt > len(delays) { + return delays[len(delays)-1] + } + return delays[attempt-1] +} + +type sendAuditInput struct { + QueueID string + UserID string + MailboxID string + SentMessageID string + Source string + MailFrom string + HeaderFrom string + Recipients []string + Error string +} + +func sendAuditInputFromQueue(item sendQueueItem, errorText string) sendAuditInput { + return sendAuditInput{ + QueueID: item.ID, + UserID: item.UserID, + MailboxID: item.MailboxID, + SentMessageID: item.SentMessageID, + Source: item.Source, + MailFrom: item.MailFrom, + HeaderFrom: item.HeaderFrom, + Recipients: item.Recipients, + Error: errorText, + } +} + +func (a *App) recordSendAudit(ctx context.Context, event, status string, in sendAuditInput) { + source := strings.TrimSpace(in.Source) + if source == "" { + source = "unknown" + } + _, err := a.db.ExecContext(ctx, `INSERT INTO send_audit_events(id,queue_id,user_id,mailbox_id,sent_message_id,source,event,status,mail_from,header_from,recipients_json,error,created_at) + VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)`, newID("audit"), in.QueueID, in.UserID, in.MailboxID, in.SentMessageID, source, event, status, normalizeEmail(in.MailFrom), normalizeEmail(in.HeaderFrom), jsonEncode(dedupeEmails(in.Recipients)), in.Error, a.now().UTC().Format(time.RFC3339Nano)) + if err != nil { + a.log.Warn("failed to record send audit", "event", event, "error", err) + } +} + +func (a *App) authorizedSender(ctx context.Context, mb *Mailbox, from string) (string, string, error) { + from = normalizeEmail(from) + if from == "" { + from = normalizeEmail(mb.Address) + } + if from == normalizeEmail(mb.Address) { + return normalizeEmail(mb.Address), mb.DisplayName, nil + } + var displayName string + var enabled int + err := a.db.QueryRowContext(ctx, `SELECT display_name,enabled FROM send_as_grants WHERE mailbox_id=? AND address=?`, mb.ID, from).Scan(&displayName, &enabled) + if err == nil { + if enabled == 0 { + return "", "", fmt.Errorf("send-as address is disabled") + } + return from, strings.TrimSpace(displayName), nil + } + if !errors.Is(err, sql.ErrNoRows) { + return "", "", err + } + var aliasDestination string + err = a.db.QueryRowContext(ctx, `SELECT destination FROM aliases WHERE source=? AND enabled=1`, from).Scan(&aliasDestination) + if err == nil && normalizeEmail(aliasDestination) == normalizeEmail(mb.Address) { + return from, mb.DisplayName, nil + } + return "", "", fmt.Errorf("send-as address is not authorized") +} diff --git a/apps/api/internal/app/submission.go b/apps/api/internal/app/submission.go index e745987..16ac7c2 100644 --- a/apps/api/internal/app/submission.go +++ b/apps/api/internal/app/submission.go @@ -185,7 +185,8 @@ func (s *submissionSession) Mail(from string, _ *smtpserver.MailOptions) error { return smtpserver.ErrAuthRequired } from = normalizeEmail(from) - if from == "" || from != normalizeEmail(s.mailbox.Address) { + authorized, _, err := s.app.authorizedSender(context.Background(), s.mailbox, from) + if err != nil || from == "" || from != authorized { return smtpError(553, smtpserver.EnhancedCode{5, 7, 1}, "sender must match authenticated mailbox") } s.mailFrom = from @@ -273,30 +274,31 @@ func (a *App) submitSMTPMessage(ctx context.Context, user *User, mb *Mailbox, ma if err != nil { return err } - prepared, msg, attachments, err := a.prepareSubmittedMessage(raw, mb.Address, mailFrom, recipients) + prepared, msg, attachments, err := a.prepareSubmittedMessage(ctx, raw, mb, mailFrom, recipients) if err != nil { return err } msg.MailboxID = mb.ID - sentID, err := a.insertSentMessageOnce(ctx, msg, attachments) + sentID, insertedSent, err := a.insertSentMessageOnce(ctx, msg, attachments) if err != nil { return err } - if a.cfg.SMTPHost != "" { - if err := a.sendSMTP(mb.Address, recipients, prepared); err != nil { - if sentID != "" { + a.recordSendAudit(ctx, sendAuditAccepted, sendQueueStatusQueued, sendAuditInput{UserID: user.ID, MailboxID: mb.ID, SentMessageID: sentID, Source: sendSourceSubmission, MailFrom: mailFrom, HeaderFrom: msg.From, Recipients: recipients}) + if sentID != "" { + if _, err := a.enqueueSend(ctx, sendQueueInput{UserID: user.ID, MailboxID: mb.ID, SentMessageID: sentID, MessageID: msg.MessageID, Source: sendSourceSubmission, MailFrom: mailFrom, HeaderFrom: msg.From, Recipients: recipients, MIMEBytes: prepared, Now: a.now().UTC()}); err != nil { + if insertedSent { a.deleteMessage(ctx, sentID) if sentFolderID, ferr := a.ensureFolder(ctx, mb.ID, "Sent"); ferr == nil { a.deleteSentDedupeKey(ctx, mb.ID, sentFolderID, msg.MessageID) } } - return smtpError(451, smtpserver.EnhancedCode{4, 4, 0}, "smtp relay failed") + return err } } return nil } -func (a *App) prepareSubmittedMessage(raw []byte, authenticatedAddress, mailFrom string, recipients []string) ([]byte, storedMessage, []AttachmentInput, error) { +func (a *App) prepareSubmittedMessage(ctx context.Context, raw []byte, mb *Mailbox, mailFrom string, recipients []string) ([]byte, storedMessage, []AttachmentInput, error) { header, body, err := readMessageHeader(raw) if err != nil { return nil, storedMessage{}, nil, smtpError(554, smtpserver.EnhancedCode{5, 6, 0}, "invalid message") @@ -305,8 +307,8 @@ func (a *App) prepareSubmittedMessage(raw []byte, authenticatedAddress, mailFrom if !ok || fromAddress == "" { return nil, storedMessage{}, nil, smtpError(550, smtpserver.EnhancedCode{5, 7, 1}, "From header must contain exactly one address") } - authAddress := normalizeEmail(authenticatedAddress) - if normalizeEmail(mailFrom) != authAddress || normalizeEmail(fromAddress) != authAddress { + authAddress, fromName, err := a.authorizedSender(ctx, mb, fromAddress) + if err != nil || normalizeEmail(mailFrom) != authAddress || normalizeEmail(fromAddress) != authAddress { return nil, storedMessage{}, nil, smtpError(553, smtpserver.EnhancedCode{5, 7, 1}, "sender must match authenticated mailbox") } now := a.now().UTC() @@ -353,10 +355,10 @@ func (a *App) prepareSubmittedMessage(raw []byte, authenticatedAddress, mailFrom return prepared, msg, attachments, nil } -func (a *App) insertSentMessageOnce(ctx context.Context, msg storedMessage, attachments []AttachmentInput) (string, error) { +func (a *App) insertSentMessageOnce(ctx context.Context, msg storedMessage, attachments []AttachmentInput) (string, bool, error) { sentFolderID, err := a.ensureFolder(ctx, msg.MailboxID, "Sent") if err != nil { - return "", err + return "", false, err } msg.FolderID = sentFolderID if msg.MessageUID == "" { @@ -367,18 +369,22 @@ func (a *App) insertSentMessageOnce(ctx context.Context, msg storedMessage, atta err := a.db.QueryRowContext(ctx, `SELECT id FROM messages WHERE mailbox_id=? AND folder_id=? AND message_id=? AND message_id <> '' LIMIT 1`, msg.MailboxID, sentFolderID, msg.MessageID).Scan(&existing) if err == nil { if err := a.insertSentDedupeKey(ctx, msg.MailboxID, sentFolderID, msg.MessageID); err != nil && !errors.Is(err, errSentDedupeExists) { - return "", err + return "", false, err } - return "", nil + return existing, false, nil } if err != nil && !errors.Is(err, sql.ErrNoRows) { - return "", err + return "", false, err } if err := a.insertSentDedupeKey(ctx, msg.MailboxID, sentFolderID, msg.MessageID); err != nil { if errors.Is(err, errSentDedupeExists) { - return "", nil + var existing string + if err := a.db.QueryRowContext(ctx, `SELECT id FROM messages WHERE mailbox_id=? AND folder_id=? AND message_id=? AND message_id <> '' LIMIT 1`, msg.MailboxID, sentFolderID, msg.MessageID).Scan(&existing); err == nil { + return existing, false, nil + } + return "", false, nil } - return "", err + return "", false, err } } id, err := a.insertMessage(ctx, msg, attachments) @@ -386,9 +392,9 @@ func (a *App) insertSentMessageOnce(ctx context.Context, msg storedMessage, atta if msg.MessageID != "" { a.deleteSentDedupeKey(ctx, msg.MailboxID, sentFolderID, msg.MessageID) } - return "", err + return "", false, err } - return id, nil + return id, true, nil } var errSentDedupeExists = errors.New("sent message already exists") diff --git a/deploy/README.md b/deploy/README.md index c217419..0a3a03e 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -128,7 +128,8 @@ docker compose -f docker-compose.stack.yml -f docker-compose.stack.build.yml up - Rspamd 会周期性从 SQLite 导出域名 DKIM 私钥到容器内 `/var/lib/rspamd/dkim`。 - Go API 是 Webmail 和管理后台入口;浏览器不直接连接 SMTP/IMAP/POP3。 - Go API 会读取 `LANQIN_MAILDIR_ROOT=/var/mail/vhosts`,周期扫描 Maildir,把 Postfix/Dovecot 入站邮件同步成 Webmail 索引。 -- 第三方客户端可通过 LanQin API 提供的 SMTP `465/587` 发信;Webmail/API 和第三方客户端的“已发送”都由 API 写入,客户端后续 IMAP APPEND 到 Sent 会按 `Message-ID` 去重。 +- 第三方客户端可通过 LanQin API 提供的 SMTP `465/587` 发信;Webmail/API 和第三方客户端的“已发送”都由 API 写入,外发投递进入发送队列并由 API worker relay/retry,客户端后续 IMAP APPEND 到 Sent 会按 `Message-ID` 去重。 +- send-as v1 支持本人邮箱、启用的别名转发 source 指向本人邮箱,或数据库表 `send_as_grants` 中显式授权的地址。 ## 邮件客户端 TLS 证书 @@ -172,13 +173,14 @@ LANQIN_SMTP_REQUIRE_TLS=false Split stack 使用 `docker-compose.stack.yml` 时,API 容器默认会把 `LANQIN_SMTP_HOST` 覆盖为 `postfix`,让 Webmail 和 SMTP 提交都 relay 到 Postfix service。只有改用外部 SMTP 时才需要在 `.env` 明确填写 `LANQIN_STACK_SMTP_HOST` / `LANQIN_STACK_SMTP_PORT`。 -如果页面提示 `smtp delivery failed: EOF`,通常是 Postfix 会话被中断。优先检查: +如果发送队列里出现 relay 失败,通常是 Postfix 会话被中断或外部 SMTP 配置错误。优先检查: ```bash docker compose exec lanqin-email supervisorctl status docker compose exec lanqin-email postconf -M smtp/inet # SMTP 提交 465/587 由 LanQin API 提供,不再由 Postfix 监听。 docker compose exec lanqin-email sqlite3 /data/lanqin.db "select key,value from system_settings where key like 'smtp%' order by key;" +docker compose exec lanqin-email sqlite3 /data/lanqin.db "select status,attempt_count,last_error from send_queue order by created_at desc limit 10;" docker compose logs --tail=200 lanqin-email ```