Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 23 additions & 2 deletions cmd/e2a/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -352,6 +352,9 @@ func main() {
// 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)
outboxWorker := webhookpub.NewOutboxWorker(pool, store).WithMetrics(metrics)
smtpRelay := outbound.NewSMTPRelay(&cfg.OutboundSMTP)
sender := outbound.NewSenderWithDKIM(smtpRelay, cfg.OutboundSMTP.FromDomain, store)
Expand Down Expand Up @@ -421,14 +424,31 @@ func main() {
// window, durable in Postgres): the cross-replica counterpart of the
// acceptance-time in-memory limiter, enforced immediately before
// provider submission so scheduled-send bursts can't exceed it.
rate: sendrate.NewStore(pool, time.Minute, 60),
rate: sendrate.NewStore(pool, time.Minute, 60),
sharedDomains: nonEmpty(cfg.SharedDomain),
})
outboundJobs := outboundSending.jobs
registrars = append(registrars, outboundJobs)
// Platform mail the API sends itself (public feedback) crosses the same
// seam with tokens from the same gate.
sendingGate, providerSubmitter := outboundSending.gate, outboundSending.submitter
registrars = append(registrars, sendramp.NewMaintenanceJobs(rampStore))
// Deletion-resistant feedback provenance (B8): refuse to start if any
// retained recipient row was signed under a key version this keyring
// does not hold — feedback for it could never be matched and the
// detector would be silently blind. The retention janitor and the
// consumer's accounting seam hang off the same module.
if outboundSending.module == nil {
log.Fatalf("sending policy gate is not the concrete module; feedback accounting cannot be wired")
}
if err := outboundSending.module.VerifyKeyringCoverage(ctx); err != nil {
log.Fatalf("Sending feedback keyring coverage: %v", err)
}
// The purge seal resolves the post-deletion horizon at purge time from
// the EFFECTIVE policy — the same accessor the retention janitor reads —
// not from the config-file policy captured here at boot.
store.SetFeedbackRetentionResolver(outboundSending.module.EffectiveFeedbackRetention)
registrars = append(registrars, outboundSending.feedbackMaintenance())
// Queue depth/age gauges: a 30s maintenance periodic sampling river_job
// per queue+state (docs/observability.md).
registrars = append(registrars, jobs.NewQueueStatsJobs(pool, metrics))
Expand Down Expand Up @@ -987,7 +1007,8 @@ func main() {
// 4b). Fail-closed: the SNS signature is verified and the TopicArn must be
// in the configured allow-list (empty allow-list → every message is
// rejected, so this is inert until ops wires the topic).
deliveryConsumer := delivery.NewConsumer(store, deliveryEventFirer(webhookOutbox), outboundSendStore.FinalizeProviderAcceptedTx)
deliveryConsumer := outboundSending.armDeliveryConsumer(
delivery.NewConsumer(store, deliveryEventFirer(webhookOutbox), outboundSendStore.FinalizeProviderAcceptedTx))
deliveryVerifier := delivery.NewVerifier(cfg.DeliveryFeedback.SNSTopicARNs, delivery.HTTPCertFetcher)
// Public webhook receiver for AWS SNS (SES delivery/bounce/complaint). Named
// /webhooks/<provider> — it's an inbound third-party callback, not an internal
Expand Down
41 changes: 39 additions & 2 deletions cmd/e2a/outbound_wiring.go
Original file line number Diff line number Diff line change
@@ -1,9 +1,12 @@
package main

import (
"strings"

"github.com/jackc/pgx/v5/pgxpool"

"github.com/tokencanopy/e2a/internal/agent"
"github.com/tokencanopy/e2a/internal/delivery"
"github.com/tokencanopy/e2a/internal/hitlnotify"
"github.com/tokencanopy/e2a/internal/identity"
"github.com/tokencanopy/e2a/internal/outbound"
Expand All @@ -25,13 +28,19 @@ type outboundSendingDeps struct {
sesConfigSet string
metrics outboundsend.Metrics
rate outboundsend.RateGate
// sharedDomains are the deployment's shared agent domains (config
// shared_domain). Deliveries to agents hosted on them never count toward
// the outcome detector's denominator.
sharedDomains []string
}

// outboundSending is the composed outbound send path.
type outboundSending struct {
gate sendingpolicy.Gate
// module is the same policy owner behind gate, exposed through its other
// narrow roles (external-sending-access preflight/status/requests).
// narrow roles: external-sending-access preflight/status/requests, the
// deletion-resistant feedback processor, the keyring coverage check, and
// the retention janitor all hang off it.
module *sendingpolicy.Module
submitter *outbound.ProviderSubmitter
jobs *outboundsend.Jobs
Expand All @@ -44,7 +53,8 @@ type outboundSending struct {
// enqueue and authorizes every worker execution through the same gate. No
// raw sender and no direct ramp store reach the worker from here.
func newOutboundSending(d outboundSendingDeps) outboundSending {
module := sendingpolicy.NewPolicyModule(d.pool, d.secrets, d.source, d.policy)
module := sendingpolicy.NewPolicyModule(d.pool, d.secrets, d.source, d.policy).
WithFeedbackExcludedDomains(d.sharedDomains...)
var gate sendingpolicy.Gate = module
submitter := outbound.NewProviderSubmitter(d.relay, gate)
// Delivery feedback: tag outbound with the SES configuration set so SES
Expand Down Expand Up @@ -99,3 +109,30 @@ func (s outboundSending) armAPI(api *agent.API) {
api.SetProviderSubmitter(s.submitter, s.gate)
api.SetExternalAccess(s.module)
}

// armDeliveryConsumer installs the deletion-resistant accounting seam on the
// SES feedback consumer. Without it the consumer still acks and still runs
// the message lifecycle, so the omission is silent: provider evidence for a
// purged message is simply dropped and the detector reads as healthy.
func (s outboundSending) armDeliveryConsumer(c *delivery.Consumer) *delivery.Consumer {
return c.WithFeedbackProcessor(s.module)
}

// feedbackMaintenance is the retention janitor for feedback provenance and
// daily outcome aggregates. Unregistered, nothing enforces the
// post-deletion horizon.
func (s outboundSending) feedbackMaintenance() *sendingpolicy.MaintenanceJobs {
return sendingpolicy.NewMaintenanceJobs(s.module)
}

// nonEmpty returns the non-blank values, so an unset config string does not
// become an empty-domain entry.
func nonEmpty(values ...string) []string {
var out []string
for _, v := range values {
if strings.TrimSpace(v) != "" {
out = append(out, v)
}
}
return out
}
45 changes: 45 additions & 0 deletions cmd/e2a/sending_policy_wiring_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (

"github.com/tokencanopy/e2a/internal/agent"
"github.com/tokencanopy/e2a/internal/config"
"github.com/tokencanopy/e2a/internal/delivery"
"github.com/tokencanopy/e2a/internal/outbound"
"github.com/tokencanopy/e2a/internal/sendingpolicy"
"github.com/tokencanopy/e2a/internal/testutil/testdb"
Expand Down Expand Up @@ -46,6 +47,9 @@ func TestSendingPolicyWiring(t *testing.T) {
if got := composed.submitter.SESConfigurationSet(); got != "e2a-delivery-test" {
t.Fatalf("submitter configuration set = %q, want the deployment's — delivery feedback must stay on", got)
}
if composed.module == nil || sendingpolicy.Gate(composed.module) != composed.gate {
t.Fatal("the concrete module (feedback processor, keyring coverage, retention janitor) is not the composed gate")
}
if composed.jobs.Gate() != composed.gate {
t.Fatal("the jobs bundle does not hold the composed gate")
}
Expand Down Expand Up @@ -125,3 +129,44 @@ func TestNotificationAndPlatformMailWiring(t *testing.T) {
t.Fatal("armAPI did not hand the API the submitter and gate")
}
}

// TestFeedbackAccountingWiring pins the two composition-root edges that are
// invisible at runtime when they are missing. A delivery consumer without
// the accounting seam still acks every SES notification and still runs the
// message lifecycle, so nothing fails — the detector simply never receives
// evidence. An unregistered retention janitor never enforces the
// post-deletion horizon, so provenance accumulates forever. Both are one
// line in main.go, and neither has any other test.
func TestFeedbackAccountingWiring(t *testing.T) {
pool := testdb.TestDB(t)
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(),
// main passes nonEmpty(cfg.SharedDomain).
sharedDomains: nonEmpty("agents.localhost", " "),
})
if !composed.module.ExcludesFeedbackDomain("agents.localhost") {
t.Fatal("the shared agent domain did not reach the feedback denominator exclusion")
}

bare := delivery.NewConsumer(nil, nil)
if bare.FeedbackProcessorWired() {
t.Fatal("a bare consumer must not claim an accounting seam")
}
if armed := composed.armDeliveryConsumer(bare); !armed.FeedbackProcessorWired() {
t.Fatal("the composition root did not install the feedback processor on the delivery consumer")
}

janitor := composed.feedbackMaintenance()
if janitor == nil {
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))
}
}
111 changes: 108 additions & 3 deletions docs/design/async-message-pipeline.md
Original file line number Diff line number Diff line change
Expand Up @@ -363,12 +363,117 @@ of Prepare itself, not of every caller.

Two consequences worth knowing. Notification and feedback mail now cross the
same submitter as customer mail, so it carries `X-SES-CONFIGURATION-SET`
and SES publishes delivery feedback for it; none of it correlates to a
message row, and the SNS consumer acks it as unknown (a log line, no
suppression). And the closure guard fences `net/smtp` and the SES v2 SDK
and SES publishes delivery feedback for it. None of it correlates to a
message row, so the SNS consumer's message-lifecycle half still acks it as
unknown; since B8 its retained sending correlation does match, which feeds
the detector but never repairs a suppression (see that addendum). And the closure guard fences `net/smtp` and the SES v2 SDK
import; a send through some other HTTP provider API would be a new
dependency, which is where review catches it.

## Addendum (2026-09-07): deletion-resistant feedback provenance (B8)

Slice B8 makes provider feedback count even when the message, the agent, or
the whole account it belongs to is gone. The SNS consumer now runs a
deletion-resistant accounting seam (`delivery.FeedbackProcessor`, implemented
by the sending-policy module) BEFORE it looks for a live message, in its own
transaction, and fails the notification if that accounting fails so the
provider retries. The seam never reads `messages`, `agent_identities`, or
`users`:

- **Correlation** is by the SES message id bound at settlement (normalized
to SES's bare form), then by the random `X-E2A-Provider-Attempt` marker SES
echoes from the submitted headers. Each recipient the event names is
matched against the keyed HMACs recorded at authorization; an address
outside the authorized envelope proves nothing and is ignored.
- **One row per provider event id** (`sending_feedback_events`), so a
redelivered notification is a zero delta.
- **Buckets** are derived from the full kind and retained subtypes:
`delivered`; `terminal_other` for a transient or undetermined bounce;
`hard_bounce` for a permanent bounce, including the global-list subtype
`Suppressed`; `complaint` for a genuine complaint. The account- and
tenant-suppression-list subtypes are `none` (SES never attempted delivery)
but still repair the local suppression list, and so does a complaint with
a suppression-list `complaintSubType`.
- **Evidence is monotonic per authorized recipient**: `none < delivered <
terminal_other < hard_bounce < complaint`. Higher evidence replaces lower
(the prior bucket is subtracted from the epoch/day it was counted in and
the new one added to the account's current `outcome_epoch` and the
ingestion UTC day, on the correlation's immutable shared/dedicated path);
a delayed lower-ranked callback never erases a hard bounce or complaint;
equal rank is a zero delta with provider time then event id deciding the
stored provenance.
- **Suppression ownership stays with the live message.** The seam only
*reports* the repairs an event proves; when a message row survives, the
consumer writes the suppression inside its own transaction exactly as
before, keeping the `source_message_id` and diagnostic reason the
suppression API returns and keeping the row's insert atomic with the
`suppression_added` event that announces it. Only when no message can own
the row — the purged-message case this slice exists for, including a
message purged between correlation and the live transaction's lock — does
the consumer apply the repair through the seam, in one transaction with a
`suppression.added` (no message id) for each row it actually inserted; an
address already suppressed is refreshed but not re-announced. Repair is
limited to customer messages: a bounce on an approval notice must not
suppress the account owner's own address. Upserts go through
`internal/suppressionsync`, which advances the row's `sync_generation` and
clears `removal_pending`; that guard has no production remover yet and is
the contract Task 11's provider reconciliation will build on.
- **After account deletion** feedback still advances the retained bucket
provenance but recreates no customer state. The 30-day retention horizon is
stamped on the account's correlations and events at **purge** (the seal
transaction that makes the account irrecoverable), not at the user's delete
click: a trashed account stays restorable for the trash window and is not
stamped, so provenance for a self-deleted account can live for the trash
window plus 30 days after purge. The horizon is read at purge time from the
effective (database-source when configured) policy. Correlation lookups take
`FOR SHARE` and skip rows past their horizon, so feedback racing the seal
inherits its expiry and a concurrent janitor delete cannot produce a
spurious "uncorrelated" result. The hourly `sending_feedback_maintenance`
job removes expired correlations together with their recipients and events,
and daily outcome rows older than the detector window plus one day. The
daily `sending_feedback_reconcile` job stamps the horizon on customer
provenance whose account no longer exists (a purge before B8, or a
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.
- **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
while the old stays active, then move the active version, and drop the old
key only once no retained row references it. Because a live account's rows
never expire, a version that has signed for a live account is effectively
permanent — treat the keyring as append-only. A deployment with no keyring
at all is not checked: it signs and matches nothing by design, and
bricking it over rows an earlier configuration wrote would turn a disabled
feature into an outage.

- **Denominator exclusion.** A delivery or terminal non-hard bounce to a
recipient a sender can generate at will is bucket `none`: the SES mailbox
simulator, the configured `shared_domain`, platform-owned verified
`domains` rows, and the sending account's own verified domains. Hard
bounces and complaints from those recipients still count. Known limits:
the match is on the exact recipient domain (a subdomain of an excluded
domain still counts), domains verified by *other* accounts still count
(a two-account dilution is not closed), and there is no DNS/MX lookup at
ingestion, so a customer domain whose MX points at this deployment is not
recognized as hosted.
- **IDN recipients.** Authorization accepts internationalized domains as
typed; both the recipient HMAC and feedback matching use the IDNA
lookup-profile ASCII form of the domain (raw form as a fallback for rows
signed earlier), and the send-time suppression lookup checks the A-label
and Unicode spellings, so a suppression repaired from provider feedback in
A-label form still blocks a Unicode-typed recipient.
- **Metric.** `e2a_sending_feedback_ingested_total{outcome,bucket}` (see
`docs/observability.md`); `uncorrelated_with_marker` is the alert signal,
and `dead_account` with bucket `complaint` is the operator's cue to consider
a manual `-escalate-deleted-account-to-abuse`.

The detector itself (thresholds, pauses, notices) is B9. B8 captures
evidence; the one new customer-visible behavior is that a hard bounce or
complaint for a message that has since been purged now repairs the account
suppression, where before it was acked and dropped. Suppressions for live
messages keep their existing shape and events exactly.

## Addendum (2026-09-26): external sending access

`internal/sendingpolicy/external_access.go` adds one permission decision that
Expand Down
Loading
Loading