Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared
- **`pipes/`** — Named query pipes: `NamedQuery` type + `BindParams` + `Source` (read per request; `settings.Store` in production, `Static(q...)` in tests)
- **`policy/`** — Hasura-style access control, **role-first**: `TablePolicy` is `map[string]RolePermissions`, and a role's grant splits by operation into `SelectPermissions` (columns, row `filter`, aggregations, the `max_*` limits) and `InsertPermissions` (columns, `check`) — so a field only one side honors does not exist on the other. `Evaluate()` resolves ONE operation and leaves the other side **nil** (`Select *ResolvedSelect` / `Insert *ResolvedInsert`), which every accessor fails closed on — nil is "not resolved", distinct from an empty side, which is "unrestricted" (what the admin return builds). Claim templating (`{{ jwt.claim.path }}`) resolves during that call. Policies come from `Source`, a `func() *Policy` read per call (`settings.Store.Policy` in production, `Static(p)` in tests)
- **`query/`** — Structured query AST types + SQL builder with schema validation, structural policy predicate/limit emission, timestamp bucketing
- **`settings/`** — the settings directory, in either shape ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)): flat (the four files: tenant `0` alone) or nested (one folder per tenant, never mixed). `Validate` detects the shape and checks it — `ValidateDir` per directory (strict JSON, per-file rules, cross-file role references), folder names against `tenant.Parse`, a nested finding's `File` led by its folder; `Store` is a passive holder (one tenant's adopted snapshot, typed accessors read per call); `Registry` (tenant id → `Store`) owns `Open`, the serialized `Reload`/`ReloadTenant`, the `AfterAdopt` hooks, and the fsnotify `Watch` (flat only). Flat refuses an invalid directory at boot and keeps the previous snapshot on a rejected reload; nested fails closed per tenant (a rejected folder stops being served, the rest carry on, a whole-tree reload mirrors the folders, down to none, and a finding about the root itself rejects the reload whole). Plus the embedded (`go:embed`) seed `wavehouse bootstrap` writes
- **`settings/`** — the settings directory, in either shape ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)): flat (the four files: tenant `0` alone) or nested (one folder per tenant, never mixed). `Validate` detects the shape and checks it — `ValidateDir` per directory (strict JSON, per-file rules, cross-file role references), folder names against `tenant.Parse`, a nested finding's `File` led by its folder; `Store` is a passive holder (one tenant's adopted snapshot, typed accessors read per call); `Registry` (tenant id → `Store`) owns `Open`, the serialized `Reload`/`ReloadTenant`, the `AfterAdopt` hooks, and the fsnotify `Watch` (flat only). Flat refuses an invalid directory at boot and keeps the previous snapshot on a rejected reload; nested fails closed per tenant (a rejected folder stops being served, the rest carry on, a whole-tree reload mirrors the folders, down to none, and a finding about the root itself rejects the reload whole; `Open` alone refuses a nested root that would serve no tenant). Plus the embedded (`go:embed`) seed `wavehouse bootstrap` writes
- **`stream/`** — SSE fan-out: rows travel POSITIONALLY, so each connection is told its projected column list in an `event: schema` frame before its first row and again on drift — **not** guaranteed after a gap-fill across a column change, which can leave a connection reading live rows against a stale list until it reconnects ([#543](https://github.com/Wave-RF/WaveHouse/issues/543)) — (tracked per connection; replay tracks its own). The event `Hub` (registers subscribers by `(mq.Topic, role)` — one tenant's table — and evaluates each event under its own tenant's policy and schema registry; `Prune` evicts the subscribers of every tenant a reload stopped serving; `Broadcast` projects + serializes each event once per role, the #294 delivery hot path — a role carrying a row-level `filter` keeps the shared projection but delivers per subscriber, each subscriber's claims evaluated against the row, #319), `Subscriber` (per-connection outbound `Frame` queue, `Send`/`Frames`; claims fixed at construction, immutable; `Evict` asks its handler to end the stream), the `Bucket` fan-out set (`subscriberSet`, one per `(topic, role)`), the `Heartbeater` keepalive wheel, and `Metrics` (the `wavehouse_sse_*` stream instruments)
- **`tenant/`** — the tenant identifier ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)): `ID` (a validated string), `Parse` (letters, digits, `_`, `-`; ≤ 64 bytes — safe as a folder name and as an MQ subject token), `Default` (`"0"`), and `Header` (`X-Tenant-ID`). Imports nothing from the rest of the repo. `api.TenantMW` resolves the header against `settings.Registry` before auth on every `/v1` route outside `/v1/ops/*` (`400` malformed, `404` unknown, a bare `503` for a nested tenant whose folder was rejected) and puts the resolved `*settings.Store` in the request context; the ops routes that address one tenant (`GET /v1/ops/pipes[/{name}]`, `POST /v1/ops/settings/reload`, `GET /v1/ops/schema`, `POST /v1/ops/schema/refresh`, `POST /v1/ops/query`, `GET /v1/ops/dlq/stats`) take a strictly parsed `?tenant=` instead; handlers read it once (`api.StoreFromContext`) and pass it down as an argument, and nothing below a handler reads context. The stream hub and the ingest worker read each message's tenant off its `mq.Topic` and their getters take it; the sweeper hands the MQ each tenant's own gap window (`gapWindows`, a rejected tenant's included); each served tenant has a schema registry of its own (story 6)

Expand Down
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),

