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
4 changes: 4 additions & 0 deletions cmd/e2a/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -1163,6 +1163,10 @@ func main() {
// (The auto-disable janitor is now a River periodic on QueueMaintenance; the
// legacy SubscriberRetryWorker is gone.)
bgCtx, bgCancel := context.WithCancel(context.Background())
if promBackend != nil {
workerWG.Add(1)
go func() { defer workerWG.Done(); observeSendingBudgets(bgCtx, outboundSending.module, metrics) }()
}

// Outbox publisher worker: drains webhook_events → subscriber_deliveries and
// enqueues River delivery jobs in-tx. Skipped when fan-out runs on River
Expand Down
1 change: 1 addition & 0 deletions cmd/e2a/outbound_wiring.go
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,7 @@ func installSendingPolicyObservers(m telemetry.Metrics) {
// Ledger retention: rows deleted per ledger table (on the shared janitor
// counter) and one run-outcome sample per pass.
sendingpolicy.SetLedgerRetentionObserver(m)
sendingpolicy.SetBudgetObserver(m.SendingBudgetDecision)
}

// nonEmpty returns the non-blank values, so an unset config string does not
Expand Down
33 changes: 33 additions & 0 deletions cmd/e2a/sending_observation.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
package main

import (
"context"
"log"
"time"

"github.com/tokencanopy/e2a/internal/sendingpolicy"
"github.com/tokencanopy/e2a/internal/telemetry"
)

// Run per process, not as a shared River periodic: every serving slot must
// publish its own effective policy and freshness. Never put DB I/O on /metrics.
func observeSendingBudgets(ctx context.Context, module *sendingpolicy.Module, metrics telemetry.Metrics) {
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
sampleCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
snapshot, err := module.BudgetSnapshot(sampleCtx)
cancel()
if err == nil {
metrics.SendingBudgetSnapshot(snapshot.UsedRatio, snapshot.Generation, snapshot.ConfigMismatch, float64(snapshot.ObservedAt.Unix()))
} else if ctx.Err() == nil {
// DB errors may include row values; no raw query/error text in this log.
log.Printf("[sending-protection] budget observation failed; gauges retain their last successful sample")
}
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}
52 changes: 52 additions & 0 deletions cmd/e2a/sending_observation_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
package main

import (
"context"
"testing"
"time"

"github.com/tokencanopy/e2a/internal/sendingpolicy"
"github.com/tokencanopy/e2a/internal/telemetry"
"github.com/tokencanopy/e2a/internal/testutil/testdb"
)

type budgetSnapshotRecorder struct {
telemetry.NoOp
samples chan int
cancel context.CancelFunc
}

func (r *budgetSnapshotRecorder) SendingBudgetSnapshot(used map[string]float64, _ int64, _ bool, _ float64) {
r.samples <- len(used)
r.cancel()
}

func TestSendingBudgetSamplerPublishesAndStops(t *testing.T) {
pool := testdb.TestDB(t)
module := sendingpolicy.NewPolicyModule(pool, sendingpolicy.Secrets{}, sendingpolicy.PolicySourceConfig, sendingpolicy.DisabledPolicy())
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
recorder := &budgetSnapshotRecorder{samples: make(chan int, 1), cancel: cancel}
done := make(chan struct{})
go func() { defer close(done); observeSendingBudgets(ctx, module, recorder) }()
select {
case <-done:
case <-time.After(10 * time.Second):
t.Fatal("sampler did not stop")
}
select {
case count := <-recorder.samples:
if count != 4 {
t.Fatalf("got %d scopes", count)
}
default:
t.Fatal("no sample")
}
// A failed sample must not overwrite the prior gauges with healthy zeros.
observeSendingBudgets(ctx, module, recorder)
select {
case <-recorder.samples:
t.Fatal("canceled read published a snapshot")
default:
}
}
43 changes: 43 additions & 0 deletions docs/observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,49 @@ Ambiguous anchors also emit a structured, process-wide rate-limited log (at
most one line per minute). It contains only candidate/thread counts: no
addresses, subjects, message content, or RFC Message-IDs.

### Sending budget observation

| Metric | Type | Labels | Meaning |
|---|---|---|---|
| `e2a_sending_budget_decisions_total` | counter | `scope`, `decision` | Committed gate evaluations. Decisions: `allow`, `would_hold` (shadow), `hold` (enforced). |
| `e2a_sending_budget_deferrals_total` | counter | `scope` | Enforced budget/trust holds; the same events as decisions with `decision="hold"`. |
| `e2a_sending_budget_used_ratio` | gauge | `scope` | Current UTC day's reserved units (including confirmed units), divided by the effective policy limit. Global pools only. Shadow demand can exceed 1. |
| `e2a_sending_protection_policy_generation` | gauge | — | Effective database policy generation; 0 for config source. |
| `e2a_sending_protection_policy_config_mismatch` | gauge | — | Database policy differs from local config (1), or matches/config source (0). Intentional database activation can produce a mismatch. |
| `e2a_sending_budget_observation_last_success_timestamp_seconds` | gauge | — | Unix timestamp of the last successful policy/counter snapshot; 0 before the first successful sample. |

Decision scopes are the closed set `global_all`, `account_daily`,
`account_shared_daily`, `global_probation`, `global_critical`, and
`global_violation`. Unknown labels normalize to `other`. No account, operation,
message, domain, recipient, or raw policy value is a label.

Reserve reports early holds only; Consume reports each evaluated scope. Shadow
mode reports every exceeded pool. Enforced evaluation stops at its first budget
hold, so an earlier pool's `allow` does not mean the whole send was allowed.
Account trust checks also report account/shared evaluations when enabled,
including when platform budgets are disabled. If both legacy account budgets
and account trust are enabled, both checks can report the same scope. Disabled
legacy account/shared/probation checks emit no decisions. These counters count
evaluations, not recipients, unique messages, or SMTP calls: retries can add
samples, and a crash after commit but before emission can lose a sample.
Rolled-back transactions do not emit samples. The former shadow-denial log
containing operation identifiers is replaced by these bounded observations.

Each process with Prometheus enabled samples immediately and every 30 seconds,
with a 5-second timeout. One read-only database snapshot reads the effective
policy and four indexed global-pool counters; it never scans account rows.
The UTC date and observation timestamp come from PostgreSQL, matching the
clock used by authorization. Missing counters report zero, and disabled legacy
probation reports zero.
Ratios use the current policy limit, not a historical limit cached in the
counter row. `/metrics` performs no database queries. A failed sample retains
the previous gauges and timestamp; zero or stale freshness is not evidence of
unused capacity. Filter out stale samples before evaluating usage, and
aggregate these shared-ledger ratios across replicas with **max**, not sum.
Inspect generation, mismatch, and freshness separately for each serving slot.
These observations do not change policy, enable enforcement, or establish
readiness of the operator-recipient registry.

### Maintenance

| Metric | Type | Labels | Meaning |
Expand Down
56 changes: 56 additions & 0 deletions internal/sendingpolicy/budget_metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
package sendingpolicy

import (
"errors"
"sync/atomic"
)

// BudgetObserver receives committed budget evaluations, never identifiers or
// recipient counts. Reserve reports only holds; Consume reports evaluated
// scopes. These are gate evaluations, not unique messages or provider sends.
type BudgetObserver func(scope, decision string)

var budgetObserver atomic.Value // BudgetObserver

// SetBudgetObserver installs the process-wide observer. nil disables it.
func SetBudgetObserver(o BudgetObserver) { budgetObserver.Store(o) }

type budgetSample struct {
scope Scope
decision string
}

func emitBudgetSamples(samples []budgetSample) {
if o, ok := budgetObserver.Load().(BudgetObserver); ok && o != nil {
for _, s := range samples {
o(string(s.scope), s.decision)
}
}
}

func observeBudgetScope(p RuntimePolicy, scope Scope) bool {
return !(p.DisableLegacyDailyBudgets && (scope == ScopeAccountDaily || scope == ScopeAccountSharedDaily || scope == ScopeGlobalProbation))
}

func appendAccountTrustSamples(samples *[]budgetSample, st authState, err error) {
if !st.ramp.account || !st.ramp.applies || st.ramp.units == 0 {
return
}
if err == nil {
*samples = append(*samples, budgetSample{ScopeAccountDaily, "allow"})
if st.ramp.shared {
*samples = append(*samples, budgetSample{ScopeAccountSharedDaily, "allow"})
}
return
}
var capacity *accountCapacityError
if errors.As(err, &capacity) {
scope := ScopeAccountDaily
// Reservation results carry usage, while SharedBinding is a preflight-only hint.
d := capacity.daily
if st.ramp.shared && d.SharedLimit-d.SharedUsed < d.Limit-d.Used {
scope = ScopeAccountSharedDaily
}
*samples = append(*samples, budgetSample{scope, "hold"})
}
}
163 changes: 163 additions & 0 deletions internal/sendingpolicy/budget_metrics_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,163 @@
package sendingpolicy_test

import (
"testing"

"github.com/tokencanopy/e2a/internal/sendingpolicy"
)

func TestBudgetObservationShadowRecordsEveryExceededScope(t *testing.T) {
f := newFixture(t)
samples := map[string]int{}
sendingpolicy.SetBudgetObserver(func(scope, decision string) { samples[scope+"/"+decision]++ })
t.Cleanup(func() { sendingpolicy.SetBudgetObserver(nil) })
p := enforcingPolicy(func(p *sendingpolicy.RuntimePolicy) {
p.BudgetMode = sendingpolicy.ModeShadow
p.AllCustomerGlobalDailyRecipients = 1
p.DefaultAccountDailyRecipients = 1
p.SharedDomainAccountDailyRecip = 1
p.ProbationGlobalDailyRecipients = 1
})
g := f.gate(p)
user := f.user("standard")
agent := f.agent(user)
for i := 0; i < 3; i++ {
if d := f.send(g, f.message(agent, "relay", 1)); !d.Allow {
t.Fatal(d)
}
}
for _, scope := range []string{"global_all", "account_daily", "account_shared_daily", "global_probation"} {
if samples[scope+"/allow"] != 1 || samples[scope+"/would_hold"] != 2 {
t.Fatalf("samples: %v", samples)
}
}
if len(samples) != 8 {
t.Fatalf("unexpected observations: %v", samples)
}
}

func TestBudgetObservationEarlyHoldAndDisabled(t *testing.T) {
f := newFixture(t)
samples := map[string]int{}
sendingpolicy.SetBudgetObserver(func(scope, decision string) { samples[scope+"/"+decision]++ })
t.Cleanup(func() { sendingpolicy.SetBudgetObserver(nil) })
user := f.user("standard")
agent := f.agent(user)
g := f.gate(enforcingPolicy(func(p *sendingpolicy.RuntimePolicy) { p.AllCustomerGlobalDailyRecipients = 1 }))
_, ref := f.prepareMessage(g, f.message(agent, "relay", 2))
d, _, err := g.Reserve(f.ctx, ref)
if err != nil || d.Allow {
t.Fatalf("reserve: %+v %v", d, err)
}
if len(samples) != 1 || samples["global_all/hold"] != 1 {
t.Fatalf("early hold: %v", samples)
}
if d := f.send(f.gate(sendingpolicy.DisabledPolicy()), f.message(agent, "relay", 1)); !d.Allow {
t.Fatal(d)
}
if len(samples) != 1 {
t.Fatalf("disabled emitted: %v", samples)
}
}

func TestBudgetObservationDoesNotPublishRolledBackHold(t *testing.T) {
f := newFixture(t)
count := 0
sendingpolicy.SetBudgetObserver(func(string, string) { count++ })
t.Cleanup(func() { sendingpolicy.SetBudgetObserver(nil) })
g := f.gate(enforcingPolicy(func(p *sendingpolicy.RuntimePolicy) { p.AllCustomerGlobalDailyRecipients = 1 }))
_, ref := f.prepareMessage(g, f.message(f.agent(f.user("standard")), "relay", 2))
_, err := f.pool.Exec(f.ctx, `CREATE FUNCTION reject_budget_metric_commit() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'synthetic commit rejection'; END $$;
CREATE CONSTRAINT TRIGGER reject_budget_metric_commit AFTER INSERT OR UPDATE ON sending_budget_reservations DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION reject_budget_metric_commit()`)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
_, _ = f.pool.Exec(f.ctx, `DROP TRIGGER IF EXISTS reject_budget_metric_commit ON sending_budget_reservations; DROP FUNCTION IF EXISTS reject_budget_metric_commit()`)
})
if _, _, err = g.Reserve(f.ctx, ref); err == nil {
t.Fatal("expected commit failure")
}
if count != 0 {
t.Fatalf("rolled-back hold emitted %d samples", count)
}
}

