From 229eabf46958cbcd0cf3a5ffc117f8ffc8cf1cab Mon Sep 17 00:00:00 2001 From: Josh Zhang <39790535+jiashuoz@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:41:05 +0800 Subject: [PATCH 1/2] feat(sending): observe budget decisions and global ledger usage --- cmd/e2a/main.go | 4 + cmd/e2a/outbound_wiring.go | 1 + cmd/e2a/sending_observation.go | 33 +++++ cmd/e2a/sending_observation_test.go | 52 +++++++ docs/observability.md | 41 ++++++ internal/sendingpolicy/budget_metrics.go | 54 +++++++ .../budget_metrics_integration_test.go | 137 +++++++++++++++++ internal/sendingpolicy/budget_snapshot.go | 69 +++++++++ .../sendingpolicy/budget_snapshot_test.go | 73 +++++++++ internal/sendingpolicy/gate.go | 40 +++-- internal/telemetry/metrics.go | 12 ++ internal/telemetry/prom.go | 139 ++++++++++++------ internal/telemetry/sending_budget_test.go | 49 ++++++ 13 files changed, 641 insertions(+), 63 deletions(-) create mode 100644 cmd/e2a/sending_observation.go create mode 100644 cmd/e2a/sending_observation_test.go create mode 100644 internal/sendingpolicy/budget_metrics.go create mode 100644 internal/sendingpolicy/budget_metrics_integration_test.go create mode 100644 internal/sendingpolicy/budget_snapshot.go create mode 100644 internal/sendingpolicy/budget_snapshot_test.go create mode 100644 internal/telemetry/sending_budget_test.go diff --git a/cmd/e2a/main.go b/cmd/e2a/main.go index 2a7050932..4e24ec487 100644 --- a/cmd/e2a/main.go +++ b/cmd/e2a/main.go @@ -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 diff --git a/cmd/e2a/outbound_wiring.go b/cmd/e2a/outbound_wiring.go index 7221b1b09..90a25a4f5 100644 --- a/cmd/e2a/outbound_wiring.go +++ b/cmd/e2a/outbound_wiring.go @@ -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 diff --git a/cmd/e2a/sending_observation.go b/cmd/e2a/sending_observation.go new file mode 100644 index 000000000..7a68419ff --- /dev/null +++ b/cmd/e2a/sending_observation.go @@ -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: + } + } +} diff --git a/cmd/e2a/sending_observation_test.go b/cmd/e2a/sending_observation_test.go new file mode 100644 index 000000000..2e5fb0272 --- /dev/null +++ b/cmd/e2a/sending_observation_test.go @@ -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: + } +} diff --git a/docs/observability.md b/docs/observability.md index 9ee1f55e5..72723a0d8 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -193,6 +193,47 @@ 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. +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 | diff --git a/internal/sendingpolicy/budget_metrics.go b/internal/sendingpolicy/budget_metrics.go new file mode 100644 index 000000000..8f3fcd99d --- /dev/null +++ b/internal/sendingpolicy/budget_metrics.go @@ -0,0 +1,54 @@ +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 + if capacity.daily.SharedBinding { + scope = ScopeAccountSharedDaily + } + *samples = append(*samples, budgetSample{scope, "hold"}) + } +} diff --git a/internal/sendingpolicy/budget_metrics_integration_test.go b/internal/sendingpolicy/budget_metrics_integration_test.go new file mode 100644 index 000000000..5c1d76d10 --- /dev/null +++ b/internal/sendingpolicy/budget_metrics_integration_test.go @@ -0,0 +1,137 @@ +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) + } +} diff --git a/internal/sendingpolicy/budget_snapshot.go b/internal/sendingpolicy/budget_snapshot.go new file mode 100644 index 000000000..274c32d40 --- /dev/null +++ b/internal/sendingpolicy/budget_snapshot.go @@ -0,0 +1,69 @@ +package sendingpolicy + +import ( + "context" + "fmt" + "time" + + "github.com/jackc/pgx/v5" +) + +// BudgetSnapshot is a bounded, identifier-free view of the four global pools. +// Reserved counts include confirmed units. Ratios deliberately exceed 1 in +// shadow mode. No account-level rows are read or exposed. +type BudgetSnapshot struct { + UsedRatio map[string]float64 + Generation int64 + ConfigMismatch bool + ObservedAt time.Time +} + +// BudgetSnapshot reads policy and counters in one consistent read-only snapshot. +// Config source reports generation 0 and no mismatch: it has no stored-policy +// authority. A failed sample returns no usable observation. +func (m *Module) BudgetSnapshot(ctx context.Context) (BudgetSnapshot, error) { + tx, err := m.pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}) + if err != nil { + return BudgetSnapshot{}, err + } + defer func() { + cleanup, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _ = tx.Rollback(cleanup) + }() + now := m.now().UTC() + out := BudgetSnapshot{UsedRatio: map[string]float64{}, ObservedAt: now} + policy := m.configPolicy + if m.source == PolicySourceDatabase { + snapshot, err := scanPolicy(tx.QueryRow(ctx, policySelect)) + if err != nil { + return BudgetSnapshot{}, err + } + policy = snapshot.Policy + out.Generation = snapshot.Generation + hash, err := Hash(m.configPolicy) + if err != nil { + return BudgetSnapshot{}, err + } + out.ConfigMismatch = hash != snapshot.PolicySHA256 + } + keys := []counterKey{{ScopeGlobalAll, scopeIDAllCustomers}, {ScopeGlobalProbation, scopeIDProbation}, {ScopeGlobalCritical, scopeIDCritical}, {ScopeGlobalViolation, scopeIDViolation}} + // Four indexed point reads, independent of the number of accounts or old days. + for _, key := range keys { + var used int64 + err := tx.QueryRow(ctx, `SELECT COALESCE((SELECT reserved_count FROM sending_budget_counters WHERE scope=$1 AND scope_id=$2 AND day=$3),0)`, key.Scope, key.ScopeID, now.Truncate(24*time.Hour)).Scan(&used) + if err != nil { + return BudgetSnapshot{}, err + } + limit := limitFor(key, policy, "") + ratio := float64(0) + if observeBudgetScope(policy, key.Scope) && limit > 0 { + ratio = float64(used) / float64(limit) + } + out.UsedRatio[string(key.Scope)] = ratio + } + if err := tx.Commit(ctx); err != nil { + return BudgetSnapshot{}, fmt.Errorf("sendingpolicy: commit budget snapshot: %w", err) + } + return out, nil +} diff --git a/internal/sendingpolicy/budget_snapshot_test.go b/internal/sendingpolicy/budget_snapshot_test.go new file mode 100644 index 000000000..10bc4ec5c --- /dev/null +++ b/internal/sendingpolicy/budget_snapshot_test.go @@ -0,0 +1,73 @@ +package sendingpolicy_test + +import ( + "context" + "testing" + "time" + + "github.com/tokencanopy/e2a/internal/sendingpolicy" +) + +func TestBudgetSnapshotUsesCurrentLimitsAndUTCDate(t *testing.T) { + f := newFixture(t) + p := enforcingPolicy(func(p *sendingpolicy.RuntimePolicy) { p.AllCustomerGlobalDailyRecipients = 10 }) + m := f.gate(p).(*sendingpolicy.Module) + day := time.Date(2026, 1, 2, 0, 0, 0, 0, time.UTC) + m.WithClock(func() time.Time { return day }) + _, err := f.pool.Exec(f.ctx, `INSERT INTO sending_budget_counters(scope,scope_id,day,daily_limit,reserved_count,confirmed_count) VALUES ('global_all','all-customers',$1,999,15,10),('global_all','all-customers',$1::date-1,999,500,500),('account_daily','usr_synthetic',$1,20,20,20)`, day) + if err != nil { + t.Fatal(err) + } + got, err := m.BudgetSnapshot(f.ctx) + if err != nil { + t.Fatal(err) + } + if got.UsedRatio["global_all"] != 1.5 || len(got.UsedRatio) != 4 { + t.Fatalf("snapshot: %+v", got) + } + if got.UsedRatio["global_violation"] != 0 { + t.Fatal("missing scope did not zero-fill") + } + m.WithClock(func() time.Time { return day.Add(24 * time.Hour) }) + got, err = m.BudgetSnapshot(f.ctx) + if err != nil || got.UsedRatio["global_all"] != 0 { + t.Fatalf("rollover: %+v %v", got, err) + } + cancelled, cancel := context.WithCancel(f.ctx) + cancel() + if _, err = m.BudgetSnapshot(cancelled); err == nil { + t.Fatal("canceled snapshot succeeded") + } +} + +func TestBudgetSnapshotDatabasePolicyAndCorruption(t *testing.T) { + f := newFixture(t) + config := sendingpolicy.DisabledPolicy() + m := sendingpolicy.NewGate(f.pool, f.secrets(), sendingpolicy.PolicySourceDatabase, config).(*sendingpolicy.Module) + p := config + p.AllCustomerGlobalDailyRecipients = 17 + raw, err := sendingpolicy.CanonicalBytes(p) + if err != nil { + t.Fatal(err) + } + hash, err := sendingpolicy.Hash(p) + if err != nil { + t.Fatal(err) + } + if _, err = f.pool.Exec(f.ctx, `UPDATE sending_protection_runtime_policy SET generation=7,policy=$1,policy_sha256=$2 WHERE singleton`, raw, hash); err != nil { + t.Fatal(err) + } + got, err := m.BudgetSnapshot(f.ctx) + if err != nil { + t.Fatal(err) + } + if got.Generation != 7 || !got.ConfigMismatch { + t.Fatalf("database policy snapshot: %+v", got) + } + if _, err = f.pool.Exec(f.ctx, `UPDATE sending_protection_runtime_policy SET policy_sha256=repeat('0',64) WHERE singleton`); err != nil { + t.Fatal(err) + } + if _, err = m.BudgetSnapshot(f.ctx); err == nil { + t.Fatal("corrupt policy accepted") + } +} diff --git a/internal/sendingpolicy/gate.go b/internal/sendingpolicy/gate.go index 01a57eb08..0a69c41ed 100644 --- a/internal/sendingpolicy/gate.go +++ b/internal/sendingpolicy/gate.go @@ -4,7 +4,6 @@ import ( "context" "errors" "fmt" - "log" "strings" "time" @@ -401,6 +400,7 @@ func (m *Module) Reserve(ctx context.Context, ref OperationRef) (Decision, Attem return Decision{}, AttemptRef{}, err } if deniedScope != "" { + emitBudgetSamples([]budgetSample{{deniedScope, "hold"}}) return holdDecision(holdReasonForScope(deniedScope), nextUTCMidnight(day)), out, nil } return allowDecision(), out, nil @@ -541,7 +541,8 @@ func (m *Module) ConsumeAttempt(ctx context.Context, ref AttemptRef) (Decision, return decision, nil, nil } - decision, err = m.reauthorizeBudget(ctx, tx, state) + var samples []budgetSample + decision, err = m.reauthorizeBudget(ctx, tx, state, &samples) if err != nil { return Decision{}, nil, err } @@ -549,6 +550,7 @@ func (m *Module) ConsumeAttempt(ctx context.Context, ref AttemptRef) (Decision, if err := m.commit(ctx, tx, "consume budget hold"); err != nil { return Decision{}, nil, err } + emitBudgetSamples(samples) return decision, nil, nil } @@ -559,6 +561,7 @@ func (m *Module) ConsumeAttempt(ctx context.Context, ref AttemptRef) (Decision, if err := m.commit(ctx, tx, "consume authorize"); err != nil { return Decision{}, nil, err } + emitBudgetSamples(samples) return allowDecision(), auth, nil } @@ -1070,7 +1073,7 @@ func (m *Module) releaseStoredUnits(ctx context.Context, tx pgx.Tx, st authState // the acquisition cannot deadlock against each other or against another worker // doing the mirror image. A denial leaves the attempt released rather than // half-charged: a later fire-time pass re-arms it from scratch. -func (m *Module) reauthorizeBudget(ctx context.Context, tx pgx.Tx, st authState) (Decision, error) { +func (m *Module) reauthorizeBudget(ctx context.Context, tx pgx.Tx, st authState, samples *[]budgetSample) (Decision, error) { stored := st.stored // Trusted first-party accounts and the disabled budget mode both mean "no @@ -1089,7 +1092,9 @@ func (m *Module) reauthorizeBudget(ctx context.Context, tx pgx.Tx, st authState) if exempt || st.op.Purpose == PurposeTrustedSystem { return allowDecision(), nil } - if err := m.rampAuthorize(ctx, tx, st.policy, st.ramp, st.day); err != nil { + rampErr := m.rampAuthorize(ctx, tx, st.policy, st.ramp, st.day) + appendAccountTrustSamples(samples, st, rampErr) + if err := rampErr; err != nil { // releaseStoredUnits above already gave back whatever an earlier // Reserve was holding, so every hold here leaves the ledger clean. if hold, ok := rampHoldFor(err, st.day); ok { @@ -1128,16 +1133,21 @@ func (m *Module) reauthorizeBudget(ctx context.Context, tx pgx.Tx, st authState) deniedScope := Scope("") for _, key := range currentKeys { + decision := "allow" if !plan.acquire(ledgerRef{counterKey: key, Day: st.day}, st.units) { deniedScope = key.Scope - if st.policy.BudgetMode == ModeEnforce { - break + decision = "hold" + if st.policy.BudgetMode == ModeShadow { + decision = "would_hold" + // Preserve actual demand beyond the cap during observation. + plan.overrun(ledgerRef{counterKey: key, Day: st.day}, st.units) } - // Shadow deliberately overruns. Clamping the counter at the limit - // would make the shadow window prove only that the limit exists; - // what the rollout gate needs is the real aggregate demand, which - // is only visible if the counter is allowed past the cap. - plan.overrun(ledgerRef{counterKey: key, Day: st.day}, st.units) + } + if st.units > 0 && observeBudgetScope(st.policy, key.Scope) { + *samples = append(*samples, budgetSample{key.Scope, decision}) + } + if deniedScope != "" && st.policy.BudgetMode == ModeEnforce { + break } } @@ -1165,10 +1175,6 @@ func (m *Module) reauthorizeBudget(ctx context.Context, tx pgx.Tx, st authState) return holdDecision(holdReasonForScope(deniedScope), nextUTCMidnight(st.day)), nil } - if deniedScope != "" { - log.Printf("[sending-protection] shadow denial: purpose=%s scope=%s operation=%s units=%d", - st.op.Purpose, deniedScope, st.op.OperationID, st.units) - } if err := plan.flush(ctx, tx); err != nil { return Decision{}, err } @@ -1178,7 +1184,9 @@ func (m *Module) reauthorizeBudget(ctx context.Context, tx pgx.Tx, st authState) // units this transaction just took are given straight back — holding them // would charge an account for a send its own domain was not allowed to // make. - if err := m.rampAuthorize(ctx, tx, st.policy, st.ramp, st.day); err != nil { + rampErr := m.rampAuthorize(ctx, tx, st.policy, st.ramp, st.day) + appendAccountTrustSamples(samples, st, rampErr) + if err := rampErr; err != nil { hold, ok := rampHoldFor(err, st.day) if !ok { return Decision{}, err diff --git a/internal/telemetry/metrics.go b/internal/telemetry/metrics.go index 661cc7b99..2ae9880bc 100644 --- a/internal/telemetry/metrics.go +++ b/internal/telemetry/metrics.go @@ -167,6 +167,10 @@ type Metrics interface { // with the ledger table as the label. SendingLedgerRetentionRun(outcome string) + // SendingBudgetDecision records a committed per-scope gate evaluation. + SendingBudgetDecision(scope, decision string) + SendingBudgetSnapshot(usedRatio map[string]float64, generation int64, mismatch bool, observedAt float64) + // WebhookAttempt records one webhook delivery attempt. outcome ∈ // {delivered, retryable_failure, exhausted, webhook_deleted, // skipped_disabled}. statusClass is the HTTP status class of the @@ -606,3 +610,11 @@ func (l *Log) SetThreadRelationshipPercent(string, float64) {} // Compile guard. var _ Metrics = NoOp{} var _ Metrics = (*Log)(nil) + +func (NoOp) SendingBudgetDecision(string, string) {} +func (l *Log) SendingBudgetDecision(scope, decision string) { + log.Printf("[metrics] event=sending_budget.decision scope=%s decision=%s", enum(sendingBudgetScopeSet, scope), enum(sendingBudgetDecisionSet, decision)) +} + +func (NoOp) SendingBudgetSnapshot(map[string]float64, int64, bool, float64) {} +func (*Log) SendingBudgetSnapshot(map[string]float64, int64, bool, float64) {} diff --git a/internal/telemetry/prom.go b/internal/telemetry/prom.go index 2c85b8380..3bad94392 100644 --- a/internal/telemetry/prom.go +++ b/internal/telemetry/prom.go @@ -18,49 +18,55 @@ import ( // addresses, URLs, or credentials — see docs/observability.md for the // full catalog and the cardinality contract. type Prom struct { - reg *prometheus.Registry - - httpRequests *prometheus.CounterVec - httpDuration *prometheus.HistogramVec - smtpInbound *prometheus.CounterVec - smtpDuration prometheus.Histogram - outQueueWait prometheus.Histogram - outTerminal *prometheus.CounterVec - outTerminalLat prometheus.Histogram - outAttempts *prometheus.CounterVec - outAttemptDur prometheus.Histogram - outRateDeferred prometheus.Counter - externalAccess *prometheus.CounterVec - sendingFeedback *prometheus.CounterVec - sendingLedgerRuns *prometheus.CounterVec - whAttempts *prometheus.CounterVec - whAttemptDur prometheus.Histogram - whTerminal *prometheus.CounterVec - whNotify *prometheus.CounterVec - whExpiredPending prometheus.Counter - whFanOutRescued prometheus.Counter - whDeliveryRescued prometheus.Counter - whFirstTryLat prometheus.Histogram - wsConnects prometheus.Counter - wsDisconnects *prometheus.CounterVec - wsRejected *prometheus.CounterVec - delegatedFailures *prometheus.CounterVec - delegatedRefresh *prometheus.CounterVec - oidcDiscovery *prometheus.CounterVec - oidcCallback *prometheus.CounterVec - provisioning *prometheus.CounterVec - wsDrained prometheus.Counter - wsSendFailures prometheus.Counter - wsActive prometheus.Gauge - inboundProcess *prometheus.CounterVec - inboundDuration prometheus.Histogram - queueDepth *prometheus.GaugeVec - queueOldestAge *prometheus.GaugeVec - threadResolution *prometheus.CounterVec - threadHeaderParse *prometheus.CounterVec - threadNull *prometheus.GaugeVec - threadViolations *prometheus.GaugeVec - threadRelationship *prometheus.GaugeVec + reg *prometheus.Registry + sendingBudgetUsed *prometheus.GaugeVec + sendingPolicyGeneration prometheus.Gauge + sendingPolicyMismatch prometheus.Gauge + sendingObservationLastSuccess prometheus.Gauge + + httpRequests *prometheus.CounterVec + httpDuration *prometheus.HistogramVec + smtpInbound *prometheus.CounterVec + smtpDuration prometheus.Histogram + outQueueWait prometheus.Histogram + outTerminal *prometheus.CounterVec + outTerminalLat prometheus.Histogram + outAttempts *prometheus.CounterVec + outAttemptDur prometheus.Histogram + outRateDeferred prometheus.Counter + externalAccess *prometheus.CounterVec + sendingFeedback *prometheus.CounterVec + sendingLedgerRuns *prometheus.CounterVec + sendingBudgetDecisions *prometheus.CounterVec + sendingBudgetDeferrals *prometheus.CounterVec + whAttempts *prometheus.CounterVec + whAttemptDur prometheus.Histogram + whTerminal *prometheus.CounterVec + whNotify *prometheus.CounterVec + whExpiredPending prometheus.Counter + whFanOutRescued prometheus.Counter + whDeliveryRescued prometheus.Counter + whFirstTryLat prometheus.Histogram + wsConnects prometheus.Counter + wsDisconnects *prometheus.CounterVec + wsRejected *prometheus.CounterVec + delegatedFailures *prometheus.CounterVec + delegatedRefresh *prometheus.CounterVec + oidcDiscovery *prometheus.CounterVec + oidcCallback *prometheus.CounterVec + provisioning *prometheus.CounterVec + wsDrained prometheus.Counter + wsSendFailures prometheus.Counter + wsActive prometheus.Gauge + inboundProcess *prometheus.CounterVec + inboundDuration prometheus.Histogram + queueDepth *prometheus.GaugeVec + queueOldestAge *prometheus.GaugeVec + threadResolution *prometheus.CounterVec + threadHeaderParse *prometheus.CounterVec + threadNull *prometheus.GaugeVec + threadViolations *prometheus.GaugeVec + threadRelationship *prometheus.GaugeVec // legacy outbox instruments (same events the Log backend emits) outboxPublished *prometheus.CounterVec @@ -92,9 +98,11 @@ const ( // Enum allowlists. Values outside these sets collapse to "other". var ( - methodSet = set("GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS") - classSet = set("1xx", "2xx", "3xx", "4xx", "5xx", "none") - smtpSet = set("accepted", "accepted_dedup", "tempfail", + sendingBudgetScopeSet = set("global_all", "account_daily", "account_shared_daily", "global_probation", "global_critical", "global_violation") + sendingBudgetDecisionSet = set("allow", "would_hold", "hold") + methodSet = set("GET", "HEAD", "POST", "PUT", "PATCH", "DELETE", "OPTIONS") + classSet = set("1xx", "2xx", "3xx", "4xx", "5xx", "none") + smtpSet = set("accepted", "accepted_dedup", "tempfail", "rejected_unknown_recipient", "rejected_unverified_domain", "rejected_quota", "rejected_line_too_long") outTermSet = set("sent", "failed_suppressed", "failed_provider", @@ -286,6 +294,16 @@ 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"}), + sendingBudgetUsed: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "e2a_sending_budget_used_ratio", Help: "Global reserved (including confirmed) units divided by the effective policy limit for the current UTC day; aggregate replicas with max, not sum."}, []string{"scope"}), + sendingPolicyGeneration: prometheus.NewGauge(prometheus.GaugeOpts{Name: "e2a_sending_protection_policy_generation", Help: "Sampled database policy generation, or zero for config source."}), + sendingPolicyMismatch: prometheus.NewGauge(prometheus.GaugeOpts{Name: "e2a_sending_protection_policy_config_mismatch", Help: "One when the sampled database policy differs from this process's config policy; zero for config source."}), + sendingObservationLastSuccess: prometheus.NewGauge(prometheus.GaugeOpts{Name: "e2a_sending_budget_observation_last_success_timestamp_seconds", Help: "Unix time of the last successful policy and global-budget snapshot; zero until the first sample."}), + sendingBudgetDecisions: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "e2a_sending_budget_decisions_total", Help: "Committed per-scope budget gate evaluations; not unique messages or sends.", + }, []string{"scope", "decision"}), + sendingBudgetDeferrals: prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "e2a_sending_budget_deferrals_total", Help: "Committed budget holds by deciding scope, including early reservation holds.", + }, []string{"scope"}), 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.", @@ -455,7 +473,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.sendingLedgerRuns, + p.outQueueWait, p.outTerminal, p.outTerminalLat, p.outAttempts, p.outAttemptDur, p.outRateDeferred, p.externalAccess, p.sendingFeedback, p.sendingLedgerRuns, p.sendingBudgetDecisions, p.sendingBudgetDeferrals, p.sendingBudgetUsed, p.sendingPolicyGeneration, p.sendingPolicyMismatch, p.sendingObservationLastSuccess, 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, @@ -470,6 +488,12 @@ func NewProm(build string) *Prom { for _, outcome := range []string{"complete", "partial", "failed"} { p.sendingLedgerRuns.WithLabelValues(outcome) } + for scope := range sendingBudgetScopeSet { + p.sendingBudgetDeferrals.WithLabelValues(scope) + for decision := range sendingBudgetDecisionSet { + p.sendingBudgetDecisions.WithLabelValues(scope, decision) + } + } return p } @@ -736,3 +760,24 @@ func (p *Prom) SetPublisherLag(sec float64) { p.publisherLag.Set(sec) } // Compile guard. var _ Metrics = (*Prom)(nil) + +func (p *Prom) SendingBudgetDecision(scope, decision string) { + scope, decision = enum(sendingBudgetScopeSet, scope), enum(sendingBudgetDecisionSet, decision) + p.sendingBudgetDecisions.WithLabelValues(scope, decision).Inc() + if decision == "hold" { + p.sendingBudgetDeferrals.WithLabelValues(scope).Inc() + } +} + +func (p *Prom) SendingBudgetSnapshot(usedRatio map[string]float64, generation int64, mismatch bool, observedAt float64) { + for _, scope := range []string{"global_all", "global_probation", "global_critical", "global_violation"} { + p.sendingBudgetUsed.WithLabelValues(scope).Set(usedRatio[scope]) + } + p.sendingPolicyGeneration.Set(float64(generation)) + value := float64(0) + if mismatch { + value = 1 + } + p.sendingPolicyMismatch.Set(value) + p.sendingObservationLastSuccess.Set(observedAt) +} diff --git a/internal/telemetry/sending_budget_test.go b/internal/telemetry/sending_budget_test.go new file mode 100644 index 000000000..5f7013b94 --- /dev/null +++ b/internal/telemetry/sending_budget_test.go @@ -0,0 +1,49 @@ +package telemetry + +import ( + "strings" + "testing" +) + +func TestSendingBudgetMetricsBoundLabelsAndCountDeferrals(t *testing.T) { + p := NewProm("") + for _, scope := range []string{"global_all", "account_daily", "account_shared_daily", "global_probation", "global_critical", "global_violation"} { + p.SendingBudgetDecision(scope, "allow") + p.SendingBudgetDecision(scope, "would_hold") + p.SendingBudgetDecision(scope, "hold") + } + p.SendingBudgetDecision("private-account@example.test", "private-message") + body := scrape(t, p) + for _, want := range []string{ + `e2a_sending_budget_decisions_total{decision="would_hold",scope="account_shared_daily"} 1`, + `e2a_sending_budget_deferrals_total{scope="global_violation"} 1`, + `e2a_sending_budget_decisions_total{decision="other",scope="other"} 1`, + } { + if !strings.Contains(body, want) { + t.Errorf("missing %s", want) + } + } + if strings.Contains(body, "private-") { + t.Fatal("unbounded label leaked") + } +} + +func TestSendingBudgetSnapshotExportsOnlyGlobalScopes(t *testing.T) { + p := NewProm("") + p.SendingBudgetSnapshot(map[string]float64{"global_all": 1.5, "account_daily": 99, "private@example.test": 99}, 7, true, 123) + body := scrape(t, p) + for _, want := range []string{`e2a_sending_budget_used_ratio{scope="global_all"} 1.5`, `e2a_sending_budget_used_ratio{scope="global_violation"} 0`, `e2a_sending_protection_policy_generation 7`, `e2a_sending_protection_policy_config_mismatch 1`, `e2a_sending_budget_observation_last_success_timestamp_seconds 123`} { + if !strings.Contains(body, want) { + t.Errorf("missing %s", want) + } + } + for _, line := range strings.Split(body, "\n") { + if strings.HasPrefix(line, "e2a_sending_budget_used_ratio") && (strings.Contains(line, "account_daily") || strings.Contains(line, "private")) { + t.Fatal(line) + } + } + p.SendingBudgetSnapshot(nil, 0, false, 124) + if !strings.Contains(scrape(t, p), `e2a_sending_budget_used_ratio{scope="global_all"} 0`) { + t.Fatal("stale gauge") + } +} From c55a3ffb33c516b30078b9e43a11c63544a4c7dc Mon Sep 17 00:00:00 2001 From: Josh Zhang <39790535+jiashuoz@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:45:40 +0800 Subject: [PATCH 2/2] fix(sending): attribute shared holds and sample the database ledger day --- docs/observability.md | 4 ++- internal/sendingpolicy/budget_metrics.go | 4 ++- .../budget_metrics_integration_test.go | 26 +++++++++++++++++++ internal/sendingpolicy/budget_snapshot.go | 8 +++++- .../sendingpolicy/budget_snapshot_test.go | 12 +++++++-- 5 files changed, 49 insertions(+), 5 deletions(-) diff --git a/docs/observability.md b/docs/observability.md index 72723a0d8..86c436c69 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -224,7 +224,9 @@ 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. -Missing counters report zero, and disabled legacy probation reports zero. +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 diff --git a/internal/sendingpolicy/budget_metrics.go b/internal/sendingpolicy/budget_metrics.go index 8f3fcd99d..249e8c788 100644 --- a/internal/sendingpolicy/budget_metrics.go +++ b/internal/sendingpolicy/budget_metrics.go @@ -46,7 +46,9 @@ func appendAccountTrustSamples(samples *[]budgetSample, st authState, err error) var capacity *accountCapacityError if errors.As(err, &capacity) { scope := ScopeAccountDaily - if capacity.daily.SharedBinding { + // 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"}) diff --git a/internal/sendingpolicy/budget_metrics_integration_test.go b/internal/sendingpolicy/budget_metrics_integration_test.go index 5c1d76d10..ff44a1b76 100644 --- a/internal/sendingpolicy/budget_metrics_integration_test.go +++ b/internal/sendingpolicy/budget_metrics_integration_test.go @@ -135,3 +135,29 @@ func TestBudgetObservationDoesNotPublishFailedConsume(t *testing.T) { 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) + } +} diff --git a/internal/sendingpolicy/budget_snapshot.go b/internal/sendingpolicy/budget_snapshot.go index 274c32d40..e8b02a8d0 100644 --- a/internal/sendingpolicy/budget_snapshot.go +++ b/internal/sendingpolicy/budget_snapshot.go @@ -31,7 +31,13 @@ func (m *Module) BudgetSnapshot(ctx context.Context) (BudgetSnapshot, error) { defer cancel() _ = tx.Rollback(cleanup) }() - now := m.now().UTC() + // Use the same clock as ledgerDay: a process clock across UTC midnight + // must not publish a fresh observation of the wrong daily pool. + var now time.Time + if err := tx.QueryRow(ctx, `SELECT clock_timestamp()`).Scan(&now); err != nil { + return BudgetSnapshot{}, err + } + now = now.UTC() out := BudgetSnapshot{UsedRatio: map[string]float64{}, ObservedAt: now} policy := m.configPolicy if m.source == PolicySourceDatabase { diff --git a/internal/sendingpolicy/budget_snapshot_test.go b/internal/sendingpolicy/budget_snapshot_test.go index 10bc4ec5c..df743e49b 100644 --- a/internal/sendingpolicy/budget_snapshot_test.go +++ b/internal/sendingpolicy/budget_snapshot_test.go @@ -12,8 +12,12 @@ func TestBudgetSnapshotUsesCurrentLimitsAndUTCDate(t *testing.T) { f := newFixture(t) p := enforcingPolicy(func(p *sendingpolicy.RuntimePolicy) { p.AllCustomerGlobalDailyRecipients = 10 }) m := f.gate(p).(*sendingpolicy.Module) - day := time.Date(2026, 1, 2, 0, 0, 0, 0, time.UTC) - m.WithClock(func() time.Time { return day }) + var day time.Time + if err := f.pool.QueryRow(f.ctx, `SELECT (clock_timestamp() AT TIME ZONE 'UTC')::date`).Scan(&day); err != nil { + t.Fatal(err) + } + // A process clock on another UTC day must not change which ledger is observed. + m.WithClock(func() time.Time { return day.Add(-24 * time.Hour) }) _, err := f.pool.Exec(f.ctx, `INSERT INTO sending_budget_counters(scope,scope_id,day,daily_limit,reserved_count,confirmed_count) VALUES ('global_all','all-customers',$1,999,15,10),('global_all','all-customers',$1::date-1,999,500,500),('account_daily','usr_synthetic',$1,20,20,20)`, day) if err != nil { t.Fatal(err) @@ -28,6 +32,10 @@ func TestBudgetSnapshotUsesCurrentLimitsAndUTCDate(t *testing.T) { if got.UsedRatio["global_violation"] != 0 { t.Fatal("missing scope did not zero-fill") } + // Moving only the previous-day counter into view must not resurrect it. + if _, err := f.pool.Exec(f.ctx, `DELETE FROM sending_budget_counters WHERE day=$1`, day); err != nil { + t.Fatal(err) + } m.WithClock(func() time.Time { return day.Add(24 * time.Hour) }) got, err = m.BudgetSnapshot(f.ctx) if err != nil || got.UsedRatio["global_all"] != 0 {