### Fixed

- **A settings directory holding only folders that define no tenant refuses boot again** (`internal/settings/{registry,tree}.go` (+ tests), `internal/app/wire.go` (comment; + tests), `docs/src/content/docs/{deployment,architecture}.md`, `settings-directory.mdx`, `AGENTS.md`): closes [#599](https://github.com/Wave-RF/WaveHouse/issues/599), part of [#583](https://github.com/Wave-RF/WaveHouse/issues/583). A root holding folders and none of the four files is read as nested, and a nested root whose every folder was rejected, or none of whose folders was named by a tenant id — a fresh volume whose only entry is `lost+found`, a parent directory mounted by mistake — opened as a registry serving no tenant, so the server came up and answered every tenant route `unknown tenant`. Boot now refuses such a root with the findings, as a flat invalid root is refused and as this root was before the nested shape existed. The rule is boot's alone: a whole-directory reload that finds every folder gone or broken still drops every tenant and keeps running, and `wavehouse validate` is unchanged. A root with one valid folder beside a rejected or stray one still boots and serves that tenant, and an empty or missing root still refuses boot as the four files, missing.
- **A tenant moved to another ClickHouse address or database discovers its new schema with the reload** (`internal/app/{wire,discoveries}.go` (+ tests), `internal/discovery/discovery.go` (comment), `internal/testutil/testutil.go`, `docs/src/content/docs/architecture.md`, `settings-directory.mdx`, `AGENTS.md`): closes [#638](https://github.com/Wave-RF/WaveHouse/issues/638). A reload that changed a tenant's `clickhouse.addr` or `clickhouse.database` orphaned the tenant's cache and left its schema registry as it was, so until the tenant's loop fired at `schema.refresh_interval` its queries and inserts were validated against the previous database's schema and run against the new one. The tenants the pools reconcile reports stale now have their registry dropped in the same hook, and the discovery reconcile that follows builds each a fresh one over the pool it is on, as it does for a tenant back after a rejection or removal: the first discovery runs at once in the tenant's own loop, so the reload never waits on ClickHouse, and until it succeeds the tenant's table lookups answer `503` with `Retry-After: 5`, as before any first discovery. A discovery that fails is logged with its tenant (`schema discovery retry failed`), counted in `wavehouse_schema_refresh_failures_total`, and retried with backoff from two seconds to sixty. A tenant whose address and database did not change keeps its registry and its loop, a flat directory's tenant `0` moves the same way, and a process without the api role, which discovers no schema, only repoints its pool.
- **A shard's owner keeps it while its ClickHouse is slow, one stuck shard no longer stalls the process's others, and a handover or clean stop keeps each table's rows in order** (`internal/mq/external_consumer.go` (new, from `external.go`), `internal/mq/{external,mq}.go`, `internal/ingest/{claims,worker}.go`, `.testcoverage.yml`, tests in `internal/mq/external_test.go`, `internal/ingest/{claims,worker}_test.go` and `tests/integration/shard_order_test.go` (new), `docs/src/content/docs/{deployment,ingest-pipeline,architecture}.md`, `AGENTS.md`): part of the external-NATS workstream of [#613](https://github.com/Wave-RF/WaveHouse/issues/613), fixing the shard ownership of the entries above before they ship. The pull consumer ran the worker's handler on its delivery goroutine and pulled again only when it returned, and the server renews a shard's pin only on a pull, so a process at its cap of held rows for longer than the pinned TTL (10 seconds) lost its pins: `wavehouse_ingest_shards_unowned` counted its shards, and a peer that judged it dead reset them, redelivering rows it still held — two writers per table. Each shard now has a puller that never runs the handler, fetches only what the shard's cap leaves room for, and at the cap renews the pin every 5 seconds with a one-row fetch, at most one `ack_wait` of such rows past the cap, 12 at the generated 1 minute (past that, and once halted, with a pull of `max_bytes` 1, which delivers nothing); measured, an owner blocked 13 seconds on a hung insert kept one pin throughout and took 2 rows past its share, and a unit drains as fast as before (160,000–192,000 rows a second, against 158,000–172,000). `ResetOrphaned` and `Unowned` also require that the shard delivered nothing and had no ack for its pinned TTL. The process-wide cap of 10,000 held rows let one stuck table take every slot and stall every shard the process owned; each shard now holds its own share, plus at most the 12 rows its renewals take (see the entry above), and with one shard's share full a healthy shard's rows kept flowing. A clean stop now halts every shard first (`mq.Halter`), keeping the pins, so the worker writes what it holds and what the shards still deliver, and releases them after its final flush; rows delivered during that flush were dropped before, and came back only once they had been waited out. A handover keeps the pin until the rows delivered are written, bounded by the worker's 60-second ack wait instead of 15 seconds, and a row the worker drops while stopping is NAKed before the release. Checked under continuous publishing, with the first process's ClickHouse slower than a batch window: every table's rows arrived exactly once and in order across a handover and a clean stop, and the next owner wrote its first rows about 2 seconds after the stop, against 9 seconds when the stop leaves rows behind. A shard fetches in 1-second pulls, so a shard stops fetching within a second of a halt, and an idle shard costs one pull a second. A shard durable deleted while no pull of it was waiting stalled silently until the five-minute topology check, since the next pull only got no responders: two such pulls in a row, or one renewal, now look the durable up, and one that is gone, or whose stream is, ends the worker (measured: 5 seconds after the delete, with the handler busy). A bind whose durable or stream is gone, or whose durable no longer fits (`mq.ErrConsumerMismatch`: `ack_wait` shorter than asked, or no `max_ack_pending`), ends the worker instead of retrying every tick with a warning, and a bind failure that lasts a minute logs an error. `wavehouse_ingest_rows_held_waits_total` is gone: nothing waits at the cap any more. A renewal at the cap takes a row rather than holding it back: a held-back redelivery goes to the back of the server's redelivery queue, and with max_bytes-1 renewals rows NAKed as 0, 1, 2, 3 came back as 2, 3, 0, 1 (measured); now in order, until a stuck shard is 12 rows past its cap. The server requeues the same way while another process's pull waits for the pin, so rows redelivered across a handover or a stop may come back out of order among themselves, though still ahead of newer rows. The unpin on release names no pin id, so it is preceded by a check that the pin is still this process's; the unit renews its pin until the unpin returns (for up to `ack_wait` and one 5-second renewal after its halt has handed on what it fetched, longer than a handover waits), so the pin cannot lapse between the two, and only a consumer leader change there could move it.
- **External NATS: the shipped manifests fit the shipped volume, boot refuses permissions narrower than the shards, and the tooling covers `js_domain`** (`internal/mq/{nats_manifests,nats_topology,external}.go`, `internal/mq/nats_permprobe.go` (new), `internal/mq/natstest/natstest.go`, `cmd/wavehouse/mq.go` (+ tests; `internal/mq/{nats_manifests,external_perms}_test.go` new), `deployments/nats/{jetstream.yaml,values.yaml}`, `docs/src/content/docs/{deployment,architecture}.md`, `configuration.mdx`): part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). The shipped streams reserved 225 GiB against the shipped 100 Gi volume, whose size the NATS Helm chart makes each server's `max_file_store`, so the server refused the third partition. `wavehouse mq manifests --file-store` (default `100Gi`) now sizes each partition at 15% of it, the history at 10% and the dead-letter stream at 5%, so the shipped four partitions reserve 75 GiB; a partition's size does not follow N, and the generator refuses streams that reserve more than the store (`--partition-max-bytes` resizes them). The test server takes its file store from the shipped values as the chart does, where it had a limit no test could reach. Boot now probes every shard durable as the connecting user, with a pull and an unpin the server rejects on their merits, and a `required` finding names each durable the user may not pull, where permissions generated for fewer shards passed boot and left those shards' rows unpulled; a consumer request the server refuses later sets `wavehouse_mq_topology_ok` to `0` at once. `wavehouse mq permissions --js-domain` allows and denies the `$JS.<domain>.API` subjects a client in a JetStream domain sends, beside the plain ones a server in the domain checks after mapping them; `wavehouse mq manifests` takes `--ingest-consumer`, `--history-stream` and `--publish-timeout`, the duplicate window following the timeout; the `wavehouse` user may ack only for its shard durables, under either ack subject layout, where it could ack for any consumer; and an `mq.nats.ingest_consumer` or `history_stream` outside `[a-zA-Z0-9_-]` refuses boot.
Expand Down
2 changes: 1 addition & 1 deletion docs/src/content/docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -220,7 +220,7 @@ The hot-reloadable half of configuration: a directory of four JSON files (`confi
- **finding.go** — `Finding` / `Severity`: errors make the directory invalid, warnings don't block adoption. The JSON shape is part of the ops API (`POST /v1/ops/settings/reload` returns them).
- **tree.go** — `Validate(root)` reads the directory's shape off its entries and checks either one: a root holding any of the four file names is flat and goes through `ValidateDir` untouched; otherwise a root holding a folder is nested, each folder name going through `tenant.Parse` and each folder through `ValidateDir`, with the folder leading every finding's `File` (`acme/policies.json`). The `Tree` it returns is nil when the finding is about the root itself — it cannot be listed, or a nested root holds a loose file or an entry that cannot be stat'ed.
- **store.go** — `Store` is a passive holder: one tenant's adopted document behind an atomic pointer, swapped by the registry, which stamps it with the id of the tenant it created the store for (`Tenant()`, how a handler names its tenant to a per-tenant resource). Consumers read typed accessors per call (`ClickHouse()`, `Auth()`, `DedupeFor(table)`, `DLQFor(table)`, `Keepalive()`, …) rather than holding values.
- **registry.go** — `Registry` maps a tenant id to its `Store` and owns everything that changes one. `Open` validates and adopts at boot; `Reload` re-validates the whole directory and `ReloadTenant` one tenant's folder, serialized with each other; `AfterAdopt` hooks run after every reload the registry applied, with the tenants it adopted — none when it only rejected or removed one, which a consumer holding a resource per tenant needs to hear of too; `For(id)` and `All()` see only the tenants being served, `Known()` every tenant held, rejected ones included, and `Resolve(id)` tells a rejected tenant from an unknown one. The shape is fixed at `Open`. Flat: an invalid directory refuses boot, and a rejected reload keeps the previous snapshot. Nested: fail closed per tenant — a folder with an error finding stops being served (the store keeps its document for requests already admitted, and gets the next good one) while the rest carry on; a whole-directory reload mirrors the folders, down to none (an emptied directory is not a change of shape); and a finding about the directory itself refuses boot or rejects the reload whole, leaving every tenant as it was. The tenant map is replaced whole by a reload, so a lookup is one lock-free load.
- **registry.go** — `Registry` maps a tenant id to its `Store` and owns everything that changes one. `Open` validates and adopts at boot; `Reload` re-validates the whole directory and `ReloadTenant` one tenant's folder, serialized with each other; `AfterAdopt` hooks run after every reload the registry applied, with the tenants it adopted — none when it only rejected or removed one, which a consumer holding a resource per tenant needs to hear of too; `For(id)` and `All()` see only the tenants being served, `Known()` every tenant held, rejected ones included, and `Resolve(id)` tells a rejected tenant from an unknown one. The shape is fixed at `Open`. Flat: an invalid directory refuses boot, and a rejected reload keeps the previous snapshot. Nested: fail closed per tenant — a folder with an error finding stops being served (the store keeps its document for requests already admitted, and gets the next good one) while the rest carry on; a whole-directory reload mirrors the folders, down to none (an emptied directory is not a change of shape); and a finding about the directory itself refuses boot or rejects the reload whole, leaving every tenant as it was; at `Open`, a nested directory that would serve no tenant (every folder rejected, or none named by a tenant id) refuses boot too, while a reload to the same tree drops every tenant and keeps running. The tenant map is replaced whole by a reload, so a lookup is one lock-free load.
- **watch.go** — `Registry.Watch`, which `internal/app` starts for a flat directory only: fsnotify on the *directory* (not the files, so atomic-writer replaces and Kubernetes ConfigMap symlink swaps aren't lost), debounced into one reload; reloads once as soon as the watch exists so an edit between the boot read and the watch is never missed. `SIGHUP` and the reload endpoint funnel through the same serialized `Reload`.
- **seed.go** / **seed/** — The embedded (`go:embed`) starter directory with every key at its default. The binary carries no compiled defaults except that a missing `dedupe.retention` means `"0"`: `wavehouse bootstrap [dir]` writes this seed, and the compose stack and e2e fixture ship copies of it.

Expand Down
Loading
Loading