From 87416546a599a2a341a873d61dcb1a5ea0b737f7 Mon Sep 17 00:00:00 2001 From: Josh Zhang <39790535+jiashuoz@users.noreply.github.com> Date: Wed, 30 Sep 2026 00:21:36 +0800 Subject: [PATCH 1/3] feat(sending): retention janitors for the sending ledger Nothing deleted the sending-ledger rows every provider-bound send has written since 1.9.0. Add a bounded, batched River periodic (sending_ledger_retention, hourly and on start) that removes: - provider operations with all their attempts 30 days after their stamped expiry, never while an attempt is reserved/authorized or unexpired, a runnable River job names the operation, its customer message is pre-terminal, or a pending notice is bound to it; - budget day counters once the day and the next have closed (today and yesterday always kept, on the gate's own DB clock); - terminal notice event/delivery pairs 30 days after event creation; - control and external-access audit at their stamped expiry, keeping an abuse-class pause event while its account exists (the purge reads it). Batches are LIMIT 1000 FOR UPDATE SKIP LOCKED under a 2s lock_timeout and 30s statement_timeout, capped at 50 per table per run; a failing table is reported and retried next run. Deletions count on e2a_janitor_rows_deleted_total{table}; each run on the new e2a_sending_ledger_retention_runs_total{outcome}. Migrations 126-128 add the missing predicate indexes concurrently. B8 feedback provenance and the policy authority tables are untouched. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_018tVLxUHk3fqQuq8C3wqyHW --- cmd/e2a/main.go | 9 +- cmd/e2a/outbound_wiring.go | 25 +- cmd/e2a/sending_policy_wiring_test.go | 68 ++- docs/data-handling.md | 10 +- docs/observability.md | 3 +- internal/sendingpolicy/export_test.go | 23 + internal/sendingpolicy/ledger_retention.go | 534 ++++++++++++++++++ .../ledger_retention_integration_test.go | 522 +++++++++++++++++ internal/sendingpolicy/maintenance.go | 16 +- internal/sendingpolicy/maintenance_test.go | 4 +- internal/telemetry/metrics.go | 13 + internal/telemetry/prom.go | 19 +- internal/telemetry/prom_test.go | 10 + ...sending_provider_operations_expiry_idx.sql | 21 + .../127_sending_budget_counters_day_idx.sql | 17 + ...ount_sending_control_events_expiry_idx.sql | 18 + 16 files changed, 1292 insertions(+), 20 deletions(-) create mode 100644 internal/sendingpolicy/export_test.go create mode 100644 internal/sendingpolicy/ledger_retention.go create mode 100644 internal/sendingpolicy/ledger_retention_integration_test.go create mode 100644 migrations/126_sending_provider_operations_expiry_idx.sql create mode 100644 migrations/127_sending_budget_counters_day_idx.sql create mode 100644 migrations/128_account_sending_control_events_expiry_idx.sql diff --git a/cmd/e2a/main.go b/cmd/e2a/main.go index 8920757c4..c529ea1e3 100644 --- a/cmd/e2a/main.go +++ b/cmd/e2a/main.go @@ -356,12 +356,9 @@ func main() { metrics = promBackend } store.SetThreadMetrics(metrics) - // External-sending-access decisions (shadow impact / enforce refusals) - // as bounded counters; a no-op while the control is disabled. - sendingpolicy.SetExternalAccessObserver(metrics.ExternalAccessDecision) - // Deletion-resistant feedback ingestion outcomes (B8): bounded - // outcome × bucket counter, no address/account/id labels. - sendingpolicy.SetFeedbackObserver(metrics.SendingFeedbackIngested) + // sendingpolicy's process-wide telemetry hooks: external-access + // decisions, feedback ingestion, and the ledger retention janitor. + installSendingPolicyObservers(metrics) outboxWorker := webhookpub.NewOutboxWorker(pool, store).WithMetrics(metrics) smtpRelay := outbound.NewSMTPRelay(&cfg.OutboundSMTP) sender := outbound.NewSenderWithDKIM(smtpRelay, cfg.OutboundSMTP.FromDomain, store) diff --git a/cmd/e2a/outbound_wiring.go b/cmd/e2a/outbound_wiring.go index 110be93ab..7221b1b09 100644 --- a/cmd/e2a/outbound_wiring.go +++ b/cmd/e2a/outbound_wiring.go @@ -12,6 +12,7 @@ import ( "github.com/tokencanopy/e2a/internal/outbound" "github.com/tokencanopy/e2a/internal/outboundsend" "github.com/tokencanopy/e2a/internal/sendingpolicy" + "github.com/tokencanopy/e2a/internal/telemetry" "github.com/tokencanopy/e2a/internal/webhooknotify" ) @@ -118,13 +119,31 @@ func (s outboundSending) armDeliveryConsumer(c *delivery.Consumer) *delivery.Con return c.WithFeedbackProcessor(s.module) } -// feedbackMaintenance is the retention janitor for feedback provenance and -// daily outcome aggregates. Unregistered, nothing enforces the -// post-deletion horizon. +// feedbackMaintenance is the sending-protection retention registrar: the +// janitor for feedback provenance and daily outcome aggregates, its reconcile +// backstop, and the sending-ledger janitor (operations, attempts, day +// counters, notice outbox, control and access audit). Unregistered, nothing +// enforces any of those horizons and the ledger grows with every send. func (s outboundSending) feedbackMaintenance() *sendingpolicy.MaintenanceJobs { return sendingpolicy.NewMaintenanceJobs(s.module) } +// installSendingPolicyObservers routes sendingpolicy's process-wide telemetry +// hooks to the process metrics backend. Each hook is a silent no-op until set, +// so a missing line here loses a signal without failing anything — which is +// why it is one function the wiring test calls, not three lines in main. +func installSendingPolicyObservers(m telemetry.Metrics) { + // External-sending-access decisions (shadow impact / enforce refusals) + // as bounded counters; a no-op while the control is disabled. + sendingpolicy.SetExternalAccessObserver(m.ExternalAccessDecision) + // Deletion-resistant feedback ingestion outcomes (B8): bounded + // outcome × bucket counter, no address/account/id labels. + sendingpolicy.SetFeedbackObserver(m.SendingFeedbackIngested) + // Ledger retention: rows deleted per ledger table (on the shared janitor + // counter) and one run-outcome sample per pass. + sendingpolicy.SetLedgerRetentionObserver(m) +} + // nonEmpty returns the non-blank values, so an unset config string does not // become an empty-domain entry. func nonEmpty(values ...string) []string { diff --git a/cmd/e2a/sending_policy_wiring_test.go b/cmd/e2a/sending_policy_wiring_test.go index 7035e9b64..a7e0ddfbc 100644 --- a/cmd/e2a/sending_policy_wiring_test.go +++ b/cmd/e2a/sending_policy_wiring_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "strings" + "sync" "testing" "github.com/riverqueue/river" @@ -11,8 +12,10 @@ import ( "github.com/tokencanopy/e2a/internal/agent" "github.com/tokencanopy/e2a/internal/config" "github.com/tokencanopy/e2a/internal/delivery" + "github.com/tokencanopy/e2a/internal/jobs" "github.com/tokencanopy/e2a/internal/outbound" "github.com/tokencanopy/e2a/internal/sendingpolicy" + "github.com/tokencanopy/e2a/internal/telemetry" "github.com/tokencanopy/e2a/internal/testutil/testdb" "github.com/tokencanopy/e2a/internal/usage" ) @@ -166,7 +169,68 @@ func TestFeedbackAccountingWiring(t *testing.T) { t.Fatal("no feedback retention janitor composed") } periodics := janitor.RegisterJobs(river.NewWorkers()) - if len(periodics) != 2 { - t.Fatalf("retention periodics = %d, want 2 (retention pass + reconcile)", len(periodics)) + if len(periodics) != 3 { + t.Fatalf("retention periodics = %d, want 3 (feedback retention + reconcile + ledger retention)", len(periodics)) + } +} + +// ledgerObserverRecorder is a full telemetry backend that records only the +// ledger janitor's samples. +type ledgerObserverRecorder struct { + telemetry.NoOp + mu sync.Mutex + deleted map[string]int + runs []string +} + +func (r *ledgerObserverRecorder) JanitorRowsDeleted(table string, count int) { + r.mu.Lock() + defer r.mu.Unlock() + r.deleted[table] += count +} + +func (r *ledgerObserverRecorder) SendingLedgerRetentionRun(outcome string) { + r.mu.Lock() + defer r.mu.Unlock() + r.runs = append(r.runs, outcome) +} + +// TestLedgerRetentionWiring: the composition root must both register the +// sending-ledger janitor and route its samples to the process metrics +// backend. Either omission is silent in production — the ledger simply grows, +// or shrinks with no signal — so this drives the composed module's janitor +// over a seeded expired row after installing the observers exactly as main +// does, and requires the deletion and the run outcome to reach the backend. +func TestLedgerRetentionWiring(t *testing.T) { + pool := testdb.TestDB(t) + if err := jobs.Migrate(context.Background(), pool); err != nil { + t.Fatalf("river migrate: %v", err) + } + relay := outbound.NewSMTPRelay(&config.OutboundSMTPConfig{Host: "relay.invalid", Port: 587, FromDomain: "test.e2a.dev"}) + composed := newOutboundSending(outboundSendingDeps{ + pool: pool, + relay: relay, + secrets: sendingpolicy.Secrets{}, + source: sendingpolicy.PolicySourceConfig, + policy: sendingpolicy.DisabledPolicy(), + }) + + rec := &ledgerObserverRecorder{deleted: map[string]int{}} + installSendingPolicyObservers(rec) + t.Cleanup(func() { installSendingPolicyObservers(telemetry.NoOp{}) }) + + if _, err := pool.Exec(context.Background(), ` + INSERT INTO sending_budget_counters (scope, scope_id, day, daily_limit) + VALUES ('global_all', 'global', (now() AT TIME ZONE 'UTC')::date - 3, 5000)`); err != nil { + t.Fatal(err) + } + if err := sendingpolicy.NewLedgerRetentionWorker(composed.module).Work(context.Background(), &river.Job[sendingpolicy.LedgerRetentionArgs]{}); err != nil { + t.Fatalf("Work: %v", err) + } + if rec.deleted[sendingpolicy.LedgerTableCounters] != 1 { + t.Fatalf("deleted samples = %v, want the closed-day counter reported", rec.deleted) + } + if len(rec.runs) != 1 || rec.runs[0] != sendingpolicy.LedgerRunComplete { + t.Fatalf("run samples = %v, want one complete", rec.runs) } } diff --git a/docs/data-handling.md b/docs/data-handling.md index 649671469..2cc62b24d 100644 --- a/docs/data-handling.md +++ b/docs/data-handling.md @@ -21,7 +21,13 @@ For vulnerability reporting and the security model, see [SECURITY.md](../SECURIT | Per-webhook signing secret (`whsec_…`) | Postgres, **plaintext** | Until the webhook is deleted. Returned once at creation; rotate via `POST /v1/webhooks/{id}/rotate-secret` (the previous secret stays valid for a 24h grace window). | | Owner-mailbox proof (the sign-in address a verified Google login confirmed, its source and timestamp) | Postgres `users.owner_email_verified_*` | Until the account is deleted; stops applying as soon as the account email changes. | | External sending access requests (use case, intended recipients, expected volume, review state) | Postgres `external_sending_access_requests` | Until the account is deleted (cascade). Private support data — never filed to a public tracker. | -| External sending access audit (operator grant/revoke: account id, old/new state, actor, reason; no addresses) | Postgres `external_sending_access_events` | Append-only; no account foreign key, so it outlives account deletion. Each row carries an `expires_at` (the `sending_control_audit_retention_days` policy value, 90 days by default); like `account_sending_control_events`, the purge of expired rows is not yet automated and is an operator task. | +| External sending access audit (operator grant/revoke: account id, old/new state, actor, reason; no addresses) | Postgres `external_sending_access_events` | Append-only; no account foreign key, so it outlives account deletion. Each row carries an `expires_at` (the `sending_control_audit_retention_days` policy value, 90 days by default) and the sending-ledger janitor deletes it then. | +| Sending pause/resume audit (account id, old/new state, reason, actor, pause class, optional operator evidence reference; no addresses) | Postgres `account_sending_control_events` | No account foreign key, so it outlives account deletion. Deleted by the sending-ledger janitor at its `expires_at` (`sending_control_audit_retention_days`, 90 days by default) — except an `abuse`-class event, which is kept past its expiry while the account still exists (live or trashed), because the account purge reads it to decide the abuse hold; it is deleted after the purge. | +| Sending ledger: provider operations and their submission attempts (opaque operation/account ids, purpose, shared-reputation flag, attempt ordinal, UTC day, recipient count, budget/call state, an HMAC commitment to a notice recipient; no addresses, subjects or bodies) | Postgres `sending_provider_operations`, `sending_budget_reservations` | No foreign keys, so message and account deletion never touch them. 30 days: each row is stamped with an `expires_at` at creation (an operation 30 days after creation or its scheduled send time, an attempt 30 days after it was first written), and the sending-ledger janitor deletes an operation together with all of its attempts once every one of them is past its expiry — but never while the send can still use it: an attempt still holding capacity (`reserved`) or holding an unredeemed provider authorization, a runnable River job that names the operation (a budget or pause hold, a retry), a customer message not yet handed to the provider, or a pending protection notice keeps it. | +| Sending feedback provenance (opaque operation/account ids, random correlation id, provider message id, keyed recipient HMACs, detector bucket, timestamps; no addresses, subjects or bodies) | Postgres `sending_feedback_correlations`, `sending_feedback_recipients`, `sending_feedback_events` | For a customer account's lifetime, then 30 days after the account is purged (`sending_feedback_post_account_retention_days`). Synthetic traffic — non-customer purposes, and customer-purpose mail from `system`/`internal` accounts such as probers and monitors, which are never deleted — keeps it for 30 days from creation. Removed by the hourly feedback retention janitor. | +| Sending budget day counters (scope, opaque account id or `global`, UTC day, reserved/confirmed counts, limit) | Postgres `sending_budget_counters` | Deleted by the sending-ledger janitor once the day and the day after it have both closed (the current and previous UTC day are always kept). Only closed days are deleted, so deletion never gives back capacity. | +| Sending protection notice outbox (opaque event/account/operation ids, closed kind/reason/scope/audience values, delivery state; no addresses) | Postgres `sending_protection_notice_events`, `sending_protection_notice_deliveries` | Survive user, message and agent deletion (their only foreign key is delivery → event). A finished event (every delivery sent, failed or skipped) is deleted with its deliveries 30 days after the event was created; an event with a pending delivery is never deleted by age. | +| Sending policy authority (runtime policy and its change events, runtime attestation and its events, operator-recipient registry, ramp grandfathering marker) | Postgres `sending_protection_runtime_policy`, `sending_protection_policy_events`, `sending_protection_runtime_attestation`, `sending_protection_runtime_attestation_events`, `sending_operator_recipient_versions`, `sending_ramp_grandfathering` | Retained indefinitely: append-only security/audit state with no customer content (the operator-recipient registry holds only keyed commitments, and its trigger refuses deletion). No janitor touches it. | | Trashed accounts (the whole account row and everything it owns, inert) | Postgres `users.deleted_at` and the owned rows | Until restored, or purged after `trash.account_retention_days` (default = `trash.retention_days`, 30 days). `0` disables account trash. | | Identity tombstones (keyed HMAC digests of a purged account's login subjects, email and — after an abuse pause — verified domains; no plaintext) | Postgres `identity_tombstones` | Only when `trash.identity_tombstones` is on. At least 30 days (longest of the trash window and the rest of the quota month); 2 years for an abuse-paused account. Swept by the janitor at `expires_at`. | | Deleted-account summaries (counts, recipient-domain histogram, per-day send counts, keyed subject/domain digests, pause state and optional operator evidence reference; no bodies, recipient addresses or owner email) | Postgres `deleted_account_summaries` | Only when `trash.identity_tombstones` is on. Same hold as the account's tombstones; swept by the janitor. Operator-readable only (`-inspect-deleted-account`). | @@ -44,7 +50,7 @@ The API exposes self-service export and deletion operations that support GDPR Ar - **`DELETE /v1/account?confirm=DELETE`** — moves the account to the **trash** (the default). The account becomes unusable at once — every API key, OAuth grant and dashboard session is revoked, every agent is trashed (inbound refused), sending stops (the sending gate refuses a trashed account; an existing abuse pause and its class are left untouched), and every custom domain loses its verification and its SES identity is torn down — while its data stays in Postgres, restorable by signing in to the dashboard, until the janitor purges it after `trash.account_retention_days` (default: `trash.retention_days`, 30 days). A restore brings the account and its agents back; API keys stay revoked and domains must be re-verified. The receipt is `mode: "trash"` with `purge_after`; `messages_deleted` is 0 and the other counts describe rows trashed, revoked or unverified. - **`DELETE /v1/account?confirm=DELETE&permanent=true`** (and "erase now" in the restore interstitial) — erases immediately, except while the account's sending is paused (any pause class), which answers `409 erase_held` (the trash path stays available, and permanent agent deletion is held the same way): the account's agents are drained in bounded chunks and then the account row and every related row are deleted (cascade through the account-owned suppression/token tables and the contacts/engagements/import-batches/templates tables, plus explicit deletion of `usage_events`, whose FK is `ON DELETE SET NULL`). The janitor's purge after the trash window runs exactly the same path. The receipt is `mode: "permanent"` with `user_deleted: true` and selected per-table counts (including `agent_suppressions_deleted` and `agent_unsubscribe_tokens_deleted`); it does not separately count contacts, engagements, import batches, templates, webhooks, or webhook event/delivery rows, even though those rows are deleted. A deployment that sets `trash.account_retention_days: 0` makes every account deletion permanent. - **Deferred erase for recent senders.** A permanent erase of an account that emailed an external recipient (anyone other than its own agents, its verified owner mailbox, or an address on the deployment's shared agent domain) within `trash.recent_sender_erase_defer_days` (default 14; `0` disables) is not performed at once: the account is moved to the trash exactly as by the default delete and purged by the janitor at the end of the account trash window. The receipt is `mode: "trash"` with `erase_deferred: true` and `purge_after`. This keeps the account's sending controls and bounce/complaint aggregates alive while late provider feedback arrives. The account is restorable until `purge_after` like any trashed account; the paused-account `409 erase_held` still takes precedence. No effect when account trash is disabled. The same rule applies one level down: `DELETE /v1/agents/{email}?permanent=true` for an agent that sent externally inside the window moves it to (or keeps it in) the agent trash, and a permanent delete of a trashed message that was sent externally inside the window leaves it in the message trash — both receipts carry `erase_deferred: true` and `purge_after`, and the janitor purges them at the end of the ordinary trash window (`trash.retention_days`). So the sent-mail records the account check reads cannot be removed on demand inside the window. A deferred agent's pending scheduled sends are canceled. System and internal account classes, the shared agent domain and the provider test domains in `trash.erase_defer_exempt_domains` (default `simulator.amazonses.com`) are exempt. Deferred items keep counting toward storage, and a deferred agent keeps its address and blocks domain deletion, until purged; operators can force the purge by backdating `deleted_at`. -- **What outlives a purge.** The sending ledger (provider operations, budget counters, feedback provenance, control audit) keeps its own retention, as before. When the deployment enables `trash.identity_tombstones` (hosted policy; off by default) every purge also writes, before any row is deleted: **identity tombstones** — HMAC-SHA256 digests under a dedicated key (`E2A_TOMBSTONE_KEY`) of the account's login subject(s) and normalized email, held for the longer of the trash window and the rest of the current quota month (at least 30 days), plus — with a 2-year hold when the account was ever paused for abuse (the current pause class or any `abuse` event in the control audit) — its verified domains, every agent address it held on the shared domain, and its email domain unless that is a public webmail provider — so that identity cannot immediately register again (`registration_refused`; domains answer `domain_taken`); and a **deleted-account summary** (`deleted_account_summaries`) — counts, first/last send, a top-20 recipient-*domain* histogram, per-day send counts for the last 30 days, keyed digests of the top-20 subjects, of verified domains and of every identifier an abuse hold would close (so an operator can escalate a purged account to an abuse hold later with `-escalate-deleted-account-to-abuse`), and the pause state/class with an optional operator evidence reference (never the free-text pause reason); send counts are taken from `usage_events` as well as messages, so permanently deleting messages first does not erase them — kept for the same hold. Neither holds message bodies, recipient addresses, or the owner email in the clear; the summary is readable only through the operator command `-inspect-deleted-account`, never over HTTP or MCP. The janitor deletes both at `expires_at`. +- **What outlives a purge.** The sending ledger (provider operations, budget counters, feedback provenance, control audit, notice outbox) keeps its own retention, listed in the table above. When the deployment enables `trash.identity_tombstones` (hosted policy; off by default) every purge also writes, before any row is deleted: **identity tombstones** — HMAC-SHA256 digests under a dedicated key (`E2A_TOMBSTONE_KEY`) of the account's login subject(s) and normalized email, held for the longer of the trash window and the rest of the current quota month (at least 30 days), plus — with a 2-year hold when the account was ever paused for abuse (the current pause class or any `abuse` event in the control audit) — its verified domains, every agent address it held on the shared domain, and its email domain unless that is a public webmail provider — so that identity cannot immediately register again (`registration_refused`; domains answer `domain_taken`); and a **deleted-account summary** (`deleted_account_summaries`) — counts, first/last send, a top-20 recipient-*domain* histogram, per-day send counts for the last 30 days, keyed digests of the top-20 subjects, of verified domains and of every identifier an abuse hold would close (so an operator can escalate a purged account to an abuse hold later with `-escalate-deleted-account-to-abuse`), and the pause state/class with an optional operator evidence reference (never the free-text pause reason); send counts are taken from `usage_events` as well as messages, so permanently deleting messages first does not erase them — kept for the same hold. Neither holds message bodies, recipient addresses, or the owner email in the clear; the summary is readable only through the operator command `-inspect-deleted-account`, never over HTTP or MCP. The janitor deletes both at `expires_at`. Both are scoped to the authenticated user — there's no path to target someone else's data. diff --git a/docs/observability.md b/docs/observability.md index 7152b9f31..9ee1f55e5 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -197,7 +197,8 @@ addresses, subjects, message content, or RFC Message-IDs. | Metric | Type | Labels | Meaning | |---|---|---|---| -| `e2a_janitor_rows_deleted_total` | counter | `table` | TTL sweep deletions. | +| `e2a_janitor_rows_deleted_total` | counter | `table` | TTL sweep deletions. The sending-ledger janitor (`internal/sendingpolicy`, River periodic `sending_ledger_retention`, hourly and on start) reports here too, with `table` ∈ `sending_provider_operations`, `sending_budget_reservations`, `sending_budget_counters`, `sending_protection_notice_events`, `sending_protection_notice_deliveries`, `account_sending_control_events`, `external_sending_access_events` (retention rules: [data handling](data-handling.md)). Table names only — never an operation, account or event id. | +| `e2a_sending_ledger_retention_runs_total` | counter | `outcome` | One sample per sending-ledger janitor pass. `complete`: every table reached its horizon. `partial`: a table hit its per-run cap (50 batches × 1,000 rows) and continues next run — expected for the first runs after rollout on a deployment with weeks of ledger backlog. `failed`: a table's batch hit its 2 s `lock_timeout` or 30 s `statement_timeout` (or lost its connection); the other tables still ran and the next hourly run retries. Alert on sustained `failed` (several consecutive hours), not on one sample; a steady `partial` that never reaches `complete` means the backlog outgrows the per-run cap. Each run also logs one `[sendingpolicy:ledger-gc]` line with per-table counts and `retained_expired_operations` — expired operations a guard is keeping (an attempt still holding capacity, a live job, an in-flight message, a pending notice); a number that only grows points at a stuck attempt or job. | | `e2a_notify_missed_total` | counter | — | Fallback-poll wakeups LISTEN/NOTIFY missed (reconnect churn indicator). | | `e2a_redeliver_requests_total` | counter | `scope` | Customer-driven webhook replays. | | `e2a_contact_due_events_total` | counter | `outcome` | `contact.due` sweep outcomes. `published` counts wake-ups committed to the durable outbox. `failed` counts a failed atomic claim/publish attempt (or the affected batch size when known); the transaction rolls back and River retries the job. A sustained zero `published` despite known overdue, armed engagements can indicate a stalled sweep. Alert on sustained `failed` growth or exhausted River jobs rather than treating one increment as a permanent miss. | diff --git a/internal/sendingpolicy/export_test.go b/internal/sendingpolicy/export_test.go new file mode 100644 index 000000000..13178c9a9 --- /dev/null +++ b/internal/sendingpolicy/export_test.go @@ -0,0 +1,23 @@ +package sendingpolicy + +import ( + "context" + "time" +) + +// LedgerTuning is the test-only batch geometry and fault hook for GCLedger. +type LedgerTuning struct { + BatchSize int + MaxBatches int + // BeforeBatch runs inside each batch transaction, after its timeouts are + // set and before the delete; exec runs SQL in that transaction. + BeforeBatch func(ctx context.Context, table string, exec func(context.Context, string) error) error +} + +// GCLedgerTuned runs GCLedger with a test batch geometry. +func (m *Module) GCLedgerTuned(ctx context.Context, now time.Time, t LedgerTuning) (LedgerRetentionStats, error) { + return m.gcLedger(ctx, now, ledgerRetentionTuning{batchSize: t.BatchSize, maxBatches: t.MaxBatches, beforeBatch: t.BeforeBatch}) +} + +// LedgerBatchSize exposes the production batch size. +const LedgerBatchSize = ledgerBatchSize diff --git a/internal/sendingpolicy/ledger_retention.go b/internal/sendingpolicy/ledger_retention.go new file mode 100644 index 000000000..cc1a788d4 --- /dev/null +++ b/internal/sendingpolicy/ledger_retention.go @@ -0,0 +1,534 @@ +package sendingpolicy + +import ( + "context" + "errors" + "fmt" + "log" + "strings" + "sync/atomic" + "time" + + "github.com/riverqueue/river" +) + +// This file is the retention janitor for the sending ledger: the tables the +// gate writes on every provider-bound send (operations, attempt reservations, +// day counters) and the protection audit/outbox tables beside them. Until it +// existed nothing deleted any of them. +// +// It deliberately does NOT touch the B8 feedback provenance (correlations, +// recipient HMACs, feedback events) or the daily outcome aggregates — GCFeedback +// owns those — and it never deletes the append-only policy authority +// (runtime policy, its events, the runtime attestation and its events, the +// operator-recipient registry, the ramp grandfathering marker). Those rows are +// few, security-relevant, and have no retention horizon. +// +// Retention rules and their sources (spec 2026-08-19-sending-abuse-prevention +// §5.1 "Reservation semantics", §5.2, §9, §10): +// +// - sending_budget_counters: "A retention janitor removes day-counter rows +// after two closed UTC days". A row for day D is deleted once D and D+1 +// have both closed, i.e. day <= today-2 on the ledger's own clock. Today +// and yesterday always survive. +// - sending_provider_operations + sending_budget_reservations: "provider- +// operation/reservation rows after 30 days", measured from the expires_at +// each row was stamped with at creation (operationTTL; migration 113's +// reservation default). An operation goes together with every attempt it +// owns, and only once every one of those attempts is past its own expiry. +// - sending_protection_notice_events/_deliveries: "sent, failed, and skipped +// rows remain for 30 days from event creation and then the event/delivery +// pair is removed". A pending delivery is never deleted by age. +// - account_sending_control_events: 90-day security retention, stamped as +// expires_at from the effective policy's sending_control_audit_retention_days +// when the event is written. +// - external_sending_access_events: the same control-audit retention, stamped +// the same way (migration 121: "reaped only by its own expires_at"). + +// ledgerCounterAgeDays: a counter for day D is deleted once day <= today - 2. +const ledgerCounterAgeDays = 2 + +// noticeTerminalRetention is how long a finished notice event and its +// deliveries are kept after the event was created (spec §5.2, §10). +const noticeTerminalRetention = 30 * 24 * time.Hour + +// ledgerBatchSize and ledgerMaxBatches bound one run: every batch is its own +// short transaction, and a table stops after ledgerMaxBatches so that the +// first run over weeks of production backlog neither holds a transaction open +// for long nor starves the tables after it. The next hourly run resumes. +const ( + ledgerBatchSize = 1000 + ledgerMaxBatches = 50 +) + +// ledgerLockTimeout and ledgerStatementTimeout bound one batch. SKIP LOCKED +// already steps around any row a sender holds; lock_timeout covers the +// remaining waits (a relation lock, a row a concurrent writer locked after the +// victim scan) so the janitor fails fast rather than queueing behind — or in +// front of — the send path. A timeout fails only that table for this run. +const ( + ledgerLockTimeout = "2s" + ledgerStatementTimeout = "30s" +) + +// ledgerRetentionInterval paces the janitor. Nothing it removes is urgent. +const ledgerRetentionInterval = time.Hour + +// Ledger retention table labels. They double as the metric's closed `table` +// label set (e2a_janitor_rows_deleted_total), so they carry no identifiers. +const ( + LedgerTableOperations = "sending_provider_operations" + LedgerTableReservations = "sending_budget_reservations" + LedgerTableCounters = "sending_budget_counters" + LedgerTableNoticeEvents = "sending_protection_notice_events" + LedgerTableNoticeDeliveries = "sending_protection_notice_deliveries" + LedgerTableControlEvents = "account_sending_control_events" + LedgerTableAccessEvents = "external_sending_access_events" +) + +// Run outcomes for the e2a_sending_ledger_retention_runs_total metric. +const ( + // LedgerRunComplete: every table was drained to its horizon. + LedgerRunComplete = "complete" + // LedgerRunPartial: at least one table hit its per-run batch cap; the + // backlog continues next run. Expected on the first runs after rollout. + LedgerRunPartial = "partial" + // LedgerRunFailed: at least one table's batch failed (lock or statement + // timeout, connection loss). Non-fatal: the next run retries. + LedgerRunFailed = "failed" +) + +// LedgerRetentionObserver receives the janitor's bounded samples. The process +// telemetry backend (telemetry.Metrics) satisfies it. Never carries an id. +type LedgerRetentionObserver interface { + JanitorRowsDeleted(table string, count int) + SendingLedgerRetentionRun(outcome string) +} + +var ledgerObserver atomic.Value // ledgerObserverBox + +type ledgerObserverBox struct{ o LedgerRetentionObserver } + +// SetLedgerRetentionObserver installs the process-wide janitor observer. Set +// once at startup by the composition root; nil disables it. +func SetLedgerRetentionObserver(o LedgerRetentionObserver) { + ledgerObserver.Store(ledgerObserverBox{o}) +} + +func currentLedgerObserver() LedgerRetentionObserver { + if b, ok := ledgerObserver.Load().(ledgerObserverBox); ok { + return b.o + } + return nil +} + +// LedgerRetentionStats reports one run. Deleted is keyed by table label. +type LedgerRetentionStats struct { + Deleted map[string]int64 + // Capped names the tables that hit the per-run batch cap. + Capped []string + // Failed maps a table label to the error that stopped it this run. + Failed map[string]error + // RetainedOperations counts expired operations kept because an attempt, + // a live job, an in-flight message, or a pending notice still needs them + // (capped at retainedCountCap). Diagnostic only. + RetainedOperations int64 +} + +// Outcome classifies the run for the run-outcome metric. +func (s LedgerRetentionStats) Outcome() string { + switch { + case len(s.Failed) > 0: + return LedgerRunFailed + case len(s.Capped) > 0: + return LedgerRunPartial + default: + return LedgerRunComplete + } +} + +// Err joins the per-table failures, nil when there were none. +func (s LedgerRetentionStats) Err() error { + if len(s.Failed) == 0 { + return nil + } + errs := make([]error, 0, len(s.Failed)) + for _, table := range ledgerTableOrder { + if err, ok := s.Failed[table]; ok { + errs = append(errs, err) + } + } + return errors.Join(errs...) +} + +// retainedCountCap bounds the diagnostic count query. +const retainedCountCap = 10000 + +// ledgerSweep is one bounded, batched delete. The statement must return +// (primary rows deleted, secondary rows deleted); the loop stops when a batch +// deletes fewer primary rows than the batch size. +type ledgerSweep struct { + table string + secondary string // label for the second count, "" when unused + sql string + args func(now time.Time) []any +} + +// ledgerTableOrder is the order sweeps run and failures are reported in. +// Notice pairs go first (cheap, tiny); operations before counters so the +// heaviest table gets the freshest budget of the job timeout. +var ledgerTableOrder = []string{ + LedgerTableNoticeEvents, LedgerTableOperations, LedgerTableCounters, + LedgerTableControlEvents, LedgerTableAccessEvents, +} + +// liveJobStates are the River states in which a job may still run and +// therefore still dereference its operation_ref. Everything else +// (completed, cancelled, discarded) is final. +const liveJobStates = `'available', 'pending', 'retryable', 'running', 'scheduled'` + +var ledgerSweeps = map[string]ledgerSweep{ + // Terminal notice pairs, 30 days after event creation. A pending delivery + // keeps its event (and, through the operation guard below, its operation) + // however old it is: the seven-day delivery deadline, not the janitor, is + // what ends a pending delivery. The FK cascade would remove the deliveries + // too; they are deleted explicitly so they can be counted. + LedgerTableNoticeEvents: { + table: LedgerTableNoticeEvents, + secondary: LedgerTableNoticeDeliveries, + sql: ` + WITH victims AS ( + SELECT e.id + FROM sending_protection_notice_events e + WHERE e.created_at <= $1 + AND NOT EXISTS ( + SELECT 1 FROM sending_protection_notice_deliveries d + WHERE d.event_id = e.id AND d.state = 'pending') + ORDER BY e.created_at, e.id + LIMIT $2 + FOR UPDATE OF e SKIP LOCKED + ), gone_deliveries AS ( + DELETE FROM sending_protection_notice_deliveries d + USING victims v WHERE d.event_id = v.id + RETURNING 1 + ), gone_events AS ( + DELETE FROM sending_protection_notice_events e + USING victims v WHERE e.id = v.id + RETURNING 1 + ) + SELECT (SELECT count(*) FROM gone_events), (SELECT count(*) FROM gone_deliveries)`, + args: func(now time.Time) []any { return []any{now.Add(-noticeTerminalRetention)} }, + }, + + // Provider operations with every attempt they own. An operation survives + // its own expiry while ANY of these still needs it: + // - an attempt that is not terminal (state reserved: capacity held; or + // call_state authorized: a token was issued and not yet redeemed or + // invalidated), or an attempt not yet past its own expiry; + // - a River job that can still run and names it in operation_ref — + // this covers budget and pause holds (snoozed jobs), retries inside + // SendRetryHorizon, and every notification kind; deleting it would + // make the worker fail closed (ErrSourceUnavailable) and drop the send; + // - its customer message is still pre-terminal (accepted/queued/sending, + // or pending review), which the terminal reconciler may still settle; + // - a pending protection-notice delivery bound to it. + // Feedback correlations are deliberately NOT a guard: they copy the + // operation's attribution at authorization and no reader joins back to + // this table (B8 is self-contained), so a correlation's account-lifetime + // retention never extends an operation's 30 days. + // The live-job set is small (the runnable backlog), so it is built once + // per batch and probed as a hashed NOT IN rather than re-scanned per + // candidate. SKIP LOCKED steps around an operation a worker holds FOR UPDATE; a + // worker that locks it after this statement's lock waits and then sees + // the row gone, which only happens for an operation no live job names. + LedgerTableOperations: { + table: LedgerTableOperations, + secondary: LedgerTableReservations, + sql: ` + WITH live_refs AS MATERIALIZED ( + SELECT DISTINCT j.args->'operation_ref'->>'id' AS operation_id + FROM river_job j + WHERE j.state IN (` + liveJobStates + `) + AND j.args ? 'operation_ref' + ), victims AS ( + SELECT o.operation_id + FROM sending_provider_operations o + WHERE o.expires_at <= $1 + AND NOT EXISTS ( + SELECT 1 FROM sending_budget_reservations r + WHERE r.operation_id = o.operation_id + AND (r.expires_at > $1 OR r.state = 'reserved' OR r.call_state = 'authorized')) + AND o.operation_id NOT IN ( + SELECT l.operation_id FROM live_refs l WHERE l.operation_id IS NOT NULL) + AND NOT EXISTS ( + SELECT 1 FROM messages m + WHERE m.id = o.operation_id AND m.direction = 'outbound' + AND (m.delivery_status IN ('accepted', 'queued', 'sending') + OR m.status = 'pending_review')) + AND NOT EXISTS ( + SELECT 1 FROM sending_protection_notice_deliveries d + WHERE d.current_operation_id = o.operation_id AND d.state = 'pending') + ORDER BY o.expires_at, o.operation_id + LIMIT $2 + FOR UPDATE OF o SKIP LOCKED + ), gone_attempts AS ( + DELETE FROM sending_budget_reservations r + USING victims v WHERE r.operation_id = v.operation_id + RETURNING 1 + ), gone_operations AS ( + DELETE FROM sending_provider_operations o + USING victims v WHERE o.operation_id = v.operation_id + RETURNING 1 + ) + SELECT (SELECT count(*) FROM gone_operations), (SELECT count(*) FROM gone_attempts)`, + args: func(now time.Time) []any { return []any{now} }, + }, + + // Day counters for closed days, on the same clock the gate stamps them + // with (ledgerDay: the database's UTC date). Only past days are touched, + // so deletion can never refund capacity: today's row is never a victim, + // and the one path that writes an older row — releasing a stale + // reservation — skips a missing counter by design (ledgerPlan.release). + LedgerTableCounters: { + table: LedgerTableCounters, + sql: ` + WITH victims AS ( + SELECT scope, scope_id, day + FROM sending_budget_counters + WHERE day <= (clock_timestamp() AT TIME ZONE 'UTC')::date - $1::int + ORDER BY day + LIMIT $2 + FOR UPDATE SKIP LOCKED + ), gone AS ( + DELETE FROM sending_budget_counters c + USING victims v + WHERE c.scope = v.scope AND c.scope_id = v.scope_id AND c.day = v.day + RETURNING 1 + ) + SELECT (SELECT count(*) FROM gone), 0`, + args: func(time.Time) []any { return []any{ledgerCounterAgeDays} }, + }, + + // Pause audit past its stamped expiry. One exception: an event recording + // an ABUSE-class pause survives while its account still exists (live or + // in the trash). The account purge derives the abuse hold from this + // history ("ever paused as abuse", identity.writePurgeRecordsTx); deleting + // it at 90 days would let an account resumed after an abuse pause wait out + // the audit and then be purged without abuse tombstones. After the purge + // has read it, the event goes at its (already passed) expiry. + LedgerTableControlEvents: { + table: LedgerTableControlEvents, + sql: ` + WITH victims AS ( + SELECT e.id + FROM account_sending_control_events e + WHERE e.expires_at <= $1 + AND (e.pause_class IS DISTINCT FROM 'abuse' + OR NOT EXISTS (SELECT 1 FROM users u WHERE u.id = e.account_ref)) + ORDER BY e.expires_at, e.id + LIMIT $2 + FOR UPDATE OF e SKIP LOCKED + ), gone AS ( + DELETE FROM account_sending_control_events e + USING victims v WHERE e.id = v.id + RETURNING 1 + ) + SELECT (SELECT count(*) FROM gone), 0`, + args: func(now time.Time) []any { return []any{now} }, + }, + + // Operator grant/revoke audit past its stamped expiry. Nothing reads it + // for a later transition: the grant itself lives on + // account_sending_controls. + LedgerTableAccessEvents: { + table: LedgerTableAccessEvents, + sql: ` + WITH victims AS ( + SELECT id + FROM external_sending_access_events + WHERE expires_at <= $1 + ORDER BY expires_at, id + LIMIT $2 + FOR UPDATE SKIP LOCKED + ), gone AS ( + DELETE FROM external_sending_access_events e + USING victims v WHERE e.id = v.id + RETURNING 1 + ) + SELECT (SELECT count(*) FROM gone), 0`, + args: func(now time.Time) []any { return []any{now} }, + }, +} + +// ledgerRetentionTuning lets tests shrink the batch geometry; production uses +// the constants. +type ledgerRetentionTuning struct { + batchSize int + maxBatches int + // beforeBatch, when set, runs inside each batch transaction after the + // timeouts are set and before the delete. Test-only fault injection. + beforeBatch func(ctx context.Context, table string, exec func(context.Context, string) error) error +} + +func defaultLedgerTuning() ledgerRetentionTuning { + return ledgerRetentionTuning{batchSize: ledgerBatchSize, maxBatches: ledgerMaxBatches} +} + +// GCLedger runs one retention pass over the sending ledger. It never returns +// early on a table's failure: each table is swept independently and failures +// are reported in the stats, so one lock timeout cannot stall the others. +// The returned error is non-nil only for a cancelled context. +func (m *Module) GCLedger(ctx context.Context, now time.Time) (LedgerRetentionStats, error) { + return m.gcLedger(ctx, now, defaultLedgerTuning()) +} + +func (m *Module) gcLedger(ctx context.Context, now time.Time, tuning ledgerRetentionTuning) (LedgerRetentionStats, error) { + now = now.UTC() + st := LedgerRetentionStats{Deleted: map[string]int64{}, Failed: map[string]error{}} + for _, table := range ledgerTableOrder { + if err := ctx.Err(); err != nil { + return st, err + } + sweep := ledgerSweeps[table] + capped, err := m.runLedgerSweep(ctx, sweep, now, tuning, &st) + if err != nil { + st.Failed[table] = err + continue + } + if capped { + st.Capped = append(st.Capped, table) + } + } + if n, err := m.countRetainedOperations(ctx, now); err == nil { + st.RetainedOperations = n + } + return st, nil +} + +func (m *Module) runLedgerSweep(ctx context.Context, sweep ledgerSweep, now time.Time, tuning ledgerRetentionTuning, st *LedgerRetentionStats) (capped bool, err error) { + args := append(sweep.args(now), tuning.batchSize) + for batch := 0; batch < tuning.maxBatches; batch++ { + primary, secondary, err := m.runLedgerBatch(ctx, sweep, args, tuning) + st.Deleted[sweep.table] += primary + if sweep.secondary != "" { + st.Deleted[sweep.secondary] += secondary + } + if err != nil { + return false, err + } + if primary < int64(tuning.batchSize) { + return false, nil + } + } + return true, nil +} + +func (m *Module) runLedgerBatch(ctx context.Context, sweep ledgerSweep, args []any, tuning ledgerRetentionTuning) (primary, secondary int64, err error) { + tx, err := m.pool.Begin(ctx) + if err != nil { + return 0, 0, fmt.Errorf("sendingpolicy: ledger retention %s: begin: %w", sweep.table, err) + } + defer func() { _ = tx.Rollback(ctx) }() + if _, err := tx.Exec(ctx, `SET LOCAL lock_timeout = '`+ledgerLockTimeout+`'`); err != nil { + return 0, 0, fmt.Errorf("sendingpolicy: ledger retention %s: lock_timeout: %w", sweep.table, err) + } + if _, err := tx.Exec(ctx, `SET LOCAL statement_timeout = '`+ledgerStatementTimeout+`'`); err != nil { + return 0, 0, fmt.Errorf("sendingpolicy: ledger retention %s: statement_timeout: %w", sweep.table, err) + } + if tuning.beforeBatch != nil { + exec := func(ctx context.Context, sql string) error { + _, err := tx.Exec(ctx, sql) + return err + } + if err := tuning.beforeBatch(ctx, sweep.table, exec); err != nil { + return 0, 0, fmt.Errorf("sendingpolicy: ledger retention %s: %w", sweep.table, err) + } + } + // The outer SELECT reads only the CTEs' RETURNING output, never a table + // a CTE modified, so the statement-snapshot hazard does not apply. + if err := tx.QueryRow(ctx, sweep.sql, args...).Scan(&primary, &secondary); err != nil { + return 0, 0, fmt.Errorf("sendingpolicy: ledger retention %s: %w", sweep.table, err) + } + if err := tx.Commit(ctx); err != nil { + return 0, 0, fmt.Errorf("sendingpolicy: ledger retention %s: commit: %w", sweep.table, err) + } + return primary, secondary, nil +} + +// countRetainedOperations counts, up to a cap, the expired operations that a +// guard kept. A number that only grows is a stuck attempt or job, which an +// operator should look at; nothing here deletes it. +func (m *Module) countRetainedOperations(ctx context.Context, now time.Time) (int64, error) { + var n int64 + err := m.pool.QueryRow(ctx, ` + SELECT count(*) FROM ( + SELECT 1 FROM sending_provider_operations WHERE expires_at <= $1 LIMIT $2 + ) AS expired`, now, retainedCountCap).Scan(&n) + return n, err +} + +// LedgerRetentionArgs is the periodic ledger retention job. +type LedgerRetentionArgs struct{} + +func (LedgerRetentionArgs) Kind() string { return "sending_ledger_retention" } + +// LedgerRetentionWorker runs one GCLedger pass. +type LedgerRetentionWorker struct { + river.WorkerDefaults[LedgerRetentionArgs] + module *Module +} + +// NewLedgerRetentionWorker builds the ledger janitor over a module. Exported +// so a test can drive the component that deletes ledger rows directly. +func NewLedgerRetentionWorker(module *Module) *LedgerRetentionWorker { + return &LedgerRetentionWorker{module: module} +} + +// Work never returns a table failure to River: a failure is logged and +// counted, and the next hourly run retries. Returning it would only make +// River re-run the whole pass on its retry schedule against the same lock. +func (w *LedgerRetentionWorker) Work(ctx context.Context, _ *river.Job[LedgerRetentionArgs]) error { + st, err := w.module.GCLedger(ctx, w.module.now()) + outcome := st.Outcome() + if err != nil { + outcome = LedgerRunFailed + } + if o := currentLedgerObserver(); o != nil { + for _, table := range ledgerMetricTables { + if n := st.Deleted[table]; n > 0 { + o.JanitorRowsDeleted(table, int(n)) + } + } + o.SendingLedgerRetentionRun(outcome) + } + log.Printf("[sendingpolicy:ledger-gc] outcome=%s %s retained_expired_operations=%d", + outcome, formatLedgerCounts(st.Deleted), st.RetainedOperations) + if joined := st.Err(); joined != nil { + log.Printf("[sendingpolicy:ledger-gc] table failures (retried next run): %v", joined) + } + if err != nil { + log.Printf("[sendingpolicy:ledger-gc] run interrupted: %v", err) + } + return nil +} + +// Timeout matches the cleanup janitor's budget: every batch autocommits, so a +// cut is safe and the next run resumes. +func (w *LedgerRetentionWorker) Timeout(*river.Job[LedgerRetentionArgs]) time.Duration { + return 5 * time.Minute +} + +// ledgerMetricTables is every label the janitor can emit, in a stable order. +var ledgerMetricTables = []string{ + LedgerTableOperations, LedgerTableReservations, LedgerTableCounters, + LedgerTableNoticeEvents, LedgerTableNoticeDeliveries, + LedgerTableControlEvents, LedgerTableAccessEvents, +} + +func formatLedgerCounts(deleted map[string]int64) string { + parts := make([]string, 0, len(ledgerMetricTables)) + for _, table := range ledgerMetricTables { + parts = append(parts, fmt.Sprintf("%s=%d", table, deleted[table])) + } + return strings.Join(parts, " ") +} diff --git a/internal/sendingpolicy/ledger_retention_integration_test.go b/internal/sendingpolicy/ledger_retention_integration_test.go new file mode 100644 index 000000000..fed23ffcd --- /dev/null +++ b/internal/sendingpolicy/ledger_retention_integration_test.go @@ -0,0 +1,522 @@ +package sendingpolicy_test + +import ( + "context" + "errors" + "fmt" + "strings" + "sync" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgconn" + "github.com/riverqueue/river" + + "github.com/tokencanopy/e2a/internal/delivery" + "github.com/tokencanopy/e2a/internal/jobs" + "github.com/tokencanopy/e2a/internal/sendingpolicy" +) + +// The sending-ledger retention janitor is the one component in the ledger +// that DELETES rows, so every guard it has is pinned here against real +// Postgres. All identifiers are synthetic. + +// ledgerFixture is a fixture whose database also has River's schema: the +// operation guard reads river_job. +func ledgerFixture(t *testing.T) (*fixture, *sendingpolicy.Module) { + t.Helper() + f := newFixture(t) + if err := jobs.Migrate(f.ctx, f.pool); err != nil { + t.Fatalf("river migrate: %v", err) + } + // River rows are not part of the harness truncation; clear ours. + f.exec(`DELETE FROM river_job WHERE kind LIKE 'ledger_test_%'`) + return f, sendingpolicy.NewModule(f.pool, f.secrets()) +} + +// op inserts one provider operation whose expiry is `expires` relative to now +// (a Postgres interval literal such as '-1 hour'). +func (f *fixture) op(id, expires string) { + f.t.Helper() + f.exec(`INSERT INTO sending_provider_operations + (operation_id, source_account_ref, policy_subject_ref, purpose, expires_at) + VALUES ($1, 'usr_ledger', 'usr_ledger', 'customer_message', now() + $2::interval)`, id, expires) +} + +// attempt inserts one reservation row in a consistent (state, call_state) +// shape with the given expiry. +func (f *fixture) attempt(opID string, n int, state, callState, expires string) { + f.t.Helper() + var nonce, started any + if callState != "none" { + nonce = fmt.Sprintf("nonce_%s_%d", opID, n) + } + if callState == "started" { + started = time.Now().UTC() + } + f.exec(`INSERT INTO sending_budget_reservations + (operation_id, submission_attempt, source_account_ref, policy_subject_ref, purpose, + day, units, probation, state, call_state, authorization_nonce, + provider_call_started_at, expires_at) + VALUES ($1, $2, 'usr_ledger', 'usr_ledger', 'customer_message', + current_date, 1, false, $3, $4, $5, $6, now() + $7::interval)`, + opID, n, state, callState, nonce, started, expires) +} + +func (f *fixture) count(sql string, args ...any) int { + f.t.Helper() + var n int + if err := f.pool.QueryRow(f.ctx, sql, args...).Scan(&n); err != nil { + f.t.Fatalf("count %q: %v", sql, err) + } + return n +} + +func (f *fixture) opExists(id string) bool { + return f.count(`SELECT count(*) FROM sending_provider_operations WHERE operation_id = $1`, id) == 1 +} + +func (f *fixture) attempts(opID string) int { + return f.count(`SELECT count(*) FROM sending_budget_reservations WHERE operation_id = $1`, opID) +} + +func gcLedger(t *testing.T, m *sendingpolicy.Module) sendingpolicy.LedgerRetentionStats { + t.Helper() + st, err := m.GCLedger(context.Background(), time.Now()) + if err != nil { + t.Fatalf("GCLedger: %v", err) + } + if st.Outcome() != sendingpolicy.LedgerRunComplete { + t.Fatalf("outcome = %s (failed %v, capped %v)", st.Outcome(), st.Failed, st.Capped) + } + return st +} + +// TestLedgerRetentionOperationsAndAttempts pins the 30-day operation/attempt +// horizon and every reason an expired operation must survive. +func TestLedgerRetentionOperationsAndAttempts(t *testing.T) { + f, m := ledgerFixture(t) + + // Inside the window. + f.op("op_fresh", "1 day") + f.attempt("op_fresh", 1, "confirmed", "started", "-1 hour") + + // Past the window, every attempt terminal and expired: deleted, with + // both of its attempts. + f.op("op_done", "-1 hour") + f.attempt("op_done", 1, "confirmed", "started", "-2 hours") + f.attempt("op_done", 2, "released", "none", "-1 hour") + + // Past the window, no attempt at all (never reserved): deleted. + f.op("op_bare", "-1 hour") + + // Non-terminal attempts hold their operation however old it is. + f.op("op_reserved", "-10 days") + f.attempt("op_reserved", 1, "reserved", "none", "-10 days") + f.op("op_authorized", "-10 days") + f.attempt("op_authorized", 1, "confirmed", "authorized", "-10 days") + + // A later attempt still inside its own window holds the operation. + f.op("op_retry", "-1 hour") + f.attempt("op_retry", 1, "confirmed", "started", "-1 hour") + f.attempt("op_retry", 2, "confirmed", "started", "5 days") + + st := gcLedger(t, m) + + for _, id := range []string{"op_done", "op_bare"} { + if f.opExists(id) || f.attempts(id) != 0 { + t.Errorf("%s: expired terminal operation or its attempts survived", id) + } + } + for id, attempts := range map[string]int{"op_fresh": 1, "op_reserved": 1, "op_authorized": 1, "op_retry": 2} { + if !f.opExists(id) || f.attempts(id) != attempts { + t.Errorf("%s: retained operation lost (exists=%v attempts=%d, want %d)", id, f.opExists(id), f.attempts(id), attempts) + } + } + if st.Deleted[sendingpolicy.LedgerTableOperations] != 2 || st.Deleted[sendingpolicy.LedgerTableReservations] != 2 { + t.Errorf("deleted = %v, want 2 operations and 2 attempts", st.Deleted) + } + if st.RetainedOperations != 3 { + t.Errorf("retained expired operations = %d, want 3 (reserved, authorized, retry)", st.RetainedOperations) + } +} + +// TestLedgerRetentionKeepsOperationsInFlight: an operation a send still +// needs is never deleted by age — a live River job (budget or pause hold, +// retry), a pre-terminal customer message, or a pending protection notice. +// Deleting any of these would make the worker fail closed and drop the send. +func TestLedgerRetentionKeepsOperationsInFlight(t *testing.T) { + f, m := ledgerFixture(t) + + // River jobs: snoozed (a budget or pause hold is a snooze), retryable, + // and running all hold; a completed job does not. + for _, state := range []string{"scheduled", "retryable", "running", "available"} { + id := "op_job_" + state + f.op(id, "-40 days") + f.exec(`INSERT INTO river_job (kind, state, args, attempted_at) + VALUES ('ledger_test_send', $1::river_job_state, jsonb_build_object('operation_ref', jsonb_build_object('v', 1, 'id', $2::text)), + CASE WHEN $1 = 'running' THEN now() END)`, state, id) + } + f.op("op_job_done", "-40 days") + f.exec(`INSERT INTO river_job (kind, state, args, finalized_at) + VALUES ('ledger_test_send', 'completed', jsonb_build_object('operation_ref', jsonb_build_object('v', 1, 'id', 'op_job_done')), now())`) + + // Customer messages: the operation id IS the message id. + user := f.user("standard") + agent := f.agent(user) + inFlight := map[string]bool{} + for _, status := range []string{"accepted", "sending", "delivered", "failed"} { + msg := f.message(agent, "relay", 1) + f.exec(`UPDATE messages SET delivery_status = $2 WHERE id = $1`, msg, status) + f.op(msg, "-40 days") + inFlight[msg] = status == "accepted" || status == "sending" + } + held := f.pendingMessage(agent, "relay") + f.op(held, "-40 days") + inFlight[held] = true + + // Protection notices: a pending delivery holds its operation; a sent one + // does not. + f.exec(`INSERT INTO sending_protection_notice_events (id, account_ref, kind, reason_code, source_event_id, expires_at) + VALUES ('spn_ledger_1', 'usr_ledger', 'pause', 'manual', 'sce_ledger_1', now() + interval '90 days'), + ('spn_ledger_2', 'usr_ledger', 'pause', 'manual', 'sce_ledger_2', now() + interval '90 days')`) + f.op("opn_pending", "-40 days") + f.op("opn_sent", "-40 days") + f.exec(`INSERT INTO sending_protection_notice_deliveries (event_id, audience, current_operation_id, state) + VALUES ('spn_ledger_1', 'owner', 'opn_pending', 'pending'), + ('spn_ledger_2', 'owner', 'opn_sent', 'sent')`) + + gcLedger(t, m) + + for _, state := range []string{"scheduled", "retryable", "running", "available"} { + if !f.opExists("op_job_" + state) { + t.Errorf("operation named by a %s River job was deleted", state) + } + } + if f.opExists("op_job_done") { + t.Error("operation named only by a completed job survived") + } + for msg, keep := range inFlight { + if f.opExists(msg) != keep { + t.Errorf("message operation %s: exists=%v, want %v", msg, f.opExists(msg), keep) + } + } + if !f.opExists("opn_pending") { + t.Error("operation bound to a pending notice delivery was deleted") + } + if f.opExists("opn_sent") { + t.Error("operation bound only to a sent notice delivery survived") + } +} + +// TestLedgerRetentionLeavesFeedbackProvenanceIntact: B8 correlations +// reference (operation_id, submission_attempt) by value with no foreign key, +// and they are kept for the account's lifetime — far past the operation's 30 +// days. Deleting the operation must not touch them, and feedback must still +// correlate afterwards: that is what makes it safe for operations NOT to +// outlive their correlations. +func TestLedgerRetentionLeavesFeedbackProvenanceIntact(t *testing.T) { + f, m := ledgerFixture(t) + g := f.gate(sendingpolicy.DisabledPolicy()) + user := f.user("standard") + agent := f.agent(user) + rcpts := []string{"carol@example.test"} + msg := f.messageTo(agent, "relay", rcpts) + _, corrID, sesID := f.authorizedSend(g, msg, rcpts) + + // Age the operation and its attempt past the 30-day horizon. + f.exec(`UPDATE sending_provider_operations SET expires_at = now() - interval '1 hour' WHERE operation_id = $1`, msg) + f.exec(`UPDATE sending_budget_reservations SET expires_at = now() - interval '1 hour' WHERE operation_id = $1`, msg) + + gcLedger(t, m) + + if f.opExists(msg) || f.attempts(msg) != 0 { + t.Fatal("expired, settled operation was not deleted") + } + if n := f.count(`SELECT count(*) FROM sending_feedback_correlations WHERE correlation_id = $1 AND expires_at IS NULL`, corrID); n != 1 { + t.Fatal("the account-lifetime correlation was touched by the ledger janitor") + } + if n := f.count(`SELECT count(*) FROM sending_feedback_recipients WHERE correlation_id = $1`, corrID); n != 1 { + t.Fatal("recipient provenance was touched by the ledger janitor") + } + res, err := m.ProcessProviderFeedback(f.ctx, delivery.ProviderFeedback{ + ProviderEventID: "evt-ledger-complaint", OccurredAt: time.Now().UTC(), Kind: delivery.KindComplaint, + ProviderMessageID: sesID, Recipients: rcpts, + }) + if err != nil { + t.Fatal(err) + } + if !res.Correlated || res.AccountRef != user { + t.Fatalf("feedback after operation deletion = %+v, want correlated to the sending account", res) + } +} + +// TestLedgerRetentionCounters: rows for today and yesterday survive; a day is +// deleted once it and the day after it have both closed. Deletion can never +// refund capacity because the open day is never a victim. +func TestLedgerRetentionCounters(t *testing.T) { + f, m := ledgerFixture(t) + for _, age := range []int{0, 1, 2, 3, 30} { + f.exec(`INSERT INTO sending_budget_counters (scope, scope_id, day, reserved_count, confirmed_count, daily_limit) + VALUES ('account_daily', 'usr_ledger', (now() AT TIME ZONE 'UTC')::date - $1::int, 5, 5, 20)`, age) + } + st := gcLedger(t, m) + + var ages []int + rows, err := f.pool.Query(f.ctx, `SELECT (now() AT TIME ZONE 'UTC')::date - day FROM sending_budget_counters ORDER BY 1`) + if err != nil { + t.Fatal(err) + } + defer rows.Close() + for rows.Next() { + var a int + if err := rows.Scan(&a); err != nil { + t.Fatal(err) + } + ages = append(ages, a) + } + if fmt.Sprint(ages) != "[0 1]" { + t.Fatalf("remaining counter ages = %v, want [0 1] (today and yesterday)", ages) + } + if st.Deleted[sendingpolicy.LedgerTableCounters] != 3 { + t.Fatalf("deleted counters = %d, want 3", st.Deleted[sendingpolicy.LedgerTableCounters]) + } +} + +// TestLedgerRetentionNoticeOutbox: terminal pairs go 30 days after event +// creation; a pending delivery keeps its event however old. +func TestLedgerRetentionNoticeOutbox(t *testing.T) { + f, m := ledgerFixture(t) + insert := func(id, age string, states ...string) { + f.exec(`INSERT INTO sending_protection_notice_events (id, account_ref, kind, reason_code, source_event_id, created_at, expires_at) + VALUES ($1, 'usr_ledger', 'pause', 'manual', $1 || '_src', now() - $2::interval, now() + interval '60 days')`, id, age) + for i, state := range states { + audience := []string{"owner", "operator"}[i] + f.exec(`INSERT INTO sending_protection_notice_deliveries (event_id, audience, state) VALUES ($1, $2, $3)`, id, audience, state) + } + } + insert("spn_old_terminal", "31 days", "sent", "failed") + insert("spn_old_skipped", "31 days", "skipped_account_deleted") + insert("spn_old_pending", "45 days", "sent", "pending") + insert("spn_young_terminal", "29 days", "sent", "sent") + + st := gcLedger(t, m) + + gone := func(id string) bool { + return f.count(`SELECT count(*) FROM sending_protection_notice_events WHERE id = $1`, id) == 0 && + f.count(`SELECT count(*) FROM sending_protection_notice_deliveries WHERE event_id = $1`, id) == 0 + } + for id, wantGone := range map[string]bool{ + "spn_old_terminal": true, "spn_old_skipped": true, + "spn_old_pending": false, "spn_young_terminal": false, + } { + if gone(id) != wantGone { + t.Errorf("%s: gone=%v, want %v", id, gone(id), wantGone) + } + } + if f.count(`SELECT count(*) FROM sending_protection_notice_deliveries WHERE event_id = 'spn_old_pending'`) != 2 { + t.Error("the pending event lost a delivery") + } + if st.Deleted[sendingpolicy.LedgerTableNoticeEvents] != 2 || st.Deleted[sendingpolicy.LedgerTableNoticeDeliveries] != 3 { + t.Errorf("deleted = %v, want 2 events and 3 deliveries", st.Deleted) + } +} + +// TestLedgerRetentionControlAndAccessAudit: both audits go at their stamped +// expiry, except an abuse-class pause event whose account still exists — +// the account purge derives the abuse hold from that history. +func TestLedgerRetentionControlAndAccessAudit(t *testing.T) { + f, m := ledgerFixture(t) + live := f.user("standard") + control := func(id, account string, class any, expires string) { + f.exec(`INSERT INTO account_sending_control_events + (id, account_ref, old_state, new_state, reason, actor, pause_class, created_at, expires_at) + VALUES ($1, $2, 'active', 'paused', 'synthetic', 'operator', $3, + now() + $4::interval - interval '90 days', now() + $4::interval)`, id, account, class, expires) + } + control("sce_expired", live, "operator", "-1 hour") + control("sce_expired_null_class", live, nil, "-1 hour") + control("sce_unexpired", live, "operator", "1 day") + control("sce_abuse_live", live, "abuse", "-1 hour") + control("sce_abuse_purged", "usr_ledger_purged", "abuse", "-1 hour") + + access := func(id, expires string) { + f.exec(`INSERT INTO external_sending_access_events + (id, account_ref, old_approved, new_approved, old_revision, new_revision, actor, reason, created_at, expires_at) + VALUES ($1, 'usr_ledger', false, true, 0, 1, 'operator', 'synthetic', + now() + $2::interval - interval '90 days', now() + $2::interval)`, id, expires) + } + access("esa_expired", "-1 hour") + access("esa_unexpired", "1 day") + + st := gcLedger(t, m) + + for id, keep := range map[string]bool{ + "sce_expired": false, "sce_expired_null_class": false, "sce_unexpired": true, + "sce_abuse_live": true, "sce_abuse_purged": false, + } { + if got := f.count(`SELECT count(*) FROM account_sending_control_events WHERE id = $1`, id) == 1; got != keep { + t.Errorf("control event %s: kept=%v, want %v", id, got, keep) + } + } + for id, keep := range map[string]bool{"esa_expired": false, "esa_unexpired": true} { + if got := f.count(`SELECT count(*) FROM external_sending_access_events WHERE id = $1`, id) == 1; got != keep { + t.Errorf("access event %s: kept=%v, want %v", id, got, keep) + } + } + if st.Deleted[sendingpolicy.LedgerTableControlEvents] != 3 || st.Deleted[sendingpolicy.LedgerTableAccessEvents] != 1 { + t.Errorf("deleted = %v", st.Deleted) + } +} + +// TestLedgerRetentionBatches: the loop drains across multiple batches, and a +// per-run cap stops a table (outcome partial) for the next run to finish. +func TestLedgerRetentionBatches(t *testing.T) { + f, m := ledgerFixture(t) + for i := 0; i < 25; i++ { + f.op(fmt.Sprintf("op_batch_%02d", i), "-1 hour") + } + + st, err := m.GCLedgerTuned(f.ctx, time.Now(), sendingpolicy.LedgerTuning{BatchSize: 10, MaxBatches: 2}) + if err != nil { + t.Fatal(err) + } + if st.Deleted[sendingpolicy.LedgerTableOperations] != 20 || st.Outcome() != sendingpolicy.LedgerRunPartial { + t.Fatalf("capped run: deleted=%d outcome=%s, want 20 and partial", st.Deleted[sendingpolicy.LedgerTableOperations], st.Outcome()) + } + if len(st.Capped) != 1 || st.Capped[0] != sendingpolicy.LedgerTableOperations { + t.Fatalf("capped = %v", st.Capped) + } + + st, err = m.GCLedgerTuned(f.ctx, time.Now(), sendingpolicy.LedgerTuning{BatchSize: 10, MaxBatches: 10}) + if err != nil { + t.Fatal(err) + } + if st.Deleted[sendingpolicy.LedgerTableOperations] != 5 || st.Outcome() != sendingpolicy.LedgerRunComplete { + t.Fatalf("resumed run: deleted=%d outcome=%s, want 5 and complete", st.Deleted[sendingpolicy.LedgerTableOperations], st.Outcome()) + } + if n := f.count(`SELECT count(*) FROM sending_provider_operations WHERE operation_id LIKE 'op_batch_%'`); n != 0 { + t.Fatalf("%d operations left after draining", n) + } +} + +// TestLedgerRetentionLockTimeoutIsNonFatal: a table whose lock cannot be +// acquired fails only itself for this run; the other tables are still swept, +// the worker returns nil, and the next run finishes the job. +func TestLedgerRetentionLockTimeoutIsNonFatal(t *testing.T) { + f, m := ledgerFixture(t) + f.exec(`INSERT INTO sending_budget_counters (scope, scope_id, day, daily_limit) + VALUES ('global_all', 'global', (now() AT TIME ZONE 'UTC')::date - 5, 5000)`) + f.exec(`INSERT INTO external_sending_access_events + (id, account_ref, old_approved, new_approved, old_revision, new_revision, actor, reason, created_at, expires_at) + VALUES ('esa_lock', 'usr_ledger', false, true, 0, 1, 'operator', 'synthetic', now() - interval '91 days', now() - interval '1 hour')`) + + // Another session holds a lock that conflicts with the janitor's + // FOR UPDATE on the counters table. + holder, err := f.pool.Begin(f.ctx) + if err != nil { + t.Fatal(err) + } + var once sync.Once + release := func() { once.Do(func() { _ = holder.Rollback(f.ctx) }) } + defer release() + if _, err := holder.Exec(f.ctx, `LOCK TABLE sending_budget_counters IN EXCLUSIVE MODE`); err != nil { + t.Fatal(err) + } + + fast := sendingpolicy.LedgerTuning{BatchSize: sendingpolicy.LedgerBatchSize, MaxBatches: 5, + BeforeBatch: func(ctx context.Context, _ string, exec func(context.Context, string) error) error { + return exec(ctx, `SET LOCAL lock_timeout = '100ms'`) + }} + st, err := m.GCLedgerTuned(f.ctx, time.Now(), fast) + if err != nil { + t.Fatalf("a table lock timeout must not fail the run: %v", err) + } + failure, failed := st.Failed[sendingpolicy.LedgerTableCounters] + var pgErr *pgconn.PgError + if !failed || !errors.As(failure, &pgErr) || pgErr.Code != "55P03" { + t.Fatalf("counters failure = %v, want a lock_not_available (55P03) error", failure) + } + if st.Outcome() != sendingpolicy.LedgerRunFailed { + t.Fatalf("outcome = %s, want failed", st.Outcome()) + } + if st.Deleted[sendingpolicy.LedgerTableAccessEvents] != 1 { + t.Fatal("a failing table stopped the tables after it") + } + + release() + st, err = m.GCLedgerTuned(f.ctx, time.Now(), fast) + if err != nil || st.Outcome() != sendingpolicy.LedgerRunComplete || st.Deleted[sendingpolicy.LedgerTableCounters] != 1 { + t.Fatalf("retry run: err=%v outcome=%s deleted=%v", err, st.Outcome(), st.Deleted) + } +} + +// recordingLedgerObserver captures the janitor's metric samples. +type recordingLedgerObserver struct { + mu sync.Mutex + deleted map[string]int + runs []string +} + +func (r *recordingLedgerObserver) JanitorRowsDeleted(table string, count int) { + r.mu.Lock() + defer r.mu.Unlock() + if r.deleted == nil { + r.deleted = map[string]int{} + } + r.deleted[table] += count +} + +func (r *recordingLedgerObserver) SendingLedgerRetentionRun(outcome string) { + r.mu.Lock() + defer r.mu.Unlock() + r.runs = append(r.runs, outcome) +} + +// TestLedgerRetentionWorkerEmitsBoundedMetrics: the worker reports rows +// deleted per ledger table and one run outcome, with table names as the only +// labels — never an operation, account, or event id. +func TestLedgerRetentionWorkerEmitsBoundedMetrics(t *testing.T) { + f, m := ledgerFixture(t) + rec := &recordingLedgerObserver{} + sendingpolicy.SetLedgerRetentionObserver(rec) + t.Cleanup(func() { sendingpolicy.SetLedgerRetentionObserver(nil) }) + + f.op("op_metric", "-1 hour") + f.attempt("op_metric", 1, "confirmed", "started", "-1 hour") + f.exec(`INSERT INTO sending_budget_counters (scope, scope_id, day, daily_limit) + VALUES ('global_all', 'global', (now() AT TIME ZONE 'UTC')::date - 3, 5000)`) + + if err := sendingpolicy.NewLedgerRetentionWorker(m).Work(f.ctx, &river.Job[sendingpolicy.LedgerRetentionArgs]{}); err != nil { + t.Fatalf("Work: %v", err) + } + want := map[string]int{ + sendingpolicy.LedgerTableOperations: 1, + sendingpolicy.LedgerTableReservations: 1, + sendingpolicy.LedgerTableCounters: 1, + } + if fmt.Sprint(rec.deleted) != fmt.Sprint(want) { + t.Fatalf("deleted samples = %v, want %v", rec.deleted, want) + } + if len(rec.runs) != 1 || rec.runs[0] != sendingpolicy.LedgerRunComplete { + t.Fatalf("run samples = %v, want one complete", rec.runs) + } + for table := range rec.deleted { + if strings.Contains(table, "op_") || strings.Contains(table, "usr_") { + t.Fatalf("metric label %q carries an identifier", table) + } + } +} + +// TestLedgerRetentionRegistersOnTheMaintenanceQueue: the janitor is one of +// the sending-protection maintenance periodics, so a registrar that dropped +// it would leave the ledger growing forever with nothing failing. +func TestLedgerRetentionRegistersOnTheMaintenanceQueue(t *testing.T) { + if kind := (sendingpolicy.LedgerRetentionArgs{}).Kind(); kind != "sending_ledger_retention" { + t.Fatalf("kind = %q", kind) + } + periodics := sendingpolicy.NewMaintenanceJobs(nil).RegisterJobs(river.NewWorkers()) + if len(periodics) != 3 { + t.Fatalf("periodic jobs = %d, want 3 (feedback retention, reconcile, ledger retention)", len(periodics)) + } +} diff --git a/internal/sendingpolicy/maintenance.go b/internal/sendingpolicy/maintenance.go index 8f1880469..aac304046 100644 --- a/internal/sendingpolicy/maintenance.go +++ b/internal/sendingpolicy/maintenance.go @@ -95,8 +95,9 @@ func (w *FeedbackReconcileWorker) Work(ctx context.Context, _ *river.Job[Feedbac return nil } -// MaintenanceJobs registers the feedback retention periodics. Implements -// jobs.Registrar. +// MaintenanceJobs registers the sending-protection retention periodics: the +// feedback provenance pass, its reconcile backstop, and the ledger janitor. +// Implements jobs.Registrar. type MaintenanceJobs struct{ module *Module } // NewMaintenanceJobs builds the registrar over the gate's module. @@ -105,6 +106,7 @@ func NewMaintenanceJobs(module *Module) *MaintenanceJobs { return &MaintenanceJo func (m *MaintenanceJobs) RegisterJobs(w *river.Workers) []*river.PeriodicJob { river.AddWorker(w, &FeedbackMaintenanceWorker{module: m.module}) river.AddWorker(w, &FeedbackReconcileWorker{module: m.module}) + river.AddWorker(w, &LedgerRetentionWorker{module: m.module}) return []*river.PeriodicJob{ river.NewPeriodicJob( river.PeriodicInterval(feedbackMaintenanceInterval), @@ -120,5 +122,15 @@ func (m *MaintenanceJobs) RegisterJobs(w *river.Workers) []*river.PeriodicJob { }, &river.PeriodicJobOpts{RunOnStart: true}, ), + // The ledger janitor also runs on start: blue/green deploys restart + // the process more often than a slow cadence would tick, and the + // first run after rollout has weeks of backlog to start on. + river.NewPeriodicJob( + river.PeriodicInterval(ledgerRetentionInterval), + func() (river.JobArgs, *river.InsertOpts) { + return LedgerRetentionArgs{}, &river.InsertOpts{Queue: jobs.QueueMaintenance} + }, + &river.PeriodicJobOpts{RunOnStart: true}, + ), } } diff --git a/internal/sendingpolicy/maintenance_test.go b/internal/sendingpolicy/maintenance_test.go index 4754d8697..fba3109b0 100644 --- a/internal/sendingpolicy/maintenance_test.go +++ b/internal/sendingpolicy/maintenance_test.go @@ -127,8 +127,8 @@ func TestFeedbackMaintenanceRegistersOnTheMaintenanceQueue(t *testing.T) { f := newFixture(t) module := sendingpolicy.NewModule(f.pool, f.secrets()) periodics := sendingpolicy.NewMaintenanceJobs(module).RegisterJobs(river.NewWorkers()) - if len(periodics) != 2 { - t.Fatalf("periodic jobs = %d, want 2 (retention pass + reconcile)", len(periodics)) + if len(periodics) != 3 { + t.Fatalf("periodic jobs = %d, want 3 (retention pass + reconcile + ledger retention)", len(periodics)) } // River keeps the periodic's constructor unexported, so assert the two // facts that are observable and load-bearing: the job kind the worker diff --git a/internal/telemetry/metrics.go b/internal/telemetry/metrics.go index 6f2df0e60..661cc7b99 100644 --- a/internal/telemetry/metrics.go +++ b/internal/telemetry/metrics.go @@ -159,6 +159,14 @@ type Metrics interface { // only: never an address, account, correlation or provider id. SendingFeedbackIngested(outcome, bucket string) + // SendingLedgerRetentionRun records one sending-ledger retention pass + // (internal/sendingpolicy). outcome ∈ {complete, partial, failed}: + // partial = a table hit its per-run batch cap and continues next run; + // failed = a table's batch failed (lock/statement timeout) and is + // retried next run. Rows it deletes are counted on JanitorRowsDeleted + // with the ledger table as the label. + SendingLedgerRetentionRun(outcome string) + // WebhookAttempt records one webhook delivery attempt. outcome ∈ // {delivered, retryable_failure, exhausted, webhook_deleted, // skipped_disabled}. statusClass is the HTTP status class of the @@ -330,6 +338,7 @@ func (NoOp) OutboundAttempt(string, float64) {} func (NoOp) OutboundRateDeferred() {} func (NoOp) ExternalAccessDecision(string, string, string) {} func (NoOp) SendingFeedbackIngested(string, string) {} +func (NoOp) SendingLedgerRetentionRun(string) {} func (NoOp) WebhookAttempt(string, string, float64) {} func (NoOp) WebhookTerminal(string, string, int) {} func (NoOp) WebhookNotify(string, string) {} @@ -468,6 +477,10 @@ func (l *Log) SendingFeedbackIngested(outcome, bucket string) { log.Printf("[metrics] event=sending_feedback.ingested outcome=%s bucket=%s", outcome, bucket) } +func (l *Log) SendingLedgerRetentionRun(outcome string) { + log.Printf("[metrics] event=sending_ledger.retention_run outcome=%s", outcome) +} + func (l *Log) WebhookAttempt(outcome, statusClass string, seconds float64) { log.Printf("[metrics] event=webhook.attempt outcome=%s status_class=%s duration=%.3f", outcome, statusClass, seconds) } diff --git a/internal/telemetry/prom.go b/internal/telemetry/prom.go index f074cc05c..33d7f9d47 100644 --- a/internal/telemetry/prom.go +++ b/internal/telemetry/prom.go @@ -32,6 +32,7 @@ type Prom struct { outRateDeferred prometheus.Counter externalAccess *prometheus.CounterVec sendingFeedback *prometheus.CounterVec + sendingLedgerRuns *prometheus.CounterVec whAttempts *prometheus.CounterVec whAttemptDur prometheus.Histogram whTerminal *prometheus.CounterVec @@ -107,6 +108,7 @@ var ( sendingFeedbackOutcomeSet = set("correlated", "dead_account", "unmatched_recipient", "uncorrelated_with_marker", "uncorrelated", "duplicate") sendingFeedbackBucketSet = set("delivered", "hard_bounce", "complaint", "terminal_other", "none") + sendingLedgerRunSet = set("complete", "partial", "failed") whSet = set("delivered", "retryable_failure", "exhausted", "webhook_deleted", "skipped_disabled") whTerminalSet = set("delivered", "e2a_failure", "endpoint_failure", "excluded") @@ -142,7 +144,12 @@ var ( scopeSet = set("single", "since") tableSet = set("webhook_events", "webhook_subscriber_deliveries", "webhook_deliveries", "messages", "agent_identities", - "user_sessions", "oauth") + "user_sessions", "oauth", + // Sending-ledger retention (internal/sendingpolicy ledger janitor). + "sending_provider_operations", "sending_budget_reservations", + "sending_budget_counters", "sending_protection_notice_events", + "sending_protection_notice_deliveries", "account_sending_control_events", + "external_sending_access_events") threadResolutionSet = set( "api_reply", "fresh_send", "forward", "rfc_in_reply_to", "rfc_references", "self_twin", "authenticated_delivery_twin", @@ -279,6 +286,10 @@ func NewProm(build string) *Prom { Name: "e2a_sending_feedback_ingested_total", Help: "Deletion-resistant SES feedback ingestion results by outcome and detector bucket (uncorrelated_with_marker = e2a-stamped mail with no retained correlation).", }, []string{"outcome", "bucket"}), + sendingLedgerRuns: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "e2a_sending_ledger_retention_runs_total", + Help: "Sending-ledger retention passes by outcome (partial = a table hit its per-run batch cap; failed = a table's batch failed and is retried next run). Rows deleted are on e2a_janitor_rows_deleted_total.", + }, []string{"outcome"}), outRateDeferred: prometheus.NewCounter(prometheus.CounterOpts{ Name: "e2a_outbound_rate_deferred_total", Help: "Outbound submissions deferred by the per-agent fire-time rate limiter (snoozed, re-fired when the window frees capacity).", @@ -444,7 +455,7 @@ func NewProm(build string) *Prom { registerer.MustRegister( p.httpRequests, p.httpDuration, p.smtpInbound, p.smtpDuration, - p.outQueueWait, p.outTerminal, p.outTerminalLat, p.outAttempts, p.outAttemptDur, p.outRateDeferred, p.externalAccess, p.sendingFeedback, + p.outQueueWait, p.outTerminal, p.outTerminalLat, p.outAttempts, p.outAttemptDur, p.outRateDeferred, p.externalAccess, p.sendingFeedback, p.sendingLedgerRuns, p.whAttempts, p.whAttemptDur, p.whTerminal, p.whNotify, p.whExpiredPending, p.whFanOutRescued, p.whDeliveryRescued, p.whFirstTryLat, p.wsConnects, p.wsDisconnects, p.wsRejected, p.wsDrained, p.wsSendFailures, p.wsActive, p.delegatedFailures, p.delegatedRefresh, p.oidcDiscovery, p.oidcCallback, p.provisioning, @@ -525,6 +536,10 @@ func (p *Prom) SendingFeedbackIngested(outcome, bucket string) { p.sendingFeedback.WithLabelValues(enum(sendingFeedbackOutcomeSet, outcome), enum(sendingFeedbackBucketSet, bucket)).Inc() } +func (p *Prom) SendingLedgerRetentionRun(outcome string) { + p.sendingLedgerRuns.WithLabelValues(enum(sendingLedgerRunSet, outcome)).Inc() +} + func (p *Prom) WebhookAttempt(outcome, statusClass string, seconds float64) { p.whAttempts.WithLabelValues(enum(whSet, outcome), enum(classSet, statusClass)).Inc() // seconds < 0 = "no duration sample" (outcomes with no HTTP POST — diff --git a/internal/telemetry/prom_test.go b/internal/telemetry/prom_test.go index 306f1bb14..b1dafa268 100644 --- a/internal/telemetry/prom_test.go +++ b/internal/telemetry/prom_test.go @@ -209,6 +209,9 @@ func TestPromEmitsSMTPOutboundWebhookWSSeries(t *testing.T) { p.SendingFeedbackIngested("uncorrelated_with_marker", "complaint") p.SendingFeedbackIngested("correlated", "delivered") p.SendingFeedbackIngested("bob@example.test", "cor_0123") + p.SendingLedgerRetentionRun("complete") + p.SendingLedgerRetentionRun("failed") + p.SendingLedgerRetentionRun("op_0123") out := scrape(t, p) for _, want := range []string{ @@ -229,6 +232,9 @@ func TestPromEmitsSMTPOutboundWebhookWSSeries(t *testing.T) { `e2a_sending_feedback_ingested_total{bucket="complaint",build="unknown",outcome="uncorrelated_with_marker"} 1`, `e2a_sending_feedback_ingested_total{bucket="delivered",build="unknown",outcome="correlated"} 1`, `e2a_sending_feedback_ingested_total{bucket="other",build="unknown",outcome="other"} 1`, + `e2a_sending_ledger_retention_runs_total{outcome="complete"} 1`, + `e2a_sending_ledger_retention_runs_total{outcome="failed"} 1`, + `e2a_sending_ledger_retention_runs_total{outcome="other"} 1`, `e2a_webhook_attempts_total{outcome="delivered",status_class="2xx"} 1`, `e2a_webhook_attempts_total{outcome="retryable_failure",status_class="5xx"} 1`, `e2a_webhook_delivery_terminal_total{outcome="delivered",scope="initial"} 1`, @@ -283,6 +289,8 @@ func TestPromEmitsLegacyOutboxSeries(t *testing.T) { p.OutboxFailures("lease") p.RedeliverRequests("single") p.JanitorRowsDeleted("webhook_events", 5) + p.JanitorRowsDeleted("sending_provider_operations", 3) + p.JanitorRowsDeleted("sending_budget_counters", 2) p.NotifyMissed() p.SetPublisherLag(2.5) @@ -295,6 +303,8 @@ func TestPromEmitsLegacyOutboxSeries(t *testing.T) { `e2a_outbox_failures_total{stage="lease"} 1`, `e2a_redeliver_requests_total{scope="single"} 1`, `e2a_janitor_rows_deleted_total{table="webhook_events"} 5`, + `e2a_janitor_rows_deleted_total{table="sending_provider_operations"} 3`, + `e2a_janitor_rows_deleted_total{table="sending_budget_counters"} 2`, `e2a_notify_missed_total 1`, `e2a_webhook_publisher_lag_seconds 2.5`, } { diff --git a/migrations/126_sending_provider_operations_expiry_idx.sql b/migrations/126_sending_provider_operations_expiry_idx.sql new file mode 100644 index 000000000..cd403c9b3 --- /dev/null +++ b/migrations/126_sending_provider_operations_expiry_idx.sql @@ -0,0 +1,21 @@ +-- 126_sending_provider_operations_expiry_idx.sql +-- e2a:no-transaction +-- +-- Access path for the sending-ledger retention janitor +-- (internal/sendingpolicy/ledger_retention.go), which selects expired +-- provider operations with `expires_at <= $1 ORDER BY expires_at LIMIT n`. +-- Migration 113 gave this table only its primary key, so without this index +-- every janitor batch would be a sequential scan of a table that gains one +-- row per provider-bound send and has never been pruned. +-- +-- CREATE INDEX CONCURRENTLY + e2a:no-transaction: the gate inserts and +-- updates these rows on the hot send path, and a production deployment +-- already holds every operation written since the gate shipped; a plain +-- CREATE INDEX would block those writes for the whole build. +-- +-- OPS NOTE — invalid-index recovery: an interrupted CONCURRENTLY build leaves +-- an INVALID index that IF NOT EXISTS then skips. To recover: +-- DROP INDEX CONCURRENTLY IF EXISTS sending_provider_operations_expiry_idx; +-- then re-run this statement. +CREATE INDEX CONCURRENTLY IF NOT EXISTS sending_provider_operations_expiry_idx + ON sending_provider_operations (expires_at, operation_id); diff --git a/migrations/127_sending_budget_counters_day_idx.sql b/migrations/127_sending_budget_counters_day_idx.sql new file mode 100644 index 000000000..c8bfde7a4 --- /dev/null +++ b/migrations/127_sending_budget_counters_day_idx.sql @@ -0,0 +1,17 @@ +-- 127_sending_budget_counters_day_idx.sql +-- e2a:no-transaction +-- +-- Access path for the sending-ledger retention janitor, which deletes day +-- counters for closed days (`day <= today - 2 ORDER BY day LIMIT n`). The +-- primary key (scope, scope_id, day) leads with scope, so it cannot serve a +-- predicate on day alone. +-- +-- CREATE INDEX CONCURRENTLY + e2a:no-transaction: counter rows are locked +-- and updated on every budgeted send; a plain CREATE INDEX would block them. +-- +-- OPS NOTE — invalid-index recovery: an interrupted CONCURRENTLY build leaves +-- an INVALID index that IF NOT EXISTS then skips. To recover: +-- DROP INDEX CONCURRENTLY IF EXISTS sending_budget_counters_day_idx; +-- then re-run this statement. +CREATE INDEX CONCURRENTLY IF NOT EXISTS sending_budget_counters_day_idx + ON sending_budget_counters (day); diff --git a/migrations/128_account_sending_control_events_expiry_idx.sql b/migrations/128_account_sending_control_events_expiry_idx.sql new file mode 100644 index 000000000..b535fbbe6 --- /dev/null +++ b/migrations/128_account_sending_control_events_expiry_idx.sql @@ -0,0 +1,18 @@ +-- 128_account_sending_control_events_expiry_idx.sql +-- e2a:no-transaction +-- +-- Access path for the sending-ledger retention janitor, which deletes pause +-- audit rows past their stamped expiry (`expires_at <= $1 ORDER BY +-- expires_at LIMIT n`). Migration 113 gave this table only its primary key. +-- The table is small today (operator and detector transitions only), but the +-- detector will write to it automatically once armed. +-- +-- CREATE INDEX CONCURRENTLY + e2a:no-transaction, matching 126/127, so the +-- build never blocks a pause or resume. +-- +-- OPS NOTE — invalid-index recovery: an interrupted CONCURRENTLY build leaves +-- an INVALID index that IF NOT EXISTS then skips. To recover: +-- DROP INDEX CONCURRENTLY IF EXISTS account_sending_control_events_expiry_idx; +-- then re-run this statement. +CREATE INDEX CONCURRENTLY IF NOT EXISTS account_sending_control_events_expiry_idx + ON account_sending_control_events (expires_at, id); From 4a2badede90414b6e9118e20d7f3f81353a49704 Mon Sep 17 00:00:00 2001 From: Josh Zhang <39790535+jiashuoz@users.noreply.github.com> Date: Wed, 30 Sep 2026 00:21:36 +0800 Subject: [PATCH 2/3] feat(sending): 30-day feedback provenance for system/internal accounts System- and internal-class accounts (prober, monitors, conformance) are never deleted, so the account-lifetime rule kept their feedback correlations forever; the standing prober makes them nearly the whole table. The gate now stamps their customer-purpose correlations with the post-account horizon (30 days) at creation, as it already does for non-customer purposes, and migration 129 backfills existing rows of those accounts with created_at + 30 days. Standard and demo accounts are unchanged. B8's existing janitor removes the expired rows. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_018tVLxUHk3fqQuq8C3wqyHW --- docs/design/async-message-pipeline.md | 12 ++ ...ck_synthetic_retention_integration_test.go | 132 ++++++++++++++++++ internal/sendingpolicy/gate.go | 10 +- ...9_sending_feedback_synthetic_retention.sql | 43 ++++++ 4 files changed, 196 insertions(+), 1 deletion(-) create mode 100644 internal/sendingpolicy/feedback_synthetic_retention_integration_test.go create mode 100644 migrations/129_sending_feedback_synthetic_retention.sql diff --git a/docs/design/async-message-pipeline.md b/docs/design/async-message-pipeline.md index 848dfd229..47d9250d9 100644 --- a/docs/design/async-message-pipeline.md +++ b/docs/design/async-message-pipeline.md @@ -436,6 +436,18 @@ provider retries. The seam never reads `messages`, `agent_identities`, or correlation authorized in a race with a purge) and sweeps events whose correlation is gone. Migration 124 ran that stamp once for the backlog with a **fixed 30 days** (the policy default), not the effective policy value. +- **Synthetic traffic keeps provenance for 30 days (2026-09-30).** System- + and internal-class accounts (the standing prober, monitors, conformance + runs) are never deleted, so the account-lifetime rule would keep their + correlations forever — and the prober sends every 30 seconds, so they are + nearly all of the table. The gate therefore stamps customer-purpose + correlations of a `system` or `internal` account (server-owned + `users.account_class`, read under the authorization lock) with + `expires_at = created_at + sending_feedback_post_account_retention_days` + (30 by default) at creation, exactly like non-customer purposes; events + inherit it and the hourly janitor removes them. Migration 129 stamped the + existing rows of those accounts once with a fixed `created_at + 30 days`. + Standard and demo accounts keep the account-lifetime rule. - **Keyring coverage is a startup gate**: a server that HAS a keyring refuses to start if any unexpired recipient row was signed under a version that keyring does not hold. Rotation is superset-first: add the new key diff --git a/internal/sendingpolicy/feedback_synthetic_retention_integration_test.go b/internal/sendingpolicy/feedback_synthetic_retention_integration_test.go new file mode 100644 index 000000000..5a50805a7 --- /dev/null +++ b/internal/sendingpolicy/feedback_synthetic_retention_integration_test.go @@ -0,0 +1,132 @@ +package sendingpolicy_test + +import ( + "testing" + "time" + + "github.com/tokencanopy/e2a/internal/sendingpolicy" + "github.com/tokencanopy/e2a/migrations" +) + +// Synthetic traffic from system and internal accounts (the standing prober, +// monitors, conformance) keeps its feedback provenance for 30 days, not for +// the account's lifetime: those accounts are never deleted, so the lifetime +// rule would retain every prober correlation forever. Standard and demo +// accounts keep the account-lifetime rule. + +// TestSyntheticAccountCorrelationsGetThirtyDayHorizon drives real authorized +// sends for each class and then the B8 janitor past the horizon. +func TestSyntheticAccountCorrelationsGetThirtyDayHorizon(t *testing.T) { + f := newFixture(t) + g := f.gate(sendingpolicy.DisabledPolicy()) + correlation := map[string]string{} + for _, class := range []string{"system", "internal", "standard", "demo"} { + agent := f.agent(f.user(class)) + rcpts := []string{class + "-rcpt@example.test"} + msg := f.messageTo(agent, "relay", rcpts) + _, corrID, _ := f.authorizedSend(g, msg, rcpts) + correlation[class] = corrID + } + + for class, corrID := range correlation { + var created time.Time + var expires *time.Time + if err := f.pool.QueryRow(f.ctx, + `SELECT created_at, expires_at FROM sending_feedback_correlations WHERE correlation_id = $1`, corrID, + ).Scan(&created, &expires); err != nil { + t.Fatal(err) + } + synthetic := class == "system" || class == "internal" + if !synthetic { + if expires != nil { + t.Errorf("%s-class correlation got expiry %v, want NULL (account lifetime)", class, expires) + } + continue + } + if expires == nil { + t.Fatalf("%s-class correlation has no expiry; synthetic traffic must expire after 30 days", class) + } + if d := expires.Sub(created); d < 30*24*time.Hour-time.Minute || d > 30*24*time.Hour+time.Minute { + t.Errorf("%s-class correlation horizon = %v, want 30 days after creation", class, d) + } + } + + // Past the horizon, B8's existing janitor removes the synthetic rows and + // their recipient provenance; the customer rows stay. + module := sendingpolicy.NewModule(f.pool, f.secrets()) + if _, err := module.GCFeedback(f.ctx, time.Now().Add(31*24*time.Hour), 7); err != nil { + t.Fatal(err) + } + for class, corrID := range correlation { + corr, recipients, _ := f.provenanceRows(corrID) + wantGone := class == "system" || class == "internal" + if gone := corr == 0 && recipients == 0; gone != wantGone { + t.Errorf("%s-class provenance after 31 days: correlations=%d recipients=%d, want gone=%v", class, corr, recipients, wantGone) + } + } +} + +// TestMigration129BackfillsOnlySyntheticAccounts: the one-time pass stamps +// created_at + 30 days on existing system/internal correlations (and their +// events) and nothing else, and is idempotent. +func TestMigration129BackfillsOnlySyntheticAccounts(t *testing.T) { + f := newFixture(t) + users := map[string]string{} + for _, class := range []string{"system", "internal", "standard", "demo"} { + users[class] = f.user(class) + } + seed := func(id, account, purpose string, expires any) { + f.exec(`INSERT INTO sending_feedback_correlations + (correlation_id, operation_id, submission_attempt, source_account_ref, policy_subject_ref, + purpose, tenant_mode, created_at, expires_at) + VALUES ($1, 'op_' || $1, 1, $2, COALESCE($2, 'system'), $3, 'none', now() - interval '10 days', $4)`, + id, account, purpose, expires) + f.exec(`INSERT INTO sending_feedback_events (provider_event_id, correlation_id, provider_occurred_at) + VALUES ('evt_' || $1, $1, now())`, id) + } + seed("cor_bf129_system", users["system"], "customer_message", nil) + seed("cor_bf129_internal", users["internal"], "customer_notification", nil) + seed("cor_bf129_standard", users["standard"], "customer_message", nil) + seed("cor_bf129_demo", users["demo"], "customer_message", nil) + stamped := time.Now().UTC().Add(5 * 24 * time.Hour).Truncate(time.Second) + seed("cor_bf129_system_stamped", users["system"], "customer_message", stamped) + + sql, err := migrations.FS.ReadFile("129_sending_feedback_synthetic_retention.sql") + if err != nil { + t.Fatal(err) + } + for pass := 1; pass <= 2; pass++ { + if _, err := f.pool.Exec(f.ctx, string(sql)); err != nil { + t.Fatalf("apply 129 (pass %d): %v", pass, err) + } + for id, want := range map[string]bool{ + "cor_bf129_system": true, "cor_bf129_internal": true, + "cor_bf129_standard": false, "cor_bf129_demo": false, + } { + var created time.Time + var expires, eventExpires *time.Time + if err := f.pool.QueryRow(f.ctx, ` + SELECT c.created_at, c.expires_at, e.expires_at + FROM sending_feedback_correlations c + JOIN sending_feedback_events e ON e.correlation_id = c.correlation_id + WHERE c.correlation_id = $1`, id).Scan(&created, &expires, &eventExpires); err != nil { + t.Fatal(err) + } + if !want { + if expires != nil || eventExpires != nil { + t.Errorf("pass %d: %s was stamped (%v / %v); only system/internal rows may be", pass, id, expires, eventExpires) + } + continue + } + if expires == nil || !expires.Equal(created.Add(30*24*time.Hour)) { + t.Errorf("pass %d: %s expiry = %v, want created_at + 30 days", pass, id, expires) + } + if eventExpires == nil || !eventExpires.Equal(*expires) { + t.Errorf("pass %d: %s event expiry = %v, want the correlation's %v", pass, id, eventExpires, expires) + } + } + if got := f.correlationExpiry("cor_bf129_system_stamped"); got == nil || !got.Equal(stamped) { + t.Errorf("pass %d: an already-stamped expiry was rewritten: %v", pass, got) + } + } +} diff --git a/internal/sendingpolicy/gate.go b/internal/sendingpolicy/gate.go index ac73cc3c7..c4706fe51 100644 --- a/internal/sendingpolicy/gate.go +++ b/internal/sendingpolicy/gate.go @@ -1356,8 +1356,16 @@ func (m *Module) recordCorrelation(ctx context.Context, tx pgx.Tx, st authState, // wait out a fixed timer before complaining. The post-deletion janitor sets // their horizon. A non-customer operation has no account to outlive, so it // receives the configured horizon at creation. + // + // So does customer-purpose traffic from a system or internal account + // (probers, monitors, conformance runs). Those accounts are never deleted, + // so "account lifetime" would mean forever, and their traffic is + // synthetic: the standing prober alone writes thousands of correlations a + // day. The class is the server-owned users.account_class, read under the + // FOR SHARE lock readAuthState already holds — the same value that exempts + // the account from budgets. Standard and demo accounts are unchanged. var expires *time.Time - if !st.op.Purpose.isCustomer() { + if !st.op.Purpose.isCustomer() || accountClassExempt(st.class) { horizon := time.Now().UTC().Add( time.Duration(st.policy.SendingFeedbackPostAcctRetention) * 24 * time.Hour) expires = &horizon diff --git a/migrations/129_sending_feedback_synthetic_retention.sql b/migrations/129_sending_feedback_synthetic_retention.sql new file mode 100644 index 000000000..9ec83fcdf --- /dev/null +++ b/migrations/129_sending_feedback_synthetic_retention.sql @@ -0,0 +1,43 @@ +-- 129_sending_feedback_synthetic_retention.sql +-- +-- Give feedback provenance of system- and internal-class accounts a fixed +-- 30-day horizon. +-- +-- The sending-policy gate writes customer-purpose correlations with +-- expires_at NULL: they are kept for the source account's lifetime, and the +-- account purge stamps the post-deletion horizon. System and internal +-- accounts (the standing prober, monitors, conformance) are never deleted, +-- so their correlations — nearly all of the table on a hosted deployment, +-- since the prober sends every 30 seconds — would be kept forever. From this +-- release the gate stamps them at creation (created_at + the post-account +-- retention, 30 days by default); this file is the one-time backlog pass for +-- rows written before it, using the same 30-day default +-- (sending_feedback_post_account_retention_days). +-- +-- Only rows whose source account's server-owned users.account_class is +-- 'system' or 'internal' are touched; standard and demo accounts keep the +-- account-lifetime rule. Events inherit their correlation's expiry; +-- recipient rows carry no expiry and are removed with their correlation. The +-- existing feedback retention janitor (GCFeedback) then deletes whatever is +-- already past its horizon on its next hourly pass. +-- +-- Idempotent: only NULL expiries are written. Row-level UPDATEs only; +-- lock_timeout bounds waiting on a row a concurrent feedback transaction +-- holds, failing the boot (which retries) rather than hanging it. + +SET LOCAL lock_timeout = '5s'; + +UPDATE sending_feedback_correlations c + SET expires_at = c.created_at + interval '30 days' + FROM users u + WHERE u.id = c.source_account_ref + AND u.account_class IN ('system', 'internal') + AND c.expires_at IS NULL + AND c.purpose IN ('customer_message', 'customer_notification'); + +UPDATE sending_feedback_events e + SET expires_at = c.expires_at + FROM sending_feedback_correlations c + WHERE c.correlation_id = e.correlation_id + AND e.expires_at IS NULL + AND c.expires_at IS NOT NULL; From 7b5aa4cdd22e5e485b1f927355c92273740f5c14 Mon Sep 17 00:00:00 2001 From: Josh Zhang <39790535+jiashuoz@users.noreply.github.com> Date: Wed, 30 Sep 2026 11:46:12 +0800 Subject: [PATCH 3/3] fix(sending): close ledger retention review gaps --- docs/data-handling.md | 2 +- internal/sendingpolicy/ledger_retention.go | 10 ++- .../ledger_retention_integration_test.go | 61 +++++++++++++++++++ internal/telemetry/prom.go | 5 ++ internal/telemetry/prom_test.go | 17 ++++++ ...9_sending_feedback_synthetic_retention.sql | 6 +- 6 files changed, 96 insertions(+), 5 deletions(-) diff --git a/docs/data-handling.md b/docs/data-handling.md index 2cc62b24d..88fe69ad3 100644 --- a/docs/data-handling.md +++ b/docs/data-handling.md @@ -23,7 +23,7 @@ For vulnerability reporting and the security model, see [SECURITY.md](../SECURIT | External sending access requests (use case, intended recipients, expected volume, review state) | Postgres `external_sending_access_requests` | Until the account is deleted (cascade). Private support data — never filed to a public tracker. | | External sending access audit (operator grant/revoke: account id, old/new state, actor, reason; no addresses) | Postgres `external_sending_access_events` | Append-only; no account foreign key, so it outlives account deletion. Each row carries an `expires_at` (the `sending_control_audit_retention_days` policy value, 90 days by default) and the sending-ledger janitor deletes it then. | | Sending pause/resume audit (account id, old/new state, reason, actor, pause class, optional operator evidence reference; no addresses) | Postgres `account_sending_control_events` | No account foreign key, so it outlives account deletion. Deleted by the sending-ledger janitor at its `expires_at` (`sending_control_audit_retention_days`, 90 days by default) — except an `abuse`-class event, which is kept past its expiry while the account still exists (live or trashed), because the account purge reads it to decide the abuse hold; it is deleted after the purge. | -| Sending ledger: provider operations and their submission attempts (opaque operation/account ids, purpose, shared-reputation flag, attempt ordinal, UTC day, recipient count, budget/call state, an HMAC commitment to a notice recipient; no addresses, subjects or bodies) | Postgres `sending_provider_operations`, `sending_budget_reservations` | No foreign keys, so message and account deletion never touch them. 30 days: each row is stamped with an `expires_at` at creation (an operation 30 days after creation or its scheduled send time, an attempt 30 days after it was first written), and the sending-ledger janitor deletes an operation together with all of its attempts once every one of them is past its expiry — but never while the send can still use it: an attempt still holding capacity (`reserved`) or holding an unredeemed provider authorization, a runnable River job that names the operation (a budget or pause hold, a retry), a customer message not yet handed to the provider, or a pending protection notice keeps it. | +| Sending ledger: provider operations and their submission attempts (opaque operation/account ids, purpose, shared-reputation flag, attempt ordinal, UTC day, recipient count, budget/call state, an HMAC commitment to a notice recipient; no addresses, subjects or bodies) | Postgres `sending_provider_operations`, `sending_budget_reservations` | No foreign keys, so message and account deletion never touch them. 30 days: each row is stamped with an `expires_at` at creation (an operation 30 days after creation or its scheduled send time, an attempt 30 days after it was first written), and the sending-ledger janitor deletes an operation together with all of its attempts once every one of them is past its expiry — but never while the send can still use it: an attempt still holding capacity (`reserved`) or holding a still-valid, unredeemed provider authorization, a runnable River job that names the operation (a budget or pause hold, a retry), a customer message not yet handed to the provider, or a pending protection notice keeps it. | | Sending feedback provenance (opaque operation/account ids, random correlation id, provider message id, keyed recipient HMACs, detector bucket, timestamps; no addresses, subjects or bodies) | Postgres `sending_feedback_correlations`, `sending_feedback_recipients`, `sending_feedback_events` | For a customer account's lifetime, then 30 days after the account is purged (`sending_feedback_post_account_retention_days`). Synthetic traffic — non-customer purposes, and customer-purpose mail from `system`/`internal` accounts such as probers and monitors, which are never deleted — keeps it for 30 days from creation. Removed by the hourly feedback retention janitor. | | Sending budget day counters (scope, opaque account id or `global`, UTC day, reserved/confirmed counts, limit) | Postgres `sending_budget_counters` | Deleted by the sending-ledger janitor once the day and the day after it have both closed (the current and previous UTC day are always kept). Only closed days are deleted, so deletion never gives back capacity. | | Sending protection notice outbox (opaque event/account/operation ids, closed kind/reason/scope/audience values, delivery state; no addresses) | Postgres `sending_protection_notice_events`, `sending_protection_notice_deliveries` | Survive user, message and agent deletion (their only foreign key is delivery → event). A finished event (every delivery sent, failed or skipped) is deleted with its deliveries 30 days after the event was created; an event with a pending delivery is never deleted by age. | diff --git a/internal/sendingpolicy/ledger_retention.go b/internal/sendingpolicy/ledger_retention.go index cc1a788d4..34d9dc63f 100644 --- a/internal/sendingpolicy/ledger_retention.go +++ b/internal/sendingpolicy/ledger_retention.go @@ -223,8 +223,11 @@ var ledgerSweeps = map[string]ledgerSweep{ // Provider operations with every attempt they own. An operation survives // its own expiry while ANY of these still needs it: // - an attempt that is not terminal (state reserved: capacity held; or - // call_state authorized: a token was issued and not yet redeemed or - // invalidated), or an attempt not yet past its own expiry; + // call_state authorized for the current or a future ordinal: a token + // may still be redeemed), or an attempt not yet past its own expiry; + // superseded authorized ordinals cannot redeem and do not hold an + // otherwise terminal operation forever; their confirmed exposure is + // retained in the day counters until those days close; // - a River job that can still run and names it in operation_ref — // this covers budget and pause holds (snoozed jobs), retries inside // SendRetryHorizon, and every notification kind; deleting it would @@ -257,7 +260,8 @@ var ledgerSweeps = map[string]ledgerSweep{ AND NOT EXISTS ( SELECT 1 FROM sending_budget_reservations r WHERE r.operation_id = o.operation_id - AND (r.expires_at > $1 OR r.state = 'reserved' OR r.call_state = 'authorized')) + AND (r.expires_at > $1 OR r.state = 'reserved' + OR (r.call_state = 'authorized' AND r.submission_attempt >= o.current_attempt))) AND o.operation_id NOT IN ( SELECT l.operation_id FROM live_refs l WHERE l.operation_id IS NOT NULL) AND NOT EXISTS ( diff --git a/internal/sendingpolicy/ledger_retention_integration_test.go b/internal/sendingpolicy/ledger_retention_integration_test.go index fed23ffcd..40c4a91c5 100644 --- a/internal/sendingpolicy/ledger_retention_integration_test.go +++ b/internal/sendingpolicy/ledger_retention_integration_test.go @@ -520,3 +520,64 @@ func TestLedgerRetentionRegistersOnTheMaintenanceQueue(t *testing.T) { t.Fatalf("periodic jobs = %d, want 3 (feedback retention, reconcile, ledger retention)", len(periodics)) } } + +// Superseded tokens can never open a socket again, even though their +// reservation preserves confirmed exposure and call_state=authorized. +func TestLedgerRetentionCollectsSupersededAuthorization(t *testing.T) { + f, m := ledgerFixture(t) + g := f.gate(enforcingPolicy(nil)) + agent := f.agent(f.user("standard")) + ref, first := f.prepareAndReserve(g, agent, 1) + _, oldAuth, err := g.ConsumeAttempt(f.ctx, first) + if err != nil || oldAuth == nil { + t.Fatalf("old auth: %v", err) + } + _, second, err := g.Reserve(f.ctx, ref) + if err != nil { + t.Fatal(err) + } + _, liveAuth, err := g.ConsumeAttempt(f.ctx, second) + if err != nil || liveAuth == nil { + t.Fatalf("live auth: %v", err) + } + if err := g.RedeemProviderCall(f.ctx, *oldAuth); !errors.Is(err, sendingpolicy.ErrAuthorizationInvalid) { + t.Fatalf("stale token: %v", err) + } + if err := g.RedeemProviderCall(f.ctx, *liveAuth); err != nil { + t.Fatal(err) + } + if err := g.SettleOperation(f.ctx, ref, sendingpolicy.SettlementProviderAccepted, "ses_synthetic_gc"); err != nil { + t.Fatal(err) + } + f.exec(`UPDATE messages SET delivery_status = 'delivered' WHERE id = $1`, ref.ID()) + f.exec(`UPDATE sending_provider_operations SET expires_at = now() - interval '1 hour' WHERE operation_id = $1`, ref.ID()) + f.exec(`UPDATE sending_budget_reservations SET expires_at = now() - interval '1 hour' WHERE operation_id = $1`, ref.ID()) + gcLedger(t, m) + if f.opExists(ref.ID()) || f.attempts(ref.ID()) != 0 { + t.Fatal("terminal operation retained forever by its superseded authorization") + } +} + +func TestLedgerRetentionSupersededAuthorizationKeepsOtherGuards(t *testing.T) { + f, m := ledgerFixture(t) + for _, tc := range []struct { + id string + attempt int + state, call, expiry string + }{ + {"op_stale_unexpired", 1, "confirmed", "authorized", "1 day"}, + {"op_current_authorized", 2, "confirmed", "authorized", "-1 hour"}, + {"op_future_authorized", 3, "confirmed", "authorized", "-1 hour"}, + {"op_old_reserved", 1, "reserved", "none", "-1 hour"}, + } { + f.op(tc.id, "-1 hour") + f.exec(`UPDATE sending_provider_operations SET current_attempt = 2 WHERE operation_id = $1`, tc.id) + f.attempt(tc.id, tc.attempt, tc.state, tc.call, tc.expiry) + } + gcLedger(t, m) + for _, id := range []string{"op_stale_unexpired", "op_current_authorized", "op_future_authorized", "op_old_reserved"} { + if !f.opExists(id) || f.attempts(id) != 1 { + t.Errorf("guarded operation %s lost", id) + } + } +} diff --git a/internal/telemetry/prom.go b/internal/telemetry/prom.go index 33d7f9d47..2c85b8380 100644 --- a/internal/telemetry/prom.go +++ b/internal/telemetry/prom.go @@ -465,6 +465,11 @@ func NewProm(build string) *Prom { p.outboxPublished, p.outboxFanOut, p.outboxMatched, p.outboxNoMatch, p.outboxFailures, p.redeliver, p.janitorDeleted, p.contactDue, p.notifyMissed, p.publisherLag, ) + // Export a zero baseline before the first run so increase() can count + // the first failure after a scrape. These outcomes are a closed set. + for _, outcome := range []string{"complete", "partial", "failed"} { + p.sendingLedgerRuns.WithLabelValues(outcome) + } return p } diff --git a/internal/telemetry/prom_test.go b/internal/telemetry/prom_test.go index b1dafa268..8a8f545c6 100644 --- a/internal/telemetry/prom_test.go +++ b/internal/telemetry/prom_test.go @@ -472,3 +472,20 @@ func itoa(i int) string { } return string(b[pos:]) } + +// A zero baseline lets increase() include the first failure after a scrape, +// instead of treating that failure as the counter's initial value. +func TestPromInitialLedgerRetentionOutcomes(t *testing.T) { + p := NewProm("") + out := scrape(t, p) + for _, outcome := range []string{"complete", "partial", "failed"} { + want := `e2a_sending_ledger_retention_runs_total{outcome="` + outcome + `"} 0` + if !strings.Contains(out, want) { + t.Errorf("initial scrape missing zero baseline: %s", want) + } + } + p.SendingLedgerRetentionRun("failed") + if !strings.Contains(scrape(t, p), `e2a_sending_ledger_retention_runs_total{outcome="failed"} 1`) { + t.Error("first failure must increment the primed counter to one") + } +} diff --git a/migrations/129_sending_feedback_synthetic_retention.sql b/migrations/129_sending_feedback_synthetic_retention.sql index 9ec83fcdf..47cf87311 100644 --- a/migrations/129_sending_feedback_synthetic_retention.sql +++ b/migrations/129_sending_feedback_synthetic_retention.sql @@ -23,9 +23,13 @@ -- -- Idempotent: only NULL expiries are written. Row-level UPDATEs only; -- lock_timeout bounds waiting on a row a concurrent feedback transaction --- holds, failing the boot (which retries) rather than hanging it. +-- holds, failing the boot (which retries) rather than hanging it. Each UPDATE +-- also has a 30-second statement budget; a timeout rolls the migration back. +-- This is a single-transaction backfill, not a batched sweep. If the backlog +-- cannot fit the budget, stop rollout and arrange a reviewed batched backfill. SET LOCAL lock_timeout = '5s'; +SET LOCAL statement_timeout = '30s'; UPDATE sending_feedback_correlations c SET expires_at = c.created_at + interval '30 days'