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
9 changes: 3 additions & 6 deletions cmd/e2a/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
25 changes: 22 additions & 3 deletions cmd/e2a/outbound_wiring.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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 {
Expand Down
68 changes: 66 additions & 2 deletions cmd/e2a/sending_policy_wiring_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,18 @@ import (
"context"
"fmt"
"strings"
"sync"
"testing"

"github.com/riverqueue/river"

"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"
)
Expand Down Expand Up @@ -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)
}
}
Loading
Loading