diff --git a/cmd/e2a/main.go b/cmd/e2a/main.go index 2b77d43ed..4b2ee9143 100644 --- a/cmd/e2a/main.go +++ b/cmd/e2a/main.go @@ -306,6 +306,10 @@ func main() { if err := store.EnsureSharedDomain(ctx, cfg.SharedDomain); err != nil { log.Fatalf("Failed to seed shared domain row: %v", err) } + // Config-controlled webhook health thresholds (issue #863): Validate has + // already rejected anything below 1, so this always carries either the + // operator's override or the compiled default. + store.SetWebhookHealthLimits(cfg.Webhook.WarnThreshold, cfg.Webhook.SweepMaxPerTick) // deliveryStore backs the legacy webhook_deliveries table. The // legacy per-agent push path (Deliverer/RetryWorker) is gone — push // now flows exclusively through the /v1/webhooks subscriber resource diff --git a/internal/config/config.go b/internal/config/config.go index b197cfdea..e89ec2ea8 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -28,6 +28,18 @@ const defaultSenderIdentityFixtureTTL = 24 * time.Hour // abandoned rather than merely idle. const defaultSenderIdentityReclaimMinAge = 168 * time.Hour +// defaultWebhookWarnThreshold and defaultWebhookSweepMaxPerTick mirror +// identity.WarnThreshold and identity.WarnSweepMaxPerTick/ +// DisableSweepMaxPerTick exactly, so an operator who configures no +// `webhook:` block at all keeps today's compiled behavior. Not imported +// directly: internal/identity does not (and should not) depend on +// internal/config, so this pair must be kept in sync with those constants +// by hand if either changes. +const ( + defaultWebhookWarnThreshold = 5 + defaultWebhookSweepMaxPerTick = 100 +) + // defaultSenderIdentityReclaimMaxPerSweep bounds deletions per reaper job. Set // so a systematic mistake costs a handful of identities and a loud log rather // than an account's worth, while still draining a realistic leak backlog @@ -427,6 +439,24 @@ type WebhookFanoutConfig struct { // (the default) disables the exemption. type WebhookConfig struct { InternalSinkURL string `yaml:"internal_sink_url"` + // WarnThreshold overrides identity.WarnThreshold: the number of + // attempt-level delivery failures in identity.WarnWindow that trips the + // early-warning notification (issue #863). Volume-dependent: a webhook + // receiving a handful of events a day can never accumulate the compiled + // default's failures in 24h, so it could never warn no matter how + // thoroughly broken. Defaults to identity.WarnThreshold (5) when unset. + // Must be >= 1 when set (Validate); a threshold of 0 would trip on any + // single recorded failure. Override with E2A_WEBHOOK_WARN_THRESHOLD. + WarnThreshold int `yaml:"warn_threshold"` + // SweepMaxPerTick overrides both identity.WarnSweepMaxPerTick and + // identity.DisableSweepMaxPerTick, the incident-response levers that cap + // how many webhooks one maintenance sweep may warn or disable (issue + // #863): during a real e2a-side outage every active webhook can satisfy + // both conditions at once, and turning this down stops a mass-mail + // without shipping a release. Defaults to 100 (both compiled constants) + // when unset. Must be >= 1 when set (Validate); 0 would silently disable + // the sweep. Override with E2A_WEBHOOK_SWEEP_MAX_PER_TICK. + SweepMaxPerTick int `yaml:"sweep_max_per_tick"` } // DeliveryFeedbackConfig controls outbound delivery feedback (decision 9 / @@ -718,6 +748,10 @@ func Load(path string) (*Config, error) { }, Inbound: InboundConfig{Mode: "sync"}, WebhookFanout: WebhookFanoutConfig{Mode: "legacy"}, + Webhook: WebhookConfig{ + WarnThreshold: defaultWebhookWarnThreshold, + SweepMaxPerTick: defaultWebhookSweepMaxPerTick, + }, SendingRamp: SendingRampConfig{ StartDaily: 50, TargetDaily: 2000, @@ -931,6 +965,16 @@ func Load(path string) (*Config, error) { if v := os.Getenv("E2A_WEBHOOK_INTERNAL_SINK_URL"); v != "" { cfg.Webhook.InternalSinkURL = v } + if v := os.Getenv("E2A_WEBHOOK_WARN_THRESHOLD"); v != "" { + if n, err := strconv.Atoi(v); err == nil { + cfg.Webhook.WarnThreshold = n + } + } + if v := os.Getenv("E2A_WEBHOOK_SWEEP_MAX_PER_TICK"); v != "" { + if n, err := strconv.Atoi(v); err == nil { + cfg.Webhook.SweepMaxPerTick = n + } + } if v := os.Getenv("E2A_OUTBOUND_SMTP_REQUIRE_TLS"); v != "" { if b, err := strconv.ParseBool(v); err == nil { cfg.OutboundSMTP.RequireTLS = &b @@ -1067,6 +1111,12 @@ func (c *Config) Validate() error { if c.Trash.RetentionDays < 1 { return fmt.Errorf("config: trash.retention_days must be at least 1 (got %d) — the stable API promises soft-deleted resources stay restorable", c.Trash.RetentionDays) } + if c.Webhook.WarnThreshold < 1 { + return fmt.Errorf("config: webhook.warn_threshold (or E2A_WEBHOOK_WARN_THRESHOLD) must be at least 1 (got %d): 0 would warn on any single recorded delivery failure", c.Webhook.WarnThreshold) + } + if c.Webhook.SweepMaxPerTick < 1 { + return fmt.Errorf("config: webhook.sweep_max_per_tick (or E2A_WEBHOOK_SWEEP_MAX_PER_TICK) must be at least 1 (got %d): 0 would silently disable the warn/auto-disable sweep", c.Webhook.SweepMaxPerTick) + } for _, cidr := range c.SMTP.ProxyTrustedCIDRs { p, err := netip.ParsePrefix(cidr) if err != nil { diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 0bbc28e63..dac77f160 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -299,8 +299,9 @@ func TestValidateAPIURL(t *testing.T) { } { t.Run(tc.name, func(t *testing.T) { cfg := &Config{ - HTTP: HTTPConfig{APIURL: tc.apiURL}, - Trash: TrashConfig{RetentionDays: 1}, + HTTP: HTTPConfig{APIURL: tc.apiURL}, + Trash: TrashConfig{RetentionDays: 1}, + Webhook: WebhookConfig{WarnThreshold: defaultWebhookWarnThreshold, SweepMaxPerTick: defaultWebhookSweepMaxPerTick}, } err := cfg.Validate() if tc.wantErr { @@ -606,6 +607,74 @@ func TestTrashRetentionDefaultOverrideAndValidation(t *testing.T) { } } +func TestWebhookHealthThresholdsDefaultOverrideAndValidation(t *testing.T) { + dir := t.TempDir() + write := func(name, body string) string { + t.Helper() + p := filepath.Join(dir, name) + if err := os.WriteFile(p, []byte(body), 0644); err != nil { + t.Fatal(err) + } + return p + } + + // Absent: compiled defaults (identity.WarnThreshold=5, + // WarnSweepMaxPerTick=DisableSweepMaxPerTick=100). + cfg, err := Load(write("default.yaml", "env: \"development\"\n")) + if err != nil { + t.Fatalf("Load: %v", err) + } + if cfg.Webhook.WarnThreshold != 5 { + t.Errorf("default Webhook.WarnThreshold = %d, want 5", cfg.Webhook.WarnThreshold) + } + if cfg.Webhook.SweepMaxPerTick != 100 { + t.Errorf("default Webhook.SweepMaxPerTick = %d, want 100", cfg.Webhook.SweepMaxPerTick) + } + + // YAML override. + cfg, err = Load(write("yaml.yaml", "env: \"development\"\nwebhook:\n warn_threshold: 3\n sweep_max_per_tick: 25\n")) + if err != nil { + t.Fatalf("Load: %v", err) + } + if cfg.Webhook.WarnThreshold != 3 { + t.Errorf("yaml Webhook.WarnThreshold = %d, want 3", cfg.Webhook.WarnThreshold) + } + if cfg.Webhook.SweepMaxPerTick != 25 { + t.Errorf("yaml Webhook.SweepMaxPerTick = %d, want 25", cfg.Webhook.SweepMaxPerTick) + } + + // Env override wins over yaml. + t.Setenv("E2A_WEBHOOK_WARN_THRESHOLD", "8") + t.Setenv("E2A_WEBHOOK_SWEEP_MAX_PER_TICK", "40") + cfg, err = Load(write("env.yaml", "env: \"development\"\nwebhook:\n warn_threshold: 3\n sweep_max_per_tick: 25\n")) + if err != nil { + t.Fatalf("Load: %v", err) + } + if cfg.Webhook.WarnThreshold != 8 { + t.Errorf("env Webhook.WarnThreshold = %d, want 8", cfg.Webhook.WarnThreshold) + } + if cfg.Webhook.SweepMaxPerTick != 40 { + t.Errorf("env Webhook.SweepMaxPerTick = %d, want 40", cfg.Webhook.SweepMaxPerTick) + } + t.Setenv("E2A_WEBHOOK_WARN_THRESHOLD", "") + t.Setenv("E2A_WEBHOOK_SWEEP_MAX_PER_TICK", "") + + // Below 1: refused (a threshold or per-tick cap of 0 silently breaks + // the feature rather than disabling it: see the field docs). + if _, err := Load(write("zero-warn.yaml", "env: \"development\"\nwebhook:\n warn_threshold: 0\n")); err == nil { + t.Error("Load should reject webhook.warn_threshold: 0") + } + if _, err := Load(write("neg-warn.yaml", "env: \"development\"\nwebhook:\n warn_threshold: -1\n")); err == nil { + t.Error("Load should reject a negative webhook.warn_threshold") + } + if _, err := Load(write("zero-sweep.yaml", "env: \"development\"\nwebhook:\n sweep_max_per_tick: 0\n")); err == nil { + t.Error("Load should reject webhook.sweep_max_per_tick: 0") + } + if _, err := Load(write("neg-sweep.yaml", "env: \"development\"\nwebhook:\n sweep_max_per_tick: -5\n")); err == nil { + t.Error("Load should reject a negative webhook.sweep_max_per_tick") + } +} + func TestSMTPProxyTrustedCIDRsEnvOverride(t *testing.T) { dir := t.TempDir() cfgPath := filepath.Join(dir, "config.yaml") diff --git a/internal/config/delegated_test.go b/internal/config/delegated_test.go index bb3307883..749c2ed17 100644 --- a/internal/config/delegated_test.go +++ b/internal/config/delegated_test.go @@ -33,6 +33,7 @@ func baseConfigWithDelegated(d DelegatedConfig, env string) *Config { Env: env, Signing: SigningConfig{HMACSecret: strings.Repeat("x", 64)}, Trash: TrashConfig{RetentionDays: 30}, + Webhook: WebhookConfig{WarnThreshold: defaultWebhookWarnThreshold, SweepMaxPerTick: defaultWebhookSweepMaxPerTick}, Delegated: d, } } diff --git a/internal/identity/store.go b/internal/identity/store.go index 9bdfef4d5..212ab622b 100644 --- a/internal/identity/store.go +++ b/internal/identity/store.go @@ -596,6 +596,14 @@ type Store struct { // feedbackRetention resolves the post-deletion horizon for retained // feedback provenance at purge time; nil means DefaultFeedbackRetention. feedbackRetention func(context.Context) (time.Duration, error) + // webhookWarnThreshold and webhookSweepMaxPerTick override WarnThreshold + // and WarnSweepMaxPerTick/DisableSweepMaxPerTick (webhooks.go) from + // operator config (see SetWebhookHealthLimits). Zero, the value every + // zero-value and NewStore-only Store has, means "use the compiled + // default", so the many tests that call NewStore(pool) directly keep + // today's behavior unchanged. + webhookWarnThreshold int + webhookSweepMaxPerTick int } // SetAccountStateHook installs the in-transaction account-state hook (see @@ -652,6 +660,19 @@ func (s *Store) SetThreadMetrics(metrics ThreadIdentityMetrics) { s.threadIdentityMetrics = metrics } +// SetWebhookHealthLimits overrides the compiled WarnThreshold and +// WarnSweepMaxPerTick/DisableSweepMaxPerTick defaults (webhooks.go) from +// operator config (issue #863: these are volume-dependent and e2a is +// self-hostable, so one compiled value cannot serve every deployment). +// Either argument left at 0 leaves the matching compiled default in place. +// cmd/e2a wires this from cfg.Webhook after Config.Validate has already +// rejected any configured value below 1, so the only zero this method sees +// there is "unset". +func (s *Store) SetWebhookHealthLimits(warnThreshold, sweepMaxPerTick int) { + s.webhookWarnThreshold = warnThreshold + s.webhookSweepMaxPerTick = sweepMaxPerTick +} + func (s *Store) cancelOutboundJobIDsTx(ctx context.Context, tx pgx.Tx, jobIDs []int64) error { if len(jobIDs) == 0 { return nil diff --git a/internal/identity/webhooks.go b/internal/identity/webhooks.go index 98a02c444..be710419a 100644 --- a/internal/identity/webhooks.go +++ b/internal/identity/webhooks.go @@ -605,6 +605,12 @@ const ( // TUNABLE: 5 / 24h are the design's proposed values, not yet frozen — // low enough to fire within one sweep of a real hard-failure burst, high // enough that a single transient blip mails nobody. +// +// WarnThreshold is also this package's compiled DEFAULT: an operator can +// override it per-deployment via Store.SetWebhookHealthLimits (issue #863), +// since a count that's right at one traffic volume is wrong at 100x or +// 1/100x. WarnWindow stays compile-time only; see the deferral note on +// DisableSweepMaxPerTick below. const ( WarnThreshold = 5 WarnWindow = 24 * time.Hour @@ -670,7 +676,12 @@ var E2AAttributableLastErrors = []string{ // disabled are never queued for that webhook, so they cannot be replayed). // Applied INSIDE the candidate subquery alongside the eligibility filter, so // the cap bounds rows we might actually disable and drains across sweeps. -// TUNABLE. +// TUNABLE, and, together with WarnSweepMaxPerTick, the compiled default for +// the one operator-facing incident-response lever: Store.SetWebhookHealthLimits +// (issue #863) lets an operator turn both caps down during a real e2a-side +// outage without shipping a release. AutoDisableThreshold/AutoDisableWindow +// above stay compile-time only: they interact with the GA-frozen 8-attempt / +// 29h21m retry envelope and want more thought before being exposed. const DisableSweepMaxPerTick = 100 // WarnSweepMaxPerTick bounds how many webhooks one warn pass may stamp + @@ -678,9 +689,28 @@ const DisableSweepMaxPerTick = 100 // active webhook satisfy the warn condition at once; an unbounded pass // would mass-mail the entire customer base copy blaming THEIR endpoints, // inside one lock-holding transaction. The cap drains legitimately over -// subsequent 5-minute sweeps. TUNABLE. +// subsequent 5-minute sweeps. TUNABLE; overridable, see DisableSweepMaxPerTick. const WarnSweepMaxPerTick = 100 +// warnThresholdOrDefault and sweepMaxPerTickOrDefault resolve the effective +// per-sweep limits: the operator override from SetWebhookHealthLimits when +// set, otherwise the compiled package default. A zero-value Store (every +// existing NewStore(pool) call site and test) never called the setter, so +// both fall through to the unchanged compiled constants. +func (s *Store) warnThresholdOrDefault() int { + if s.webhookWarnThreshold > 0 { + return s.webhookWarnThreshold + } + return WarnThreshold +} + +func (s *Store) sweepMaxPerTickOrDefault(compiledDefault int) int { + if s.webhookSweepMaxPerTick > 0 { + return s.webhookSweepMaxPerTick + } + return compiledDefault +} + // WebhookNotifyTx enqueues one webhook health-notification job inside the // sweep's transaction, so the state transition and its notification commit // atomically (a row cannot be disabled/warn-stamped without its job, and @@ -740,7 +770,7 @@ func (s *Store) AutoDisableFailingWebhooks(ctx context.Context, notifyTx Webhook ) AND enabled = true RETURNING id`, - AutoDisableThreshold, AutoDisableWindow, E2AAttributableLastErrors, DisableSweepMaxPerTick, + AutoDisableThreshold, AutoDisableWindow, E2AAttributableLastErrors, s.sweepMaxPerTickOrDefault(DisableSweepMaxPerTick), ) } @@ -813,7 +843,7 @@ func (s *Store) WarnFailingWebhooks(ctx context.Context, notifyTx WebhookNotifyT AND enabled = true AND warn_notified_at IS NULL RETURNING id`, - WarnThreshold, WarnWindow, E2AAttributableLastErrors, WarnSweepMaxPerTick, + s.warnThresholdOrDefault(), WarnWindow, E2AAttributableLastErrors, s.sweepMaxPerTickOrDefault(WarnSweepMaxPerTick), ) } diff --git a/internal/identity/webhooks_health_limits_test.go b/internal/identity/webhooks_health_limits_test.go new file mode 100644 index 000000000..874ba3e5c --- /dev/null +++ b/internal/identity/webhooks_health_limits_test.go @@ -0,0 +1,110 @@ +package identity_test + +import ( + "context" + "testing" + "time" + + "github.com/tokencanopy/e2a/internal/identity" + "github.com/tokencanopy/e2a/internal/testutil" +) + +// TestSetWebhookHealthLimits_OverridesWarnThreshold is the regression for +// issue #863: WarnThreshold used to be a compiled constant, so a deployment +// whose traffic could never accumulate 5 attempt-level failures in 24h could +// never warn no matter how thoroughly broken. SetWebhookHealthLimits lets an +// operator lower it per-deployment. +func TestSetWebhookHealthLimits_OverridesWarnThreshold(t *testing.T) { + pool := testutil.TestDB(t) + store := identity.NewStore(pool) + store.SetWebhookHealthLimits(2, 0) // sweepMaxPerTick left at its default + ctx := context.Background() + user, _ := store.CreateOrGetUser(ctx, "wh-limit-warn@example.com", "Owner", "google-wh-limit-warn") + wh, _ := store.CreateWebhook(ctx, user.ID, "https://example.com/limit-warn", "", []string{"email.received"}, identity.WebhookFilters{}) + + // 2 attempt-level failures: below the compiled WarnThreshold (5), at the + // configured override (2). + seedFailedDeliveries(t, pool, ctx, wh.ID, "limit_warn", 2, time.Minute) + + rec := ¬ifyRecorder{} + n, err := store.WarnFailingWebhooks(ctx, rec.enqueue) + if err != nil { + t.Fatalf("WarnFailingWebhooks: %v", err) + } + if n != 1 || len(rec.ids) != 1 || rec.ids[0] != wh.ID { + t.Fatalf("warned %d, enqueued %v, want 1 and [%s] under the threshold=2 override", n, rec.ids, wh.ID) + } +} + +// TestWarnFailingWebhooks_DefaultWarnThresholdUnaffectedWithoutSetter pins the +// backward-compatible half: a Store nobody configured (every existing +// NewStore(pool) call site, including every other test in this package) +// keeps behaving exactly as the compiled WarnThreshold constant, even though +// the override mechanism now exists. +func TestWarnFailingWebhooks_DefaultWarnThresholdUnaffectedWithoutSetter(t *testing.T) { + pool := testutil.TestDB(t) + store := identity.NewStore(pool) + ctx := context.Background() + user, _ := store.CreateOrGetUser(ctx, "wh-limit-default@example.com", "Owner", "google-wh-limit-default") + wh, _ := store.CreateWebhook(ctx, user.ID, "https://example.com/limit-default", "", []string{"email.received"}, identity.WebhookFilters{}) + + // 2 failures: below the compiled WarnThreshold (5). With no override in + // effect this must NOT warn. + seedFailedDeliveries(t, pool, ctx, wh.ID, "limit_default", 2, time.Minute) + + rec := ¬ifyRecorder{} + if n, err := store.WarnFailingWebhooks(ctx, rec.enqueue); err != nil || n != 0 { + t.Errorf("warned %d (err %v), want 0 below the unconfigured default threshold of %d", n, err, identity.WarnThreshold) + } +} + +// TestSetWebhookHealthLimits_OverridesSweepMaxPerTick is the regression for +// the incident-response half of #863: WarnSweepMaxPerTick and +// DisableSweepMaxPerTick used to be compiled constants, so an operator could +// not turn either cap down during a real e2a-side outage without shipping a +// release. A single override lowers both the warn and disable per-tick caps. +func TestSetWebhookHealthLimits_OverridesSweepMaxPerTick(t *testing.T) { + pool := testutil.TestDB(t) + store := identity.NewStore(pool) + store.SetWebhookHealthLimits(0, 2) // warnThreshold left at its default + ctx := context.Background() + user, _ := store.CreateOrGetUser(ctx, "wh-limit-sweep@example.com", "Owner", "google-wh-limit-sweep") + + const total = 5 // comfortably over the override cap of 2, well under the compiled 100 + for i := 0; i < total; i++ { + wh, err := store.CreateWebhook(ctx, user.ID, + "https://example.com/limit-sweep-"+string(rune('a'+i)), "", []string{"email.received"}, identity.WebhookFilters{}) + if err != nil { + t.Fatalf("CreateWebhook %d: %v", i, err) + } + seedFailedDeliveries(t, pool, ctx, wh.ID, "limit_sweep", identity.WarnThreshold, time.Minute) + } + + rec := ¬ifyRecorder{} + n, err := store.WarnFailingWebhooks(ctx, rec.enqueue) + if err != nil { + t.Fatalf("WarnFailingWebhooks: %v", err) + } + if n != 2 { + t.Fatalf("warned %d webhooks in one tick, want the overridden cap of 2 (total eligible = %d)", n, total) + } + + // The remaining webhooks drain across further ticks (each still capped at + // the same override) rather than being dropped: same shape as + // TestWarnFailingWebhooks_CapDrainsAcrossTicks, against the override + // instead of the compiled default. + n2, err := store.WarnFailingWebhooks(ctx, rec.enqueue) + if err != nil { + t.Fatalf("WarnFailingWebhooks (tick 2): %v", err) + } + if n2 != 2 { + t.Fatalf("second tick warned %d, want the override cap of 2 again", n2) + } + n3, err := store.WarnFailingWebhooks(ctx, rec.enqueue) + if err != nil { + t.Fatalf("WarnFailingWebhooks (tick 3): %v", err) + } + if n3 != total-4 { + t.Fatalf("third tick warned %d, want the final remaining %d", n3, total-4) + } +}