Phase 6: journal-milter + send-log status tailer + retention

Implement the structured send log (spec 7.3), the project's highest-risk
component since a milter bug can break the relay itself.

- internal/milter: go-milter v0.4.1 journal-milter. Per-connection session
  collects SASL login, From, recipients and Subject across callbacks and
  writes one send_log "queued" row per (queue-id, recipient) at EOM
  (spec 7.3.3). Monitoring only: callbacks return Continue/Accept, recorder
  errors are logged never propagated, so it can never block mail.
- internal/logtail: polling mail.log tailer with rotation handling (inode
  change / truncation), parses sent/deferred/bounced/expired by queue-id +
  recipient and advances rows; background retention sweep prunes rows past
  SEND_LOG_RETENTION_DAYS (default 90) at startup and every 6h.
- internal/store/sendlog.go: InsertQueued, UpdateStatus (case-insensitive
  recipient match), DeleteSendLogBefore + status constants.
- cmd/panel: open the store once and share it across http/milter/tailer;
  replace the journal/logtail stubs with the real roles.
- build/postfix-config.sh: bounded milter timeouts (15/15/30s) so a hung
  milter also fails open in seconds, not the 300s default.

Fix found in-container: SASL login (app_login) was empty because go-milter
keys macros exactly as Postfix sends them, and multi-character macro names
arrive brace-wrapped ({auth_authen}); the SASL-less Phase 0 spike could not
observe this. Added a brace-tolerant macro lookup.