func TestBudgetObservationAccountTrustWithLegacyBudgetsDisabled(t *testing.T) {
f := newFixture(t)
samples := map[string]int{}
sendingpolicy.SetBudgetObserver(func(scope, decision string) { samples[scope+"/"+decision]++ })
t.Cleanup(func() { sendingpolicy.SetBudgetObserver(nil) })
p := sendingpolicy.DisabledPolicy()
p.AccountTrustEnabled = true
p.DisableLegacyDailyBudgets = true
g := f.gate(p)
user := f.user("standard")
f.plan(user, "scale")
agent := f.agent(user)
if d := f.send(g, f.message(agent, "relay", 20)); !d.Allow {
t.Fatal(d)
}
if d := f.send(g, f.message(agent, "relay", 1)); d.Allow {
t.Fatal("expected trust hold")
}
if samples["account_daily/allow"] != 1 || samples["account_shared_daily/allow"] != 1 || samples["account_daily/hold"] != 1 {
t.Fatalf("trust samples: %v", samples)
}
if len(samples) != 3 {
t.Fatalf("disabled platform budgets emitted: %v", samples)
}
}

func TestBudgetObservationDoesNotPublishFailedConsume(t *testing.T) {
f := newFixture(t)
count := 0
sendingpolicy.SetBudgetObserver(func(string, string) { count++ })
t.Cleanup(func() { sendingpolicy.SetBudgetObserver(nil) })
off := f.gate(sendingpolicy.DisabledPolicy())
_, ref := f.prepareMessage(off, f.message(f.agent(f.user("standard")), "relay", 1))
_, attempt, err := off.Reserve(f.ctx, ref)
if err != nil {
t.Fatal(err)
}
_, err = f.pool.Exec(f.ctx, `CREATE FUNCTION reject_budget_consume_commit() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'synthetic consume commit rejection'; END $$;
CREATE CONSTRAINT TRIGGER reject_budget_consume_commit AFTER UPDATE ON sending_budget_reservations DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION reject_budget_consume_commit()`)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
_, _ = f.pool.Exec(f.ctx, `DROP TRIGGER IF EXISTS reject_budget_consume_commit ON sending_budget_reservations; DROP FUNCTION IF EXISTS reject_budget_consume_commit()`)
})
if _, _, err = f.gate(enforcingPolicy(nil)).ConsumeAttempt(f.ctx, attempt); err == nil {
t.Fatal("expected commit failure")
}
if count != 0 {
t.Fatalf("failed consume emitted %d samples", count)
}
}

func TestBudgetObservationSharedTrustCeiling(t *testing.T) {
f := newFixture(t)
samples := map[string]int{}
sendingpolicy.SetBudgetObserver(func(scope, decision string) { samples[scope+"/"+decision]++ })
t.Cleanup(func() { sendingpolicy.SetBudgetObserver(nil) })
p := sendingpolicy.DisabledPolicy()
p.AccountTrustEnabled = true
p.DisableLegacyDailyBudgets = true
g := f.gate(p)
user := f.user("standard")
f.plan(user, "scale")
if _, err := f.pool.Exec(f.ctx, `INSERT INTO account_sending_trust(user_id,archived_clean_days) VALUES($1,30)`, user); err != nil {
t.Fatal(err)
}
agent := f.agent(user)
if d := f.send(g, f.message(agent, "relay", 50)); !d.Allow {
t.Fatal(d)
}
if d := f.send(g, f.message(agent, "relay", 1)); d.Allow {
t.Fatal("expected shared ceiling hold")
}
if samples["account_shared_daily/hold"] != 1 || samples["account_daily/hold"] != 0 {
t.Fatalf("shared ceiling mislabeled: %v", samples)
}
}
Loading
Loading