Skip to content
Open
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 @@ -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
Expand Down
50 changes: 50 additions & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 /
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
73 changes: 71 additions & 2 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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")
Expand Down
1 change: 1 addition & 0 deletions internal/config/delegated_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}
Expand Down
21 changes: 21 additions & 0 deletions internal/identity/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
38 changes: 34 additions & 4 deletions internal/identity/webhooks.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -670,17 +676,41 @@ 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 +
// enqueue. A systemic failure on OUR side (egress outage) makes every
// 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
Expand Down Expand Up @@ -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),
)
}

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

Expand Down
Loading
Loading