Verified on selfpost.example.com: gofmt/vet/unit tests green; container e2e
records rows with correct fields and advances status via the tailer; fail-open
confirmed for both an unreachable and a hung milter; retention prunes at start.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-13 22:58:34 +03:00
parent 4f1f7761a2
commit 6ebb6f56d6
13 changed files with 1033 additions and 51 deletions
+82
View File
@@ -0,0 +1,82 @@
package store
import (
"fmt"
"time"
)
// Send-log status values (spec 7.3). "queued" is written by the journal-milter
// when a message is accepted; the log-tailer advances it to one of the final
// states as Postfix reports delivery per recipient.
const (
StatusQueued = "queued"
StatusSent = "sent"
StatusDeferred = "deferred"
StatusBounced = "bounced"
)
// SendLogEntry is a single queued send-log row. The journal-milter creates one
// per (queue-id, recipient) pair at end-of-message (spec 7.3.3); every field
// except the status/timestamps comes from the accepted message.
type SendLogEntry struct {
QueueID string
Domain string
AppLogin string
From string
To string
Subject string
}
// InsertQueued records an accepted message in the send log with status
// "queued". It is called from the journal-milter hot path, so it returns any
// error for the caller to log rather than deciding policy here; the milter must
// stay fail-open regardless (spec 7.3).
func (s *Store) InsertQueued(e SendLogEntry) error {
now := time.Now().UTC().Format(time.RFC3339)
_, err := s.db.Exec(
`INSERT INTO send_log
(queue_id, domain, app_login, from_addr, to_addr, subject, status, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
e.QueueID, e.Domain, e.AppLogin, e.From, e.To, e.Subject, StatusQueued, now, now,
)
if err != nil {
return fmt.Errorf("insert send_log: %w", err)
}
return nil
}
// UpdateStatus advances the delivery status of the send-log rows matching a
// (queue-id, recipient) pair, which the log-tailer parses out of mail.log.
// Recipient matching is case-insensitive because Postfix may normalise address
// case between the milter (envelope) and the delivery log. It returns the
// number of rows updated so the caller can tell whether the line matched a
// journal entry.
func (s *Store) UpdateStatus(queueID, recipient, status string) (int64, error) {
now := time.Now().UTC().Format(time.RFC3339)
res, err := s.db.Exec(
`UPDATE send_log SET status = ?, updated_at = ?
WHERE queue_id = ? AND to_addr = ? COLLATE NOCASE`,
status, now, queueID, recipient,
)
if err != nil {
return 0, fmt.Errorf("update send_log status: %w", err)
}
n, _ := res.RowsAffected()
return n, nil
}
// DeleteSendLogBefore removes send-log rows created before cutoff, implementing
// the configurable retention window (spec 7.3, SEND_LOG_RETENTION_DAYS). It
// returns the number of rows pruned. created_at is stored as RFC3339 UTC, so a
// lexical comparison against the same format is chronologically correct.
func (s *Store) DeleteSendLogBefore(cutoff time.Time) (int64, error) {
res, err := s.db.Exec(
`DELETE FROM send_log WHERE created_at < ?`,
cutoff.UTC().Format(time.RFC3339),
)
if err != nil {
return 0, fmt.Errorf("prune send_log: %w", err)
}
n, _ := res.RowsAffected()
return n, nil
}
+140
View File
@@ -0,0 +1,140 @@
package store
import (
"testing"
"time"
)
// readSendLog returns every send_log row ordered by id. Phase 6 has no read
// query yet (the monitoring UI is Phase 7), so tests read the table directly.
type sendLogRow struct {
QueueID string
Domain string
AppLogin string
From string
To string
Subject string
Status string
}
func readSendLog(t *testing.T, s *Store) []sendLogRow {
t.Helper()
rows, err := s.db.Query(
`SELECT queue_id, domain, app_login, from_addr, to_addr, subject, status
FROM send_log ORDER BY id`)
if err != nil {
t.Fatalf("query send_log: %v", err)
}
defer rows.Close()
var out []sendLogRow
for rows.Next() {
var r sendLogRow
if err := rows.Scan(&r.QueueID, &r.Domain, &r.AppLogin, &r.From, &r.To, &r.Subject, &r.Status); err != nil {
t.Fatalf("scan: %v", err)
}
out = append(out, r)
}
return out
}
func TestInsertQueuedAndUpdateStatus(t *testing.T) {
st := openTestStore(t)
// Two recipients on the same queue-id → two independent rows (spec 7.3.3).
for _, to := range []string{"a@example.net", "b@example.net"} {
if err := st.InsertQueued(SendLogEntry{
QueueID: "ABC123",
Domain: "example.com",
AppLogin: "app1",
From: "noreply@example.com",
To: to,
Subject: "Hello",
}); err != nil {
t.Fatalf("InsertQueued: %v", err)
}
}
rows := readSendLog(t, st)
if len(rows) != 2 {
t.Fatalf("want 2 rows, got %d: %+v", len(rows), rows)
}
for _, r := range rows {
if r.Status != StatusQueued {
t.Fatalf("new row should be queued, got %q", r.Status)
}
}
// One recipient goes to sent; the other stays queued.
n, err := st.UpdateStatus("ABC123", "a@example.net", StatusSent)
if err != nil {
t.Fatalf("UpdateStatus: %v", err)
}
if n != 1 {
t.Fatalf("want 1 row updated, got %d", n)
}
rows = readSendLog(t, st)
if rows[0].Status != StatusSent || rows[1].Status != StatusQueued {
t.Fatalf("unexpected statuses: %+v", rows)
}
}
func TestUpdateStatusRecipientCaseInsensitive(t *testing.T) {
st := openTestStore(t)
if err := st.InsertQueued(SendLogEntry{QueueID: "Q1", To: "User@Example.NET"}); err != nil {
t.Fatalf("InsertQueued: %v", err)
}
// mail.log may report a differently-cased recipient; matching must still hit.
n, err := st.UpdateStatus("Q1", "user@example.net", StatusBounced)
if err != nil {
t.Fatalf("UpdateStatus: %v", err)
}
if n != 1 {
t.Fatalf("case-insensitive match failed, updated %d rows", n)
}
}
func TestUpdateStatusNoMatch(t *testing.T) {
st := openTestStore(t)
if err := st.InsertQueued(SendLogEntry{QueueID: "Q1", To: "a@example.net"}); err != nil {
t.Fatalf("InsertQueued: %v", err)
}
// A queue-id/recipient the milter never recorded must be a no-op, not an error.
n, err := st.UpdateStatus("Q1", "unknown@example.net", StatusSent)
if err != nil {
t.Fatalf("UpdateStatus: %v", err)
}
if n != 0 {
t.Fatalf("want 0 rows updated, got %d", n)
}
}
func TestDeleteSendLogBefore(t *testing.T) {
st := openTestStore(t)
// Insert one row, then backdate it beyond the retention window by rewriting
// created_at directly (InsertQueued always stamps "now").
if err := st.InsertQueued(SendLogEntry{QueueID: "OLD", To: "a@example.net"}); err != nil {
t.Fatalf("InsertQueued: %v", err)
}
old := time.Now().UTC().AddDate(0, 0, -100).Format(time.RFC3339)
if _, err := st.db.Exec(`UPDATE send_log SET created_at = ? WHERE queue_id = 'OLD'`, old); err != nil {
t.Fatalf("backdate: %v", err)
}
if err := st.InsertQueued(SendLogEntry{QueueID: "NEW", To: "b@example.net"}); err != nil {
t.Fatalf("InsertQueued: %v", err)
}
cutoff := time.Now().UTC().AddDate(0, 0, -90)
n, err := st.DeleteSendLogBefore(cutoff)
if err != nil {
t.Fatalf("DeleteSendLogBefore: %v", err)
}
if n != 1 {
t.Fatalf("want 1 row pruned, got %d", n)
}
rows := readSendLog(t, st)
if len(rows) != 1 || rows[0].QueueID != "NEW" {
t.Fatalf("retention kept wrong rows: %+v", rows)
}
}