From 651956a836beddc1cd44c8e5a670d76addce4fd9 Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Thu, 24 Sep 2026 23:28:38 -0400 Subject: [PATCH 1/2] fix(discovery): jitter the refresh retry backoff RetryRefresh slept exactly 2s * 2^n capped at 60s, so instances retrying against one recovering ClickHouse fired in lockstep. Each sleep is now drawn uniformly below the backoff (full jitter), via an injectable retryDelay so the tests pin the backoff sequence without wall-clock sleeps. Closes #141. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- AGENTS.md | 2 +- CHANGELOG.md | 1 + docs/src/content/docs/api.md | 2 +- docs/src/content/docs/architecture.md | 2 +- internal/app/wire.go | 2 +- internal/discovery/discovery.go | 16 ++++--- internal/discovery/discovery_test.go | 62 ++++++++++++++++++++------- 7 files changed, 63 insertions(+), 24 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 59f08ef36..7b241821e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -67,7 +67,7 @@ The invariant index — what must stay true. Full narrative and rationale live i 14. **TypeScript SDK** — `@wavehouse/sdk`: typed query builder, real-time SSE over `fetch`, live queries (incrementable/decomposable/poll aggregation), codegen CLI. Exactly one runtime dependency — `eventsource-parser` (SSE framing, itself dependency-free); adding a second needs the same scrutiny the first got. The canonical client (see §SDK Sync). 15. **Observability invariants** — stdout always 100% (sampling is OTLP-push-only); WARN+ERROR always export at 100% (a non-configurable floor — don't expose it); gRPC OTel exporters dial lazily so an unreachable collector never blocks startup; the OTel Prometheus exporter uses a **private** `prometheus.Registry`. The OTLP endpoint/TLS/custom-CA/mTLS/headers are delegated to the OpenTelemetry SDK's standard `OTEL_EXPORTER_OTLP_*` env vars — `InitProvider` passes **no** endpoint/header options. Known gap, intentionally not patched in WaveHouse app code: the pinned gRPC logs exporter (`otlploggrpc` v0.19/v0.20) ignores the env TLS-cert vars, so a custom/private CA and mutual TLS apply to traces/metrics but **not** the logs signal (public-CA/system-roots TLS and plaintext still work for logs) — upstream bug open-telemetry/opentelemetry-go#6661. A malformed `OTEL_EXPORTER_OTLP_HEADERS` is logged and skipped by the SDK (fail-soft), not fatal. Preserve when touching the logger/sampler/provider. Detail: architecture.md § `observability/`. 16. **Bearer-token-only CORS posture (security)** — Bearer JWT on every request, no cookies/sessions; `corsMiddleware` deliberately **never** emits `Access-Control-Allow-Credentials` (not needed, and `*` + credentials is a spec violation browsers reject). `cors.allowed_origins` (settings directory, per tenant: a tenant route is decorated from the list of the tenant it names, everything else from tenant `0`'s — `corsOrigins`) controls who can *read* responses, not cookie scope; CSRF protection is structural. Don't reintroduce cookie auth or `Allow-Credentials` without a design discussion — answers GitHub #29/#30. Code: `internal/api/router.go`. -17. **Non-fatal boot** — schema-discovery failure on boot is non-fatal: `internal/app` records an `api.BootState`, binds `:8080`, serves 503 on `/livez`/`/readyz` with the diagnostic, and retries via `SchemaRegistry.RetryRefresh` (backoff 2s → 60s), per tenant over a nested directory: `/livez` is 503 while no tenant has completed a first discovery, then sticky 200, and a tenant's outage after that is its log line and counter, never a probe failure. Until a tenant's first discovery its table lookups are a 503 with `Retry-After`, not a 404. Bounds supervisor restart loops. +17. **Non-fatal boot** — schema-discovery failure on boot is non-fatal: `internal/app` records an `api.BootState`, binds `:8080`, serves 503 on `/livez`/`/readyz` with the diagnostic, and retries via `SchemaRegistry.RetryRefresh` (jittered backoff 2s → 60s), per tenant over a nested directory: `/livez` is 503 while no tenant has completed a first discovery, then sticky 200, and a tenant's outage after that is its log line and counter, never a probe failure. Until a tenant's first discovery its table lookups are a 503 with `Retry-After`, not a 404. Bounds supervisor restart loops. 18. **Health endpoints** — liveness `/livez`, readiness `/readyz` (k8s convention; `/readyz` pings every open ClickHouse pool at once and is ready at the first answer, 503 naming each when none answers); `/healthz` is a permanent alias of `/livez`; `/health` + `/ready` are deprecated (removal v0.2.0, CHANGELOG #144). `/v1/health` is the SDK's content-free public ping (no ClickHouse check), a `/v1` route so it survives reverse-proxy probe-path filtering. Point k8s at `/livez`/`/readyz`, SDK/online-checks at `/v1/health`, never the deprecated aliases. 19. **Canonical timestamp wire form (fail-open at ingest)** — the HTTP ingest handler rewrites every top-level `DateTime`/`DateTime64` column value it can parse to RFC 3339 UTC (`discovery.CanonicalizeTimestamps`; per-column precision + zone precomputed at schema refresh) after validation + policy checks and **before** the NATS publish, so the one payload every consumer shares — SSE subscribers, the ClickHouse insert, the DLQ — carries the same spelling `/v1/query` renders: live and query reads can't drift on the instant (#372). Zone-less inputs are read in the column's declared zone, else the discovered server default — ClickHouse's own rule, so the spelling changes but never the instant. Deliberately **fail-open**: an unparseable value or unresolvable zone (no tzdata embedded — never a failed refresh, never a silent UTC reinterpretation, which would move instants) publishes verbatim; ingest must not reject a record over its timestamp spelling — fail-closed enforcement belongs to the stream row-filter (#381). Don't re-spell timestamps downstream. Preserve when touching `internal/discovery`, the ingest handler, or the SSE fan-out. Detail: architecture.md § `discovery/` + §Ingest Path; the exact spelling spec (truncation, zero-trimming, `Z`-only) lives in api.md §Timestamp canonicalization — keep it in sync with `canonicalTimestamp`. 20. **Sealed MQ boundary** — only `internal/mq` imports NATS/JetStream (`github.com/nats-io/…`), enforced by the `depguard` rule in `.golangci.yml`, so `make lint` fails on a leak in every package it builds (the `integration`-tagged files under `tests/` are outside lint's build context — keep them clean by convention, through `mq.Broker`). The boundary is semantic as well: everything else addresses events by `mq.Topic` and states intent through mq-owned interfaces (`Publisher`, `Consumer`, `DeadLetterer`, `Purger`, `Replayer`, …), and never builds a subject, names a stream, or reasons in sequences — so a subject, stream, or broker change lands in one package ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 4; story 5's tenant token landed there alone — `Topic.Tenant`, first in every subject). Don't add a raw accessor (`JetStream()`, `NatsConn()`, `GetServer()`) back, and don't hand-build `"ingest."`/`"dlq."` subjects outside `internal/mq` — widen the mq surface with an intent-level method instead. diff --git a/CHANGELOG.md b/CHANGELOG.md index 5bf580196..28520b6b1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -76,6 +76,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Fixed +- **Schema discovery's retry loop jitters its backoff** (`internal/discovery/discovery.go` (+ tests), `internal/app/wire.go`, `AGENTS.md`, `docs/src/content/docs/{architecture,api}.md`): `RetryRefresh` slept exactly `2s * 2^n` capped at 60s, so instances retrying against one recovering ClickHouse fired in lockstep, every 60s on the same second. Each sleep is now drawn uniformly from below the backoff (full jitter), spreading the retries over the whole window and halving the mean wait ([#141](https://github.com/Wave-RF/WaveHouse/issues/141)). - **An insert invalidates a table's cached results under every tenant the directory holds** (`internal/app/wire.go` (+ tests), `internal/settings/registry.go` (+ tests), `AGENTS.md`, `docs/src/content/docs/{deployment,architecture,ingest-pipeline}.md`): until [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6 gives each tenant its own ClickHouse, every tenant reads the same tables, but the ingest worker — which writes every event as tenant `0`'s until story 5 — bumped only tenant `0`'s cache namespaces after an insert, so another tenant's cached query could answer stale rows for up to its TTL (an hour at most). The cache the worker invalidates through now fans each bumped namespace out to every tenant the registry knows (the new `Registry.Known`), the named one and a rejected one included — a rejected tenant comes back into service with the entries it has, so leaving it out would let a folder repaired inside a TTL serve pre-insert rows; reads are untouched, so a tenant is still never served another's cached rows. The residual, a folder removed and restored inside a TTL, is closed since #610: a tenant back on a pool after an absence has its structured-query results orphaned at once (`Cache.InvalidateTenant`). Raised by CodeRabbit on #602. - **The `?token=` strip no longer repairs a query string that does not parse** (`internal/auth/auth.go`): `bearerToken` removed a query-string token by parsing the query, deleting `token`, and re-encoding what was left — and `url.ParseQuery` skips a pair it cannot read, so the re-encoding erased that pair. A handler that parses the query strictly in order to refuse a malformed one would then see a clean query: `GET /v1/ops/pipes?tenant=acme;x=1&token=…` would have answered `200` with the default tenant's pipes. The token is read exactly as before and a query that parses is rewritten exactly as before; a query that does not parse now loses its token pairs and nothing else, byte for byte. Pinned through `api.NewRouter` with the real authenticator, since a handler-level test never runs the middleware that rewrote the URL. diff --git a/docs/src/content/docs/api.md b/docs/src/content/docs/api.md index 75311090c..8a4b7cc8c 100644 --- a/docs/src/content/docs/api.md +++ b/docs/src/content/docs/api.md @@ -109,7 +109,7 @@ Returns `200 OK` once the gateway has discovered ClickHouse table schemas at lea Status code: `503 Service Unavailable` -The boot-degraded response lets an operator `curl /livez` to learn why the gateway isn't ready to serve traffic yet, instead of grepping a restart-loop log. The binary is bound on `:8080` and serves diagnostics, but is not yet accepting ingest/query traffic. Schema discovery retries with exponential backoff (2s → 60s); once a Refresh succeeds, `/livez` flips to `200` and stays there for the rest of the process lifetime — transient ClickHouse blips after that point are reflected in `/readyz`, not `/livez`. +The boot-degraded response lets an operator `curl /livez` to learn why the gateway isn't ready to serve traffic yet, instead of grepping a restart-loop log. The binary is bound on `:8080` and serves diagnostics, but is not yet accepting ingest/query traffic. Schema discovery retries with jittered exponential backoff (each wait a random time below a bound that doubles from 2s to 60s); once a Refresh succeeds, `/livez` flips to `200` and stays there for the rest of the process lifetime — transient ClickHouse blips after that point are reflected in `/readyz`, not `/livez`. Over a [nested settings directory](/deployment#the-nested-settings-directory) the probe reads every tenant together: `/livez` is `503` while **no** tenant has completed a first discovery — the diagnostic names the tenant whose attempt it reports (`schema discovery: tenant acme: …`), and reads `no tenant has completed a first discovery yet` before any attempt, when the directory serves no tenant, and once the tenant it named stops being served — and `200` from the first tenant's success on, for the rest of the process lifetime. A tenant whose ClickHouse is unreachable after that is a log line and the `wavehouse_schema_refresh_failures_total{tenant}` counter, never a probe failure. A tenant that has not completed its own first discovery answers `503` (`schema not loaded yet`) on its schema-aware routes until it does; one whose ClickHouse goes down after that answers query errors, as a single-tenant server does. diff --git a/docs/src/content/docs/architecture.md b/docs/src/content/docs/architecture.md index 9d1dc6363..5a910b579 100644 --- a/docs/src/content/docs/architecture.md +++ b/docs/src/content/docs/architecture.md @@ -129,7 +129,7 @@ The SSE fan-out, factored out of `api/` so the delivery hot path ([#294](https:/ ### `discovery/` — Schema Discovery & Validation -- **discovery.go** — `SchemaRegistry`, one per tenant since [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6, over a `Source` read once per refresh: the tenant's pool's connection and the database that pool was opened for, one snapshot (`App.discoverySource` over `chconn.Pools.For` in production, so a reload that repoints the tenant applies to the next refresh, and a refused move keeps discovering the database the tenant's queries and inserts still use; no pool is `ErrNoConnection`), queries `system.columns` to discover ClickHouse table schemas, keeping each column's `default_kind` so `IsInsertable` / `InsertableColumns` / `InsertableColumnNames` (memoized per table at refresh) can decide the insertable subset the ingest envelope and the SSE announcement are both built from. Each refresh also records the server version (`SELECT version()`), joins `system.tables` for each table's `create_table_query` (kept in-process as `TableSchema.DDL` and marked `json:"-"` — an external-engine table renders its wiring in that statement — endpoint, bucket/host, database, username, access key id — so it must never reach `/v1/ops/schema`; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so what is withheld here is the topology), reads each column's `default_expression` and 1-based `position` alongside its type, discovers the server's default time zone (`SELECT timezone()`) and bakes every `DateTime`/`DateTime64` column's canonicalization spec (precision + resolved zone) into the cached schema, so the per-record ingest path parses no type strings and loads no zones ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)). `Lookup` tells the two misses apart — `ErrNotLoaded` before the first successful refresh, `ErrUnknownTable` after — where `Get` answers nil for both (the stream hub's fail-closed reading). Supports periodic auto-refresh (`StartAutoRefresh`, the first tick at a random point within the interval so tenants adopted together do not refresh together, the cadence re-read after every tick), on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop used by `internal/app` so a transiently unreachable ClickHouse doesn't crash-loop the binary); a loop's failed attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`. Thread-safe via `sync.RWMutex`. +- **discovery.go** — `SchemaRegistry`, one per tenant since [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6, over a `Source` read once per refresh: the tenant's pool's connection and the database that pool was opened for, one snapshot (`App.discoverySource` over `chconn.Pools.For` in production, so a reload that repoints the tenant applies to the next refresh, and a refused move keeps discovering the database the tenant's queries and inserts still use; no pool is `ErrNoConnection`), queries `system.columns` to discover ClickHouse table schemas, keeping each column's `default_kind` so `IsInsertable` / `InsertableColumns` / `InsertableColumnNames` (memoized per table at refresh) can decide the insertable subset the ingest envelope and the SSE announcement are both built from. Each refresh also records the server version (`SELECT version()`), joins `system.tables` for each table's `create_table_query` (kept in-process as `TableSchema.DDL` and marked `json:"-"` — an external-engine table renders its wiring in that statement — endpoint, bucket/host, database, username, access key id — so it must never reach `/v1/ops/schema`; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so what is withheld here is the topology), reads each column's `default_expression` and 1-based `position` alongside its type, discovers the server's default time zone (`SELECT timezone()`) and bakes every `DateTime`/`DateTime64` column's canonicalization spec (precision + resolved zone) into the cached schema, so the per-record ingest path parses no type strings and loads no zones ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)). `Lookup` tells the two misses apart — `ErrNotLoaded` before the first successful refresh, `ErrUnknownTable` after — where `Get` answers nil for both (the stream hub's fail-closed reading). Supports periodic auto-refresh (`StartAutoRefresh`, the first tick at a random point within the interval so tenants adopted together do not refresh together, the cadence re-read after every tick), on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop, each sleep drawn uniformly below the backoff so instances retrying one ClickHouse do not retry in lockstep, used by `internal/app` so a transiently unreachable ClickHouse doesn't crash-loop the binary); a loop's failed attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`. Thread-safe via `sync.RWMutex`. - **timestamp.go** — `CanonicalizeTimestamps(schema, data)` rewrites every parseable value in a top-level `DateTime`/`DateTime64` column to the canonical RFC 3339 UTC wire form before the event is published ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)): zone-less values are interpreted in the column's declared zone, else the discovered server default — ClickHouse's own rule, so the spelling changes but never the instant. Fail-open: an unparseable value or unresolvable zone passes through verbatim for ClickHouse's own parser to judge; ingest never rejects a record over its timestamp spelling. `Column.TimeParser()` exposes the same grammar as a value→instant parser (nil for a column with no resolved timestamp spec — a non-timestamp column, or one whose declared zone couldn't be loaded), which the stream row-filter uses so filter constants and canonicalized payloads can't disagree on the instant ([#381](https://github.com/Wave-RF/WaveHouse/issues/381)). - **validation.go** — `Validate(schema, data)` checks incoming JSON against the discovered schema: unknown fields, type compatibility, missing required columns, null handling. Also exports the type classifiers `IsNumericType` / `IsStringType` and the storage-model classifier `NumericStorageOf` (all unwrapping `Nullable`/`LowCardinality`; the latter yields a numeric column's float width, `Decimal` scale, or integer exactness), which — together with `Column.TimeParser` from timestamp.go — seed the stream row-filter's `policy.ColumnSpec` comparison. - **discovery_test.go** — Unit tests for validation logic. diff --git a/internal/app/wire.go b/internal/app/wire.go index 009f9920e..bf1e643df 100644 --- a/internal/app/wire.go +++ b/internal/app/wire.go @@ -400,7 +400,7 @@ func queryTimeout(s *settings.Store) time.Duration { return s.ClickHouse().Query // again. Non-fatal either way. A flat // directory's tenant 0 is refreshed synchronously here, as before, so the // port binds with the state known; a failure marks the binary degraded and -// leaves the retry (backoff 2s → 60s) to its loop. A nested directory's +// leaves the retry (jittered backoff 2s → 60s) to its loop. A nested directory's // tenants refresh in their loops from the start, so boot never waits on a // tenant's ClickHouse, and a nested directory serving no tenant stays // degraded until a reload adopts one that loads. The process still binds its diff --git a/internal/discovery/discovery.go b/internal/discovery/discovery.go index f7bfdbd1d..90fe9b0ef 100644 --- a/internal/discovery/discovery.go +++ b/internal/discovery/discovery.go @@ -211,6 +211,9 @@ type SchemaRegistry struct { // firstTick picks how long StartAutoRefresh waits before its first // refresh, within the interval; rand.N, substituted by tests. firstTick func(interval time.Duration) time.Duration + // retryDelay picks how long RetryRefresh sleeps before its next attempt, + // within the current backoff; rand.N, substituted by tests. + retryDelay func(backoff time.Duration) time.Duration // loaded is set by the first successful Refresh and never cleared: the // line between "no schema known yet" and "this table is unknown". loaded atomic.Bool @@ -237,6 +240,7 @@ func NewSchemaRegistry(source Source, id tenant.ID, refreshInterval func(tenant. tenant: id, refreshInterval: refreshInterval, firstTick: rand.N[time.Duration], + retryDelay: rand.N[time.Duration], tables: make(map[string]*TableSchema), } } @@ -454,10 +458,12 @@ func clampBackoff(initialBackoff, maxBackoff time.Duration) (time.Duration, time // attempt with the resulting error, letting callers surface the latest // diagnostic (e.g. via /livez) while the registry is still degraded. // -// The first attempt fires immediately. After a failure the loop sleeps for -// initialBackoff, then doubles up to maxBackoff between attempts. Returns -// nil on success or ctx.Err() on cancellation. Zero/negative bounds are -// clamped via clampBackoff rather than busy-looping. +// The first attempt fires immediately. After a failure the loop sleeps a +// uniformly random time below the backoff ("full jitter"), which starts at +// initialBackoff and doubles up to maxBackoff, so instances retrying against +// one recovering ClickHouse spread over the whole window rather than firing +// in lockstep (#141). Returns nil on success or ctx.Err() on cancellation. +// Zero/negative bounds are clamped via clampBackoff rather than busy-looping. func (sr *SchemaRegistry) RetryRefresh(ctx context.Context, initialBackoff, maxBackoff time.Duration, onAttempt func(err error)) error { initialBackoff, maxBackoff = clampBackoff(initialBackoff, maxBackoff) backoff := initialBackoff @@ -481,7 +487,7 @@ func (sr *SchemaRegistry) RetryRefresh(ctx context.Context, initialBackoff, maxB select { case <-ctx.Done(): return ctx.Err() - case <-time.After(backoff): + case <-time.After(sr.retryDelay(backoff)): } backoff *= 2 if backoff > maxBackoff { diff --git a/internal/discovery/discovery_test.go b/internal/discovery/discovery_test.go index 6e7864153..1c8b8d19d 100644 --- a/internal/discovery/discovery_test.go +++ b/internal/discovery/discovery_test.go @@ -582,31 +582,63 @@ func TestRetryRefresh_DoesNotFireOnAttemptDuringCancel(t *testing.T) { assert.Equal(t, int32(1), conn.calls.Load(), "Refresh should have been called exactly once before the select caught ctx.Done()") } -// TestRetryRefresh_BackoffIsBounded verifies that maxBackoff caps the -// exponential growth. We use small bounds so the test stays fast. +// TestRetryRefresh_BackoffIsBounded verifies that the backoff each sleep is +// drawn within doubles from initialBackoff and is capped at maxBackoff. func TestRetryRefresh_BackoffIsBounded(t *testing.T) { t.Parallel() - // Five failures then success; with initial 1ms and max 4ms backoff, - // sleeps are 1, 2, 4, 4, 4 = 15ms total. The unbounded-doubling worst - // case would be 1+2+4+8+16 = 31ms. We leave generous headroom on the - // upper bound because shared CI runners can stall the scheduler enough - // to drag a 15ms sleep budget past 100ms; 250ms still catches a real - // unbounded backoff regression (which would balloon by orders of - // magnitude) without flaking on noisy hosts. errs := make([]error, 5) for i := range errs { errs[i] = errors.New("transient") } - sr, _ := newFakeRegistry(t, errs) + sr, conn := newFakeRegistry(t, errs) + var asked []time.Duration + sr.retryDelay = func(backoff time.Duration) time.Duration { + asked = append(asked, backoff) + return 0 + } - start := time.Now() err := sr.RetryRefresh(context.Background(), time.Millisecond, 4*time.Millisecond, nil) require.NoError(t, err) - elapsed := time.Since(start) - // Lower bound proves we actually slept; upper bound proves capping. - assert.GreaterOrEqual(t, elapsed, 10*time.Millisecond) - assert.Less(t, elapsed, 250*time.Millisecond) + ms := time.Millisecond + assert.Equal(t, []time.Duration{ms, 2 * ms, 4 * ms, 4 * ms, 4 * ms}, asked) + assert.Equal(t, int32(6), conn.calls.Load()) +} + +// TestRetryRefresh_SleepsTheJitteredDelay pins that the loop sleeps what +// retryDelay picks, not the backoff it was handed: a backoff of an hour with +// a 1ms delay still retries at once. +func TestRetryRefresh_SleepsTheJitteredDelay(t *testing.T) { + t.Parallel() + sr, conn := newFakeRegistry(t, []error{errors.New("once")}) + sr.retryDelay = func(time.Duration) time.Duration { return time.Millisecond } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + require.NoError(t, sr.RetryRefresh(ctx, time.Hour, time.Hour, nil)) + assert.Equal(t, int32(2), conn.calls.Load()) +} + +// TestRetryRefresh_DelayIsSpreadOverTheBackoff: the production retryDelay +// draws within [0, backoff) and spreads across it, so instances capped at +// the same maxBackoff do not retry on the same tick (#141). +func TestRetryRefresh_DelayIsSpreadOverTheBackoff(t *testing.T) { + t.Parallel() + sr, _ := newFakeRegistry(t, nil) + const backoff = time.Minute + seen := make(map[time.Duration]struct{}) + lo, hi := backoff, time.Duration(0) + for range 200 { + d := sr.retryDelay(backoff) + require.GreaterOrEqual(t, d, time.Duration(0)) + require.Less(t, d, backoff) + seen[d] = struct{}{} + lo, hi = min(lo, d), max(hi, d) + } + // Each bound fails with probability 0.75^200 for a uniform draw. + assert.Less(t, lo, backoff/4, "draws reach the bottom quarter") + assert.Greater(t, hi, backoff*3/4, "draws reach the top quarter") + assert.Greater(t, len(seen), 190, "draws are not clustered on a few values") } // TestRetryRefresh_NilOnAttemptIsSafe verifies the loop tolerates a nil From 2a55fb12b9615fe916f5a50a0728ef829a9be3fe Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Thu, 24 Sep 2026 23:50:36 -0400 Subject: [PATCH 2/2] docs(discovery): sweep the remaining unjittered backoff claims Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- CHANGELOG.md | 2 +- docs/src/content/docs/deployment.md | 2 +- internal/api/errors.go | 4 ++-- internal/discovery/discovery_test.go | 3 +-- 4 files changed, 5 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 28520b6b1..48ccb8af4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -76,7 +76,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Fixed -- **Schema discovery's retry loop jitters its backoff** (`internal/discovery/discovery.go` (+ tests), `internal/app/wire.go`, `AGENTS.md`, `docs/src/content/docs/{architecture,api}.md`): `RetryRefresh` slept exactly `2s * 2^n` capped at 60s, so instances retrying against one recovering ClickHouse fired in lockstep, every 60s on the same second. Each sleep is now drawn uniformly from below the backoff (full jitter), spreading the retries over the whole window and halving the mean wait ([#141](https://github.com/Wave-RF/WaveHouse/issues/141)). +- **Schema discovery's retry loop jitters its backoff** (`internal/discovery/discovery.go` (+ tests), `internal/app/wire.go`, `internal/api/errors.go`, `AGENTS.md`, `docs/src/content/docs/{architecture,api,deployment}.md`): `RetryRefresh` slept exactly `2s * 2^n` capped at 60s, so instances retrying against one recovering ClickHouse fired in lockstep, every 60s on the same second. Each sleep is now drawn uniformly from below the backoff (full jitter), spreading the retries over the whole window and halving the mean wait — so a failing tenant's retries, their log lines and `wavehouse_schema_refresh_failures_total` come about twice as often ([#141](https://github.com/Wave-RF/WaveHouse/issues/141)). - **An insert invalidates a table's cached results under every tenant the directory holds** (`internal/app/wire.go` (+ tests), `internal/settings/registry.go` (+ tests), `AGENTS.md`, `docs/src/content/docs/{deployment,architecture,ingest-pipeline}.md`): until [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6 gives each tenant its own ClickHouse, every tenant reads the same tables, but the ingest worker — which writes every event as tenant `0`'s until story 5 — bumped only tenant `0`'s cache namespaces after an insert, so another tenant's cached query could answer stale rows for up to its TTL (an hour at most). The cache the worker invalidates through now fans each bumped namespace out to every tenant the registry knows (the new `Registry.Known`), the named one and a rejected one included — a rejected tenant comes back into service with the entries it has, so leaving it out would let a folder repaired inside a TTL serve pre-insert rows; reads are untouched, so a tenant is still never served another's cached rows. The residual, a folder removed and restored inside a TTL, is closed since #610: a tenant back on a pool after an absence has its structured-query results orphaned at once (`Cache.InvalidateTenant`). Raised by CodeRabbit on #602. - **The `?token=` strip no longer repairs a query string that does not parse** (`internal/auth/auth.go`): `bearerToken` removed a query-string token by parsing the query, deleting `token`, and re-encoding what was left — and `url.ParseQuery` skips a pair it cannot read, so the re-encoding erased that pair. A handler that parses the query strictly in order to refuse a malformed one would then see a clean query: `GET /v1/ops/pipes?tenant=acme;x=1&token=…` would have answered `200` with the default tenant's pipes. The token is read exactly as before and a query that parses is rewritten exactly as before; a query that does not parse now loses its token pairs and nothing else, byte for byte. Pinned through `api.NewRouter` with the real authenticator, since a handler-level test never runs the middleware that rewrote the URL. diff --git a/docs/src/content/docs/deployment.md b/docs/src/content/docs/deployment.md index 020c676c6..f6216319c 100644 --- a/docs/src/content/docs/deployment.md +++ b/docs/src/content/docs/deployment.md @@ -244,7 +244,7 @@ Configure your load balancer or orchestrator to use these endpoints. ### Boot-time degraded mode -If ClickHouse is unreachable when WaveHouse starts (connection refused, missing database, DNS failure, etc.), the gateway no longer exits — it binds `:8080` and serves `/livez` 503 with the latest schema-discovery error as the diagnostic. Schema discovery retries in the background with exponential backoff (2s → 60s cap). Once a Refresh succeeds, `/livez` flips to 200 and normal serving begins automatically. +If ClickHouse is unreachable when WaveHouse starts (connection refused, missing database, DNS failure, etc.), the gateway no longer exits — it binds `:8080` and serves `/livez` 503 with the latest schema-discovery error as the diagnostic. Schema discovery retries in the background with jittered exponential backoff (each wait a random time below a bound that doubles from 2s to a 60s cap). Once a Refresh succeeds, `/livez` flips to 200 and normal serving begins automatically. This means: diff --git a/internal/api/errors.go b/internal/api/errors.go index d03ccf2a4..808bbb2c7 100644 --- a/internal/api/errors.go +++ b/internal/api/errors.go @@ -25,8 +25,8 @@ func writeJSONError(w http.ResponseWriter, status int, message string) { } // The Retry-After hints of the two 503s a tenant's ClickHouse side answers -// with: a schema not discovered yet, which discovery retries on a 2s → 60s -// backoff, and no pool — one that could not be opened, such as one the +// with: a schema not discovered yet, which discovery retries on a jittered +// 2s → 60s backoff, and no pool — one that could not be opened, such as one the // connection ceiling refused — which the next settings reload retries (the // ingest backpressure hint). const ( diff --git a/internal/discovery/discovery_test.go b/internal/discovery/discovery_test.go index 1c8b8d19d..00bc40733 100644 --- a/internal/discovery/discovery_test.go +++ b/internal/discovery/discovery_test.go @@ -483,8 +483,7 @@ func TestRetryRefresh_SucceedsOnFirstAttempt(t *testing.T) { }) require.NoError(t, err) - // Same 250ms headroom as TestRetryRefresh_BackoffIsBounded — the - // expected wall-clock budget here is ~0 (no sleep at all), but a + // The expected wall-clock budget here is ~0 (no sleep at all), but a // scheduler stall on a contended CI runner can drag a no-sleep test // past 100ms. 250ms is still orders of magnitude under any real-sleep // regression (the misbehaviour would sleep `initialBackoff` = 1h).