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 @@ -37,7 +37,7 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared
- **`config/`** — YAML + env var config loading (cleanenv); strict on both sides (undeclared YAML key, unbound `WH_*` variable) and probes `data_dir` writability when a selected backend keeps state there (`NeedsDataDir`); `backends.go` holds each layer's `<layer>.backend` (the in-process value by default; `mq.backend` also takes `nats`, with its `mq.nats` sub-block of file-path-only credentials; `coord.backend` takes `nats`, whose `coord.nats` block names only the lease bucket and rides `mq.nats`'s connection (both blocks, their rules and warnings are `mq_nats.go`); `cache.backend` takes `redis`, whose sub-block is `cache_redis.go`; and `dedupe.backend` takes `dynamodb`, with its `dedupe.dynamodb` sub-block) and `Warnings`, the valid combinations boot logs at `WARN`; `config.go` holds `roles` (`Has(Role)`) and `instance_id`, and `Validate` refuses a role split the backends cannot serve (any split over the embedded MQ; `api` without `ingest`, or the reverse, over a local cache; `coord.backend=nats` without `mq.backend=nats`; `mq.backend=nats` with `coord.backend=local` in a process running `ingest`; a process running only `sweeper` under `mq.backend=nats`) — boot is the validator, there is no dry run
- **`coord/`** — leases for work that must run in one process at a time (`Observer.Held` reads whether one is held without campaigning): `Coordinator.TryAcquire(ctx, name)` → a `Term` (fencing `Token`, strictly increasing per name; `Done`/`Err`, `ErrLost` on loss; `Resign`), `ErrHeld` while another holder's — or this coordinator's own — term is live; `RunElected` runs a loop only while holding its lease, resigning when the loop returns and campaigning again every `RetryPeriod`. `Local` is the in-process implementation (first taker wins, never expires; `Peer` is a second handle over the same table for tests); every implementation runs `coordtest.Conformance`. Imports only the standard library, so a distributed backend lives beside its connection: `coord.backend: nats` is `internal/mq/lease.go` (`ExternalNATS.Leases`), a key per lease in the operator's KV bucket, the KV revision as the fencing token, and expiry judged on the candidate's own clock (the same revision seen unchanged for 15s), never by a server TTL. `internal/app`'s `wireCoord` opens the one `coord.backend` selects and the sweeper runs through `RunElected` under the `sweeper` lease
- **`dedupe/`** — `Deduplicator` interface (two-phase `Reserve`/`Commit`/`Release` over `Key{Table, ID}`; every backend passes the `dedupetest` conformance suite) → `Embedded` (Pebble: every tenant's seen ids in one instance at `data_dir/pebble`, each key led by its tenant and table, pending claims in memory, committed ids stored with their expiry and deleted by an hourly background sweep along with the version-0 keys from before the table joined the key, open while any tenant's store is — the layout is the implementation's call, and the wiring hands it `data_dir` once; its `Stats` feed the system gauges) or `Dynamo` (one shared DynamoDB table, conditional `PutItem` claims; conformance-tested against dynamodb-local, selected by `dedupe.backend: dynamodb`; boot checks the table and never creates it outside dynamodb-local), wrapped by `Managed` whose open/closed state follows the hot-reloadable `dedupe.enabled` in the settings directory's `config.json`; `Stores` holds one `Managed` per tenant, built through a `Factory` (`func(tenant.ID) *Managed`, `Embedded.Tenant` or, gated on the table check (`Factory.Gated`), `Dynamo.Tenant` in production; `Managed` opens its store through a function, so every backend gets the same switch), and reconciled from the registry's `AfterAdopt` hook — open exactly when the tenant is served with its switch on, closed with its seen ids kept otherwise ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) stories 7 and 3)
- **`discovery/`** — `SchemaRegistry`, one per served tenant over a `Source` read once per refresh — the tenant's pool's connection and the database that pool was opened for, one snapshot, so a refused move keeps discovering the database the tenant's queries still use (`internal/app`'s `discoveries` builds, runs and stops them from `AfterAdopt` and `App.Close`: `RetryRefresh` until the first success, then `StartAutoRefresh` with a random first tick; `Lookup` answers `ErrNotLoaded` before the first success — the handlers' `503` with `Retry-After` — and `ErrUnknownTable` after; a failed loop attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`), that introspects ClickHouse `system.columns` (name/type/nullability plus `default_expression` and 1-based `position`) and `system.tables` (each table's `create_table_query`, kept in-process and never serialized — an external-engine table renders its wiring there unconditionally — endpoint, bucket/host, database, username, S3 access key id; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so the exposure is the topology, not the secret), records the server version, + `Validate()` for ingest payloads + `CanonicalizeTimestamps()` rewriting top-level `DateTime`/`DateTime64` column values to the canonical RFC 3339 UTC wire form pre-publish (Key Design Decision #19)
- **`discovery/`** — `SchemaRegistry`, one per served tenant over a `Source` read once per refresh — the tenant's pool's connection and the database that pool was opened for, one snapshot, so a refused move keeps discovering the database the tenant's queries still use (`internal/app`'s `discoveries` builds, runs and stops them from `AfterAdopt` and `App.Close`, and starts a tenant over on a fresh registry when a reload moves it to another address or database: `RetryRefresh` until the first success, then `StartAutoRefresh` with a random first tick; `Lookup` answers `ErrNotLoaded` before the first success — the handlers' `503` with `Retry-After` — and `ErrUnknownTable` after; a failed loop attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`), that introspects ClickHouse `system.columns` (name/type/nullability plus `default_expression` and 1-based `position`) and `system.tables` (each table's `create_table_query`, kept in-process and never serialized — an external-engine table renders its wiring there unconditionally — endpoint, bucket/host, database, username, S3 access key id; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so the exposure is the topology, not the secret), records the server version, + `Validate()` for ingest payloads + `CanonicalizeTimestamps()` rewriting top-level `DateTime`/`DateTime64` column values to the canonical RFC 3339 UTC wire form pre-publish (Key Design Decision #19)
- **`ingest/`** — Ingest worker pipeline (`worker.go`: JetStream input → per-table batch INSERT with DLQ output; over a `Sharded` queue, `claims.go`'s `ClaimShards` narrows the worker to the units this process is assigned — membership leases, capped rendezvous, halt-drain-then-release handover and stop, reset at takeover from a dead owner, each unit's share of a 10,000-row budget of unsettled rows). The pipeline is **insert-only**. The wire format `EventMessage` (`types.go`) carries `{table_name, scope, received_timestamp, format, columns, row}` and nothing else — `row` is one positional `JSONCompactEachRow` line and `columns` names its slots, the table's **insertable** columns (a `MATERIALIZED`/`ALIAS` column cannot be named in an `INSERT`); the worker batches per (tenant, table, column list), the tenant read off each message's `mq.Topic`, and inserts each batch into its tenant's own ClickHouse (`chconn.Pools.Target`); the worker accepts whatever table name the envelope carries (table existence was already checked by the HTTP ingest handler, which `404`s an unknown table before publish; the worker doesn't re-validate), then bulk-INSERTs. In the embedded-NATS deployment (the default), the server runs with `DontListen: true` (`internal/mq/embedded.go`), so the only Publishers reachable on the `ingest.>` subjects are in-process Go code — today, only the HTTP `/v1/ingest?table={table}` handler. Non-insert mutations (`DELETE`/`UPDATE`/`TRUNCATE`/…) must go through `POST /v1/ops/query` under the admin role (the same `RequireAdmin` gate as the rest of `/v1/ops/*`, so non-admin callers never reach the proxy) or through an operator-authored pipe that writes, gated only by its `allowed_roles` (#386). A request with no token (or an invalid one) resolves to the `default_role`, which in a production config is not the admin role (setting them equal is a loudly-warned dev-only setting), so it can't reach this endpoint. Plus `Sweeper` (Active Sweeper for NATS message lifecycle) + `EventMessage`/`BufferConsumerName` types (`types.go`)
- **`keyenc/`** — the one escaping composite keys are built from: `Escape` keeps `[A-Za-z0-9_-]` (exactly the tenant-id grammar, so a tenant id is its own escaped form) and writes every other byte as `%XX`, `Unescape` is `url.PathUnescape` (lenient: either hex case, and a byte left unescaped reads as itself, so a `%2D` an earlier build wrote still reads), `Join`/`AppendJoin` escape each field and put a separator between them (they panic on no fields, and on a separator the escaping could write or one outside ASCII) and `Split` reverses them. The package that builds a key takes raw names and escapes them itself, so no caller has to and no field reaches a key unescaped: NATS subjects (`Join`/`Split` after the verbatim tenant) and the cache's keys — the version index and the shared backend's Redis keys (`internal/cache`) — and the dedupe keys (`<tenant>/<table>/<id>`) use it; changing what it keeps orphans every stored key (an orphaned dedupe key lets a seen id through again), and on the shared backend, whose keys every process builds for itself, splits them between builds for the length of a rolling upgrade (a bump one build makes misses the entries the other filed, served until their TTL)
- **`mq/`** — the message-queue boundary: the **only** package that imports NATS/JetStream (Key Design Decision #20), and the only one that knows how the broker works. Everything else addresses events by `Topic{Tenant, Table, Scope}` (a validated tenant id and raw names — the tenant leads every subject, `ingest.<tenant>.<table>`, so one wildcard selects a tenant's traffic, and a topic without one is refused) and states intent through the interfaces — `Publisher` (`ErrQueueFull` is the backpressure signal, `ErrUnavailable` a broker that cannot be reached — both a `503`, with `Retry-After` `30` and `5`; `WithIdempotencyKey` makes a republish inside the queue's duplicate window a no-op), `Subscriber`, `ConsumerManager`/`Consumer`/`ConsumerConfig` (the ingest worker's durable consumer), `DeadLetterer` and `DeadLetterStats` (park a message, count what is parked), `Purger` (drop what is both acked and older than a cutoff — the sweeper), `Replayer` (SSE gap-fill) — composed into `Broker`, which adds each tenant's byte budget (`SetMaxBytes`/`MaxBytes`: the `mq.max_bytes_gb` reload, which opens a tenant's queue the first time) and `Stats` (the system gauges' source). A broker whose ingest queue is split into units one consumer at a time owns implements `Sharded` too (`IngestUnits`, `ResetOrphaned`, `Unowned`; `ConsumerConfig.Units` narrows a consumer to some of them, and its consumer is a `Releaser` and a `Halter`, capped per unit by `ConsumerConfig.MaxHeld`, with `ErrConsumerMismatch` when an operator durable no longer fits), and `Message.OnSettled` runs a hook once, at the first ack or nak attempt, confirmed or not. Every interface speaks per tenant, never per stream: the embedded implementation gives each tenant a queue of its own (a stream pair, `INGEST_<tenant>`/`DLQ_<tenant>`), and nothing outside the package may assume that layout — an external implementation may keep one shared stream. Subjects, prefixes, wildcards, stream names, sequences, and ack floors are private to the implementations, whose subject tokens are escaped by the shared `internal/keyenc`: `EmbeddedNATS` (`embedded.go`, `subject.go`, `purge.go`, `deadletter.go`), which `internal/app` constructs and hands everything else as a `mq.Broker`, and `ExternalNATS` (`external.go`, `subject_nats.go`, `nats_topology.go`: an operator-owned cluster whose streams, durables and lease bucket it never creates, changes, purges or deletes; `lease.go` holds `coord.backend: nats`'s leases in that bucket), which `internal/app` constructs from the `mq.nats` block when `mq.backend` is `nats` ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)). Every implementation passes the conformance suite in `internal/mq/mqtest` (`mqtest.Run`), which states the `Broker` contract as behavior; a new backend runs it from its own test, with `mqtest.Caps` only where its semantics legitimately differ
Expand Down
Loading
Loading