diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 06ab8b1..edb2e56 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -2,10 +2,10 @@ ## High-Level Overview -Meridian is a distributed time-series database written in Go, inspired by -Prometheus TSDB internals and Facebook's Gorilla paper. It is designed as a -single-binary system that handles ingestion, compression, storage, querying, -clustering, and visualization. +Meridian is a distributed time-series database in Go, drawing on Prometheus TSDB +internals and Facebook's Gorilla paper. It runs as a **single binary** (ingestion, +compression, storage, query, clustering, and visualization in one process), or +splits the same code across a **microservices cluster**. ``` ┌─────────────────────────────────────────────────────────┐ @@ -25,7 +25,7 @@ clustering, and visualization. │ Hash Ring │ Quorum Replication │ Read-Repair │ ├─────────────────────────────────────────────────────────┤ │ Retention & Downsampling │ -│ TTL Enforcer │ 5s→1m→1h Rollups │ +│ TTL Enforcer │ raw → 1m → 1h Rollups │ └─────────────────────────────────────────────────────────┘ ``` @@ -33,329 +33,294 @@ clustering, and visualization. ### Storage Engine -**Head Block** (`internal/storage/head.go`): In-memory storage for the most -recent data. Maintains an inverted index mapping label name/value pairs to -sorted series ID slices. Ingestion enforces a monotonic-per-series order: -samples older than a series' last are rejected (counted in -`meridian_out_of_order_samples_total`), an exact duplicate of the last sample is -deduplicated, and a conflicting value at the same timestamp is rejected — so each -series stays sorted and time bounds never invert (see ADR-015). A real `ts=0` -sample is a valid datum, not an "unset" sentinel. The head is periodically -flushed to a persistent block. - -**Write-Ahead Log** (`internal/storage/wal.go`): CRC32-framed WAL with 8-byte -aligned entries and automatic segment rotation at 128 MB. Every write is fsynced -before acknowledgment. Under **group commit** (default on, ADR-026) a single -committer goroutine coalesces concurrently-submitted frames behind one fsync — a -write still returns only after the fsync covering its frame, so durability and the -on-disk format are unchanged, but concurrent writers no longer serialize one fsync -at a time. Recovery is resilient to corruption: a bad length field or CRC mismatch -re-anchors at the next 8-byte frame boundary and keeps scanning, so one corrupt -frame no longer discards the rest of a segment. Replay can start past a given -segment so that data already captured in a durable block is not replayed again -(see below). - -**Persistent Blocks** (`internal/storage/block.go`): Gorilla-compressed blocks -with ULID-named directories. Each block contains a binary index file mapping -series IDs to byte offsets in the compressed chunks file, plus a recorded WAL -low-water-mark. Blocks are written crash-safely: into a temp directory, fsynced -(files and dirs), atomically renamed into place, then the parent directory is -fsynced — the rename is the single durable commit point. - -**Crash-consistent flush** (`internal/storage/tsdb.go`): A flush (1) under a -write lock captures the old head, installs a fresh head, and rotates the WAL so -in-flight writes land in a new segment — the rotation point is the block's -low-water-mark; (2) writes the old head to a durable block outside the lock; (3) -only then deletes the now-covered WAL segments (best-effort). On open, replay -skips WAL segments at or below the maximum low-water-mark across all blocks, so a -crash that leaves both a block and its source segments present recovers exactly -once — no loss, no double-count (see ADR-016). +**Head block** (`internal/storage/head.go`): in-memory store for the most recent data. +- Inverted index maps label name/value pairs to sorted series-ID slices. +- Monotonic per-series order (ADR-015): a sample older than the series' last is + rejected (counted in `meridian_out_of_order_samples_total`), an exact duplicate is + deduped, a same-timestamp conflicting value is rejected. Series stay sorted; time + bounds never invert. +- `ts=0` is a valid datum, not an "unset" sentinel. +- Periodically flushed to a persistent block. + +**Write-ahead log** (`internal/storage/wal.go`): CRC32-framed, 8-byte-aligned entries, +auto-rotated at 128 MB. Every write is fsynced before acknowledgment. +- **Group commit** (default on, ADR-026): one committer goroutine coalesces concurrent + frames behind a single fsync. A write still returns only after the fsync covering its + frame, so durability and the on-disk format are unchanged, but concurrent writers no + longer serialize. +- Corruption-resilient recovery: a bad length field or CRC re-anchors at the next + 8-byte frame boundary and keeps scanning, so one corrupt frame no longer discards the + rest of a segment. +- Replay can start past a given segment, so data already in a durable block is not + replayed twice. + +**Persistent blocks** (`internal/storage/block.go`): Gorilla-compressed, ULID-named +directories. +- Each block holds a binary index (series ID to byte offset in the compressed chunks) + plus a recorded WAL low-water-mark. +- Written crash-safely: temp dir → fsync files+dirs → atomic rename → fsync parent. The + rename is the single durable commit point. + +**Crash-consistent flush** (`internal/storage/tsdb.go`, ADR-016): +1. Under a write lock, capture the old head, install a fresh head, and rotate the WAL so + in-flight writes land in a new segment. The rotation point is the block's + low-water-mark. +2. Write the old head to a durable block outside the lock. +3. Only then delete the now-covered WAL segments (best-effort). + +On open, replay skips segments at or below the maximum low-water-mark across all blocks, +so a crash that leaves both a block and its source segments recovers exactly once: no +loss, no double-count. ### Compression -**Gorilla Encoding** (`internal/compress/gorilla.go`): Implements Facebook's -Gorilla compression for time-series data: -- Delta-of-delta encoding for timestamps -- XOR-based encoding for float64 values -- 4-byte count header for decoder bootstrapping -- Achieves 20-30x compression on regular metric data +**Gorilla encoding** (`internal/compress/gorilla.go`): Facebook's Gorilla codec. +- Delta-of-delta encoding for timestamps, XOR encoding for float64 values. +- 4-byte count header for decoder bootstrapping. +- ~20-30x compression on regular metric data. ### Query Engine -**Lexer** (`internal/query/lexer.go`): Tokenizes PromQL-subset expressions -including durations (5m, 1h), label matchers, operators, and aggregations. +**Lexer** (`internal/query/lexer.go`): tokenizes the PromQL subset (durations like +`5m`/`1h`, label matchers, operators, aggregations). -**Parser** (`internal/query/parser.go`): Recursive descent parser producing an -AST. Supports vector selectors, range selectors, function calls, aggregations -(sum, avg, min, max, count, topk, bottomk), binary expressions, and -sub-expressions. +**Parser** (`internal/query/parser.go`): recursive descent producing an AST: vector/range +selectors, function calls, aggregations (sum, avg, min, max, count, topk, bottomk), binary +expressions, sub-expressions. -**Planner** (`internal/query/planner.go`): Extracts label matchers for predicate -pushdown, adjusts time ranges for range selectors, selects a rollup resolution from the -span/step, and annotates each selector with the rollup aggregate its wrapping function -needs for function-aware coarse reads (ADR-025). +**Planner** (`internal/query/planner.go`): extracts label matchers for predicate pushdown, +adjusts time ranges for range selectors, selects a rollup resolution from the span/step, and +annotates each selector with the rollup aggregate its wrapping function needs (ADR-025). -**Executor** (`internal/query/executor.go`): Evaluates the AST against the TSDB. -Implements rate(), the `*_over_time` range-aggregation functions, histogram_quantile(), -and all aggregation functions. At a coarse resolution it reads the column that matches -the operation and serves rate() from the stored counter-increase column. +**Executor** (`internal/query/executor.go`): evaluates the AST against the TSDB (`rate()`, +the `*_over_time` range aggregations, `histogram_quantile()`, all aggregations). At a coarse +resolution it reads the column matching the operation and serves `rate()` from the stored +counter-increase column. ### Ingestion -**TCP Server** (`internal/ingestion/server.go`): JSON-over-TCP ingestion -protocol. Accepts WriteRequest messages containing batched time-series samples, -under a per-message size cap, a per-message read deadline, and a concurrent- -connection bound. A write that the bounded queue sheds is NACKed in the response -(`Shed`/`Throttled`). - -**Batch Writer** (`internal/ingestion/batch.go`): Coalesces incoming samples into -batches and drains them to the TSDB through a **bounded queue** with block-then-shed -flow control (ADR-023). Producers enqueue full batches; a single drain goroutine -writes them in FIFO order; a full queue blocks the producer up to a deadline (the -backpressure) and then sheds — dropping and counting the batch — instead of growing -without bound. Queue bounds come from ingestion config; a single drain preserves -FIFO so an in-order producer is not reordered into out-of-order rejections. - -**Backpressure primitive** (`internal/backpressure/queue.go`): the shared -cost-bounded, block-then-shed FIFO behind every ingest path. The cost is a sample -count, so queue depth is a memory bound (depth ≤ capacity); `Enqueue` blocks up to a -deadline when full, then sheds. - -**Admission shaper** (`internal/backpressure/admission.go`): an optional, off-by-default -layer consulted *before* the queue that makes shedding selective instead of uniform -(ADR-027). It sheds by **priority class** (a label or `__name__` match → a capacity -ceiling, so low priority sheds before high) and **per-series token-bucket fair share** -(a hot/high-cardinality series is throttled rather than starving well-behaved ones). -Both gates engage only under contention; per-series token buckets live in a fixed-size -shard array, so a cardinality flood cannot grow the tracking state. It holds no samples -— the queue capacity is still the hard memory bound — and an admission drop is folded -into `meridian_dropped_samples_total` while also being attributed by class, reason, and -series-hash bucket. The monolith `BatchWriter` and the service `WritePool` both apply it -per-series (a multi-series batch/request is filtered, not classified as a whole), and -order within a series is preserved. - -**Service write pool** (`internal/service/pool.go`): the ingestor and storage node -bound in-flight writes with a fixed worker pool draining a bounded queue. `Submit` -blocks while the queue is full and sheds past the deadline (`ErrShed` → HTTP 429 + -`Retry-After` / TCP NACK), so a stalled quorum write or a slow WAL fsync caps -concurrency rather than piling up unbounded goroutines — quorum semantics are -unchanged, only the submission rate is bounded. +**TCP server** (`internal/ingestion/server.go`): JSON-over-TCP. Accepts `WriteRequest` +batches under a per-message size cap, read deadline, and concurrent-connection bound. A +write the bounded queue sheds is NACKed (`Shed`/`Throttled`). + +**Batch writer** (`internal/ingestion/batch.go`, ADR-023): coalesces samples and drains +them to the TSDB through a **bounded queue** with block-then-shed flow control. +- Producers enqueue full batches; a single drain goroutine writes them FIFO. +- A full queue blocks the producer up to a deadline (backpressure), then sheds: it drops + and counts the batch instead of growing without bound. +- The single drain preserves FIFO, so an in-order producer is not reordered into + out-of-order rejections. + +**Backpressure primitive** (`internal/backpressure/queue.go`): the shared cost-bounded, +block-then-shed FIFO behind every ingest path. Cost is a sample count, so depth ≤ capacity +is a memory bound; `Enqueue` blocks up to a deadline when full, then sheds. + +**Admission shaper** (`internal/backpressure/admission.go`, ADR-027): optional, off by +default; consulted before the queue to make shedding selective instead of uniform. +- **Priority class**: a label or `__name__` match sets a capacity ceiling, so low priority + sheds before high. +- **Per-series fair share**: a token bucket throttles a hot or high-cardinality series + rather than starving well-behaved ones. +- Both gates engage only under contention. Per-series buckets live in a fixed-size shard + array, so a cardinality flood cannot grow the tracking state. +- Holds no samples (the queue capacity is still the hard memory bound). Drops fold into + `meridian_dropped_samples_total`, attributed by class, reason, and series-hash bucket. +- Applied per-series by both the monolith `BatchWriter` and the service `WritePool`; order + within a series is preserved. + +**Service write pool** (`internal/service/pool.go`): the ingestor and storage node bound +in-flight writes with a fixed worker pool draining a bounded queue. `Submit` blocks while +the queue is full and sheds past the deadline (`ErrShed`, mapped to HTTP 429 + `Retry-After` +or TCP NACK), so a stalled quorum write or slow WAL fsync caps concurrency rather than +piling up goroutines. Quorum semantics are unchanged; only the submission rate is bounded. ### HTTP & WebSocket -**HTTP Server** (`internal/server/http.go`): REST API for queries, label -browsing, and health checks, plus the dashboard SPA. Hardened at the boundary: a -traversal guard rejects any `..` path with 400 before `http.ServeMux` can clean and -redirect it, and static reads are confined to the dashboard directory; -`/api/v1/query` runs under a configurable deadline (`server.query_timeout`) with -`start ≤ end` validation, strict `start`/`end`/`step` parsing, and a panic-recovery -guard; and CORS echoes only configured origins — default localhost, never `*` -(ADR-018). `/api/v1/cluster` probes peers concurrently under a request-scoped -deadline and reports each reachable peer's real series/samples, or — in single-node -mode — just the one real node. - -**WebSocket Hub** (`internal/server/websocket.go`): Broadcasts live metrics and -system stats to connected dashboard clients. Each tick's payload is marshaled once -and the same bytes are sent to every client; a client whose send buffer stays full -across `maxClientDrops` consecutive broadcasts is force-disconnected so a stalled -reader cannot linger and leak its read/write goroutines. - -**Metrics** (`internal/server/metrics.go`): shared Prometheus exposition helpers. -The monolith and all five microservices serve `/metrics`; storage nodes expose the -full storage metric set, and the cumulative `meridian_samples_ingested_total` counter -is kept distinct from the windowed ingestion rate reported on `/api/v1/stats` -(ADR-017, ADR-019). Every ingest-bounding node (monolith, ingestor, storage) also -exports the write-path flow-control families via `WriteQueueMetrics`: -`meridian_dropped_samples_total`, `meridian_ingest_shed_events_total`, and -`meridian_ingest_backpressure_events_total` (cumulative counters), plus -`meridian_ingest_queue_depth`/`_capacity`/`_high_watermark` (gauges) — ADR-023. When -per-series/priority admission is enabled, `WriteAdmissionMetrics` adds the -`meridian_admission_admitted_samples_total{class}`, -`meridian_admission_dropped_samples_total{class,reason}`, and -`meridian_admission_series_bucket_dropped_total{bucket}` families (ADR-027); they are -omitted entirely when it is off, so the default scrape is unchanged. The -monolith and the gateway, where the anomaly detector runs, additionally export -`meridian_anomalies_total` (counter) and `meridian_active_anomalies` (gauge) via -`WriteAnomalyMetrics` (ADR-024). - -### Streaming anomaly detection - -**Detector** (`internal/anomaly/detector.go`): a pure, per-series online detector -fed from the same per-series stream the broadcaster emits each tick, so it runs -uniformly in the monolith (`cmd/meridian/serve.go`, from `Head().SeriesInfos()`) and -the cluster gateway (`cmd/gateway`, from the aggregated `FetchSeries`). Per-series -scoring sits behind a small **model interface**, selected by `Config.Mode`; the -detector owns the shared machinery (dedup, debounce/hysteresis, scale floor, event -emission, eviction) and a model only supplies a `(baseline, score)`. - -- **EWMA** (default): each series keeps O(1) state — an **EWMA level** (baseline) and - **EWMA variance** (dispersion); the score `|value − level| / dispersion` is a *local* - z-score, so the slow diurnal swing and memory drift are tracked (small residual → no - alert) while a spike departs sharply (large score → alert) — where a naive global - z-score would flag the diurnal peak (ADR-024). A Welford warmup seeds the baseline. -- **Holt-Winters** (`holtwinters.go`, opt-in via `mode: holt_winters`): an additive - level+trend+seasonal model (ADR-028). It derives the seasonal **phase from the sample - timestamp** (`ts mod season_period`), warms up over one full season, and scores each - value against the band for its own time of day — so it flags a value normal globally - but abnormal for that phase, which EWMA cannot. Per-series state is O(season) (the - seasonal array). It reuses the shared scale floor and the Huber clamp. - -A Huber clamp bounds a spike's pull on the baseline/dispersion; a relative scale floor -stops a flat series from collapsing the dispersion; debounce + a hysteresis clear band -and timestamp dedup (on `SeriesInfo.LastTS`) avoid alert storms; stale series are -evicted so memory follows live cardinality. Raise/clear transitions broadcast as a -distinct `anomaly` WebSocket frame and into a bounded recent-events ring exposed at -`/api/v1/anomalies` (which also reports the active model) for late-joining clients; the -dashboard's Anomalies strip lists them most-recent-first and shows the model. - -### Cluster (microservices tier) - -**Hash Ring** (`internal/cluster/ring.go`): a 64-bit SHA-256 consistent hash ring -with configurable virtual nodes and a deterministic nodeID tie-break, so a key's -replica set is a function of membership rather than of node-join order. `GetNodes` -returns the first N distinct **routable** replicas walking the ring, skipping nodes in -the Dead, Leaving, or Joining (catching-up) state; `PreferenceList` returns the natural -N owners *including* the skipped ones, so hinted handoff can tell which owner a write -missed. `SetState`/`State`/`LiveNodes` are the hooks the health monitor and replay loop -drive. - -**Replicated storage client** (`internal/service/client.go`): the live routing -source. It seeds a ring from the configured storage addresses and applies the -quorum model (ADR-022): - -- **Write** — for each series, `replicas = ring.GetNodes(key, N)`; the series is sent - to all live replicas (batched per node) and the write succeeds only at ≥W acks. - Too few live replicas or acks is a quorum error, not a partial write. With hinted - handoff on, a natural owner the write missed (down, catching up, or its write failed) - has the series buffered as a durable hint (below). -- **Read** — a label-matcher query scatters to all live nodes (a superset of any - matched series' replicas), merges and dedupes by (series, timestamp), enforces read - quorum R globally and per series, and asynchronously **read-repairs** any responding - replica missing points relative to the merged truth. -- **Health** — a background `/health` monitor (`StartHealthMonitor`) sets ring node - state, so routing excludes dead nodes and re-includes recovered ones. The same refresh - discovers each live node's rollup tier availability, which the resolution planner - intersects across live nodes. With hinted handoff on, a node returning with a hint - backlog is routed through `joining` (catching up) before `active`. - -**Hinted handoff** (`internal/service/handoff.go`, `hints.go`; ADR-029): the ingestor +**HTTP server** (`internal/server/http.go`, ADR-018): REST API for queries, label browsing, +and health checks, plus the dashboard SPA. Hardened at the boundary: +- A traversal guard rejects any `..` path with 400 before `http.ServeMux` can clean and + redirect it; static reads are confined to the dashboard directory. +- `/api/v1/query` runs under a configurable deadline (`server.query_timeout`) with + `start ≤ end` validation, strict `start`/`end`/`step` parsing, and panic recovery. +- CORS echoes only configured origins (default localhost, never `*`). +- `/api/v1/cluster` probes peers concurrently under a request-scoped deadline, reporting + each reachable peer's real series/samples (or, single-node, just the one node). + +**WebSocket hub** (`internal/server/websocket.go`): broadcasts live metrics and stats to +dashboard clients. Each tick's payload is marshaled once and the same bytes sent to every +client. A client whose send buffer stays full across `maxClientDrops` consecutive +broadcasts is force-disconnected, so a stalled reader cannot leak its goroutines. + +**Metrics** (`internal/server/metrics.go`): shared Prometheus exposition helpers; the +monolith and all five microservices serve `/metrics`. +- Storage nodes expose the full storage set; the cumulative `meridian_samples_ingested_total` + counter stays distinct from the windowed rate on `/api/v1/stats` (ADR-017, ADR-019). +- Ingest-bounding nodes export the flow-control families (`WriteQueueMetrics`, ADR-023): + dropped/shed/backpressure counters plus queue depth/capacity/high-watermark gauges. +- With admission on, `WriteAdmissionMetrics` adds admitted/dropped/series-bucket breakdowns + by class and reason (ADR-027); omitted when off. +- The monolith and gateway (where the detector runs) export `meridian_anomalies_total` and + `meridian_active_anomalies` (`WriteAnomalyMetrics`, ADR-024). + +### Streaming Anomaly Detection + +**Detector** (`internal/anomaly/detector.go`, ADR-024): a pure, per-series online detector +fed from the per-series tick stream the broadcaster emits, so it runs uniformly in the +monolith (`Head().SeriesInfos()`) and the cluster gateway (aggregated `FetchSeries`). +Scoring sits behind a small **model interface** selected by `Config.Mode`; the detector owns +the shared machinery (dedup, debounce/hysteresis, scale floor, event emission, eviction) and +a model supplies only a `(baseline, score)`. + +- **EWMA** (default): O(1) state per series, an EWMA **level** (baseline) and EWMA + **variance** (dispersion). The score `|value - level| / dispersion` is a *local* z-score, + so slow diurnal swing and drift are tracked (no alert) while a spike departs sharply + (alert), where a global z-score would flag the diurnal peak. A Welford warmup seeds the + baseline. +- **Holt-Winters** (`holtwinters.go`, opt-in `mode: holt_winters`, ADR-028): additive + level+trend+seasonal. It derives the seasonal phase from the sample timestamp + (`ts mod season_period`), warms up over one full season, and scores each value against the + band for its own time of day, so it flags a value normal globally but abnormal for its + phase, which EWMA cannot. State is O(season). + +Shared machinery: a Huber clamp bounds a spike's pull on the baseline/dispersion; a relative +scale floor stops a flat series collapsing the dispersion; debounce, a hysteresis clear band, +and timestamp dedup (on `SeriesInfo.LastTS`) avoid alert storms; stale series are evicted so +memory follows live cardinality. Raise/clear transitions broadcast as a distinct `anomaly` +WebSocket frame and into a bounded recent-events ring at `/api/v1/anomalies` (which reports +the active model); the dashboard's Anomalies strip lists them most-recent-first. + +### Cluster (Microservices Tier) + +**Hash ring** (`internal/cluster/ring.go`): a 64-bit SHA-256 consistent hash ring with +configurable virtual nodes and a deterministic nodeID tie-break, so a key's replica set is a +function of membership rather than join order. +- `GetNodes` returns the first N distinct **routable** replicas, skipping Dead/Leaving/ + Joining nodes. +- `PreferenceList` returns the natural N owners *including* the skipped ones, so hinted + handoff can tell which owner a write missed. +- `SetState`/`State`/`LiveNodes` are the hooks the health monitor and replay loop drive. + +**Replicated storage client** (`internal/service/client.go`, ADR-022): the live routing +source; seeds a ring from the configured storage addresses and applies the quorum model. +- **Write**: `replicas = ring.GetNodes(key, N)`, sent to all live replicas (batched per + node), succeeding only at ≥W acks. Too few is a quorum error, not a partial write. With + handoff on, a missed natural owner has the series buffered as a durable hint. +- **Read**: scatters to all live nodes, merges and dedupes by (series, timestamp), enforces + read quorum R globally and per series, and asynchronously **read-repairs** any replica + missing points relative to the merged truth. +- **Health**: `StartHealthMonitor` sets ring node state so routing excludes dead nodes and + re-includes recovered ones. The same refresh discovers each node's rollup-tier + availability (intersected by the resolution planner); a node returning with a hint backlog + routes through `joining` before `active`. + +**Hinted handoff** (`internal/service/handoff.go`, `hints.go`, ADR-029): the ingestor buffers writes a replica misses while down and replays them on its return, closing the interior-gap case read-repair cannot reach. - -- **Buffer** — `Write` diffs `PreferenceList` (natural owners, incl. down) against the - live replicas it reached and, once every series met quorum, records a durable hint per - missed owner in a bounded, per-target `HintStore` (one crash-safe file per hint, FIFO, - capped in samples per target with drop-oldest, rebuilt from disk on restart). -- **Replay** — a background loop drains each reachable target's hints in order through - the storage node's out-of-order-tolerant `/api/internal/backfill` endpoint - (`TSDB.Backfill`), deleting each on ack. Backfill inserts historical samples in sorted - position (gap-fill only) and logs them under a distinct WAL frame, so they survive a - storage crash while the live in-order policy (ADR-015) is untouched. -- **Catch-up lifecycle** — a replica returning from Dead with a backlog enters `joining` - (excluded from live read/write routing) and is promoted to `active` only once its hints - drain, so it is made whole before it can take a fresh write that would strand an - interior gap. An already-Active node is never demoted by a transient hint. +- **Buffer**: `Write` diffs `PreferenceList` against the live replicas it reached and, once + every series met quorum, records a durable hint per missed owner in a bounded, per-target + `HintStore` (one crash-safe file per hint, FIFO, capped per target with drop-oldest, + rebuilt from disk on restart). +- **Replay**: a background loop drains each reachable target's hints in order through the + out-of-order-tolerant `/api/internal/backfill` (`TSDB.Backfill`), deleting each on ack. + Backfill inserts historical samples in sorted position (gap-fill only) under a distinct WAL + frame, so they survive a crash while the live in-order policy (ADR-015) is untouched. +- **Catch-up**: a replica returning from Dead with a backlog enters `joining` (out of + routing) and is promoted to `active` only once its hints drain, so it is made whole before + it can strand an interior gap. An already-Active node is never demoted by a transient hint. **Proactive anti-entropy** (`internal/service/antientropy.go`, `internal/storage/merkle.go`, -`internal/cluster/antientropy.go`; ADR-030): a rate-limited, jittered background sweep on -the ingestor converges co-replicas neither read-repair nor hinted handoff reaches — cold -data, an unobserved partial write, a hint dropped past the cap, a series no longer written. - -- **Digest** — each storage node summarises the series whose ring position falls in a set - of hash arcs, bucketed by time window, as a Merkle tree - (`/api/internal/antientropy/digest`). Equal roots mean the replicas agree over the whole - span and nothing transfers; a differing root localises divergence to specific windows. - The ring hash is injected (`cluster.HashKey`), so the storage node stays ring-agnostic. -- **Sweep** — `ring.ReplicaGroups` partitions the ring into the arc sets that share a - replica set; the loop round-robins them (`groups_per_round` per round), fetches each live - replica's digest, and for a divergent window reads it from every replica - (`/api/internal/antientropy/range`) and pushes each the points it lacks back through the - same `/api/internal/backfill` apply hinted handoff uses — bidirectional gap-fill to the +`internal/cluster/antientropy.go`, ADR-030): a rate-limited, jittered background sweep on the +ingestor that converges co-replicas neither read-repair nor hinted handoff reaches (cold +data, an unobserved partial write, a hint dropped past the cap, a series no longer written). +- **Digest**: each storage node summarises the series whose ring position falls in a set of + hash arcs, bucketed by time window, as a Merkle tree (`/api/internal/antientropy/digest`). + Equal roots mean agreement (nothing transfers); a differing root localises divergence to + specific windows. The ring hash is injected (`cluster.HashKey`), so storage stays + ring-agnostic. +- **Sweep**: `ring.ReplicaGroups` partitions the ring into the arc sets sharing a replica + set; the loop round-robins them (`groups_per_round` per round). For a divergent window it + reads that window from every replica (`/api/internal/antientropy/range`) and pushes each + the points it lacks through `/api/internal/backfill`, a bidirectional gap-fill to the window union. -- **Scoped to shared arcs** — comparison runs per replica group, so a node's data outside - an arc it actually shares with the peer is never mistaken for divergence and re-shipped. +- **Scoped to shared arcs**: comparison runs per replica group, so a node's data outside an + arc it shares with the peer is never mistaken for divergence and re-shipped. **Rebalancing on membership change** (`internal/cluster/rebalance.go`, -`internal/service/rebalance.go`, `internal/storage/rebalance.go`; ADR-031): adding or -removing a storage node re-derives placement *and moves the data*. Where read-repair, hinted -handoff, and anti-entropy all converge nodes that *share* an owner set, rebalancing **moves** -the owner set — the piece those layers assume but never provide. - -- **Diff** — `cluster.PlacementDelta` diffs two ring snapshots (before/after the change) at - the union of their virtual-node boundaries and groups the arcs whose owner set changed into - `OwnershipChange{Before, After, Added, Removed, Arcs}`. A joining node is already a target - owner; a leaving/dead node is excluded as a target but still a source. -- **Migrate** — for each change the coordinator reads the moved arcs from a current owner - (`/api/internal/antientropy/range`) and pushes them to each new owner through the same - `/api/internal/backfill` apply handoff and anti-entropy use; a new owner is confirmed only - on its ack. -- **GC** — once data has moved and the ring is flipped to the target placement (so read-repair - cannot re-add it), the old owner drops the arcs it no longer owns - (`TSDB.DropSeriesInRanges` via `/api/internal/rebalance/drop`): head flushed, then each - block left, deleted whole, or rewritten to keep only its owned series. GC runs only after - the new owners are confirmed at quorum and never drops the last copy. -- **Lifecycle** — a join catches up out of routing (`joining`) then promotes to `active` then - GCs the displaced owner; a leave migrates while still `active` (reads stay complete), then - goes `leaving` → removed. A node returning from `dead` has the over-replication its absence +`internal/service/rebalance.go`, `internal/storage/rebalance.go`, ADR-031): adding or +removing a storage node re-derives placement *and moves the data*. Read-repair, hinted +handoff, and anti-entropy all converge nodes that *share* an owner set; rebalancing supplies +the piece they assume but never provide, moving the owner set itself. +- **Diff**: `cluster.PlacementDelta` diffs two ring snapshots at the union of their + virtual-node boundaries and groups the arcs whose owner set changed into + `OwnershipChange{Before, After, Added, Removed, Arcs}`. A joining node is a target owner; a + leaving/dead node is excluded as a target but still a source. +- **Migrate**: the coordinator reads the moved arcs from a current owner + (`/api/internal/antientropy/range`) and pushes them to each new owner through + `/api/internal/backfill`; a new owner is confirmed only on its ack. +- **GC**: once data has moved and the ring is flipped to the target placement (so read-repair + cannot re-add it), the old owner drops the arcs it no longer owns (`TSDB.DropSeriesInRanges` + via `/api/internal/rebalance/drop`): head flushed, then each block left, deleted whole, or + rewritten to keep only its owned series. GC runs only after new owners are confirmed at + quorum and never drops the last copy. +- **Lifecycle**: a join catches up out of routing (`joining`), promotes to `active`, then GCs + the displaced owner; a leave migrates while still `active` (reads stay complete), then goes + `leaving`, then removed. A node returning from `dead` has the over-replication its absence created reclaimed off the fallbacks. Driven by `POST /api/internal/cluster/{join,leave}`. -The ingestor and querier each build this client from `REPLICATION_FACTOR`/ -`WRITE_QUORUM`/`READ_QUORUM` and run the health monitor; only the ingestor wires the -hint store and the rebalance coordinator (the write owner). The `Coordinator` -(`internal/cluster/coordinator.go`) wraps the same ring with `RouteWrite`/`RouteRead` -helpers. The monolith `serve` path is single-node and does not use the ring. +Wiring: the ingestor and querier build this client from `REPLICATION_FACTOR`/`WRITE_QUORUM`/ +`READ_QUORUM` and run the health monitor; only the ingestor wires the hint store and the +rebalance coordinator (the write owner). The `Coordinator` (`internal/cluster/coordinator.go`) +wraps the same ring with `RouteWrite`/`RouteRead` helpers. The monolith `serve` path is +single-node and does not use the ring. ### Retention & Downsampling -**Downsampler** (`internal/retention/downsampler.go`): Runs the live raw → 1m → 1h -cascade (ADR-011). Each background pass advances a tier as far as its source is -durably closed, rolling sealed raw blocks into resolution-tagged **rollup blocks** -(`storage.RollupBlock`, format v2) that store six Gorilla-compressed aggregate columns -(min/max/sum/count/avg + a reset-aware counter increase) per series under -`/rollups//`. The 1h tier is *chained* from the 1m tier — a -count-weighted average and an additive increase — so it equals a 1h rollup built -directly from raw. A per-resolution `covered_through` watermark makes passes idempotent -and crash-recoverable (rollups are regenerable). At query time the planner -(`internal/query`) selects a resolution from the span and step and the executor reads -the column that matches the operation (`TSDB.QueryResolution` with a `RollupAggregate`): -avg for a bare value, max/min/sum/count for the matching `*_over_time` function, and the -increase column for `rate()` (ADR-025) — rolling up the freshest not-yet-closed tail on -the fly so the coarse series stays current. This query-time selection runs in both -deployments: the single-binary `serve`, whose TSDB implements the `ResolutionDataSource` -capability, **and** the cluster, where `service.StorageClient` implements the same -capability — the engine and planner are reused unchanged. In the cluster the querier picks -a resolution from the span/step against the intersection of the live nodes' advertised tiers -(`/api/internal/resolutions`), asks the replicas for it (`QueryAtMostResolution` on each node -serves the coarsest tier it holds at or below the request, reporting what it served), and -merges the coarse results by (series, window-timestamp). The coarse merge skips read-repair — -rollups are node-local derivations, so raw replication + read-repair is the convergence layer -(ADR-022) — and a series whose replicas served different resolutions falls back to a raw read, -so heterogeneity never breaks a query (ADR-011). Older v1 blocks (no increase column) load and -serve their five columns, with `rate()` falling back to raw for the spans they cover. - -**Enforcer** (`internal/retention/enforcer.go`): Per-resolution TTL cleanup. Raw -blocks expire first (and only once the finest rollup tier has captured them); each -rollup tier is kept longer, so a long-range query is still answerable from the 1h -tier after the raw behind it is gone. +**Downsampler** (`internal/retention/downsampler.go`, ADR-011): runs the live raw → 1m → 1h +rollup cascade. +- Each pass advances a tier as far as its source is durably closed, rolling sealed raw blocks + into resolution-tagged **rollup blocks** (`storage.RollupBlock`, format v2). Each stores six + Gorilla-compressed aggregate columns (min/max/sum/count/avg plus a reset-aware counter + increase) per series under `/rollups//`. +- The 1h tier is *chained* from 1m (count-weighted average, additive increase), so it equals a + 1h rollup built directly from raw. A per-resolution `covered_through` watermark makes passes + idempotent and crash-recoverable (rollups are regenerable). +- **Query-time selection**: the planner picks a resolution from the span/step; the executor + reads the column matching the operation (`TSDB.QueryResolution` with a `RollupAggregate`): + avg for a bare value, max/min/sum/count for the matching `*_over_time`, the increase column + for `rate()` (ADR-025). It rolls up the freshest not-yet-closed tail on the fly so the coarse + series stays current. +- Runs in both deployments: the single-binary `serve` (TSDB implements + `ResolutionDataSource`) and the cluster, where `service.StorageClient` implements the same + capability, so engine and planner are reused unchanged. The cluster querier picks a + resolution against the intersection of live nodes' advertised tiers + (`/api/internal/resolutions`), requests it (`QueryAtMostResolution` serves the coarsest + tier at or below), and merges by (series, window-timestamp). +- The coarse merge skips read-repair (rollups are node-local derivations, so raw replication + + read-repair is the convergence layer, ADR-022). A series whose replicas served different + resolutions falls back to a raw read, so heterogeneity never breaks a query. Older v1 blocks + (no increase column) serve their five columns, with `rate()` falling back to raw. + +**Enforcer** (`internal/retention/enforcer.go`): per-resolution TTL cleanup. Raw blocks expire +first (and only once the finest rollup tier has captured them); each rollup tier is kept +longer, so a long-range query is still answerable from the 1h tier after the raw behind it is +gone. ## Data Flow -1. **Ingest**: Samples arrive via TCP → (optional per-series/priority admission shaper, - ADR-027) → BatchWriter's bounded queue (block-then-shed backpressure) → drain → WAL - (group-commit fsync, ADR-026) → HeadBlock (in-order policy). The TCP server - bounds message size, applies a per-message read deadline, and caps concurrent - connections; a full queue sheds and NACKs the producer rather than growing memory - (ADR-023), and when admission is enabled the shed victims are chosen by lowest - priority / most over-budget series rather than uniformly. -2. **Flush**: Head swap + WAL rotate (cut) → Gorilla-compressed block written via - temp-dir + fsync + atomic rename → covered WAL segments reclaimed. -3. **Query**: Parser → Planner (also selects a rollup resolution from the span/step) - → merge(HeadBlock, Blocks) or read the chosen rollup tier → Executor → Result -4. **Stream**: WebSocket hub broadcasts metrics + stats to dashboard at 60fps. The - same per-tick per-series stream feeds the anomaly detector, whose raise/clear - transitions broadcast as `anomaly` frames (ADR-024). -5. **Downsample & retain**: the downsampler rolls sealed raw blocks into the 1m/1h - tiers (1h chained count-weighted from 1m); the enforcer expires each tier on its - own TTL, raw first (only once it has been rolled up) -6. **Recover**: On open, load blocks, then replay only the WAL beyond the blocks' - max low-water-mark → exactly-once reconstruction of the head. +1. **Ingest**: TCP → optional admission shaper (ADR-027) → BatchWriter's bounded + block-then-shed queue (ADR-023) → drain → WAL (group-commit fsync, ADR-026) → head + (in-order policy, ADR-015). A full queue sheds and NACKs the producer rather than growing + memory; with admission on, the shed victims are the lowest-priority or most-over-budget + series. +2. **Flush**: head swap + WAL rotate → Gorilla block via temp-dir + fsync + atomic rename → + covered WAL segments reclaimed. +3. **Query**: Parser → Planner (also selects a rollup resolution from the span/step) → + merge(head, blocks) or read the chosen rollup tier → Executor → result. +4. **Stream**: the WebSocket hub broadcasts metrics + stats at 60fps; the same per-tick + per-series stream feeds the anomaly detector, whose raise/clear transitions broadcast as + `anomaly` frames (ADR-024). +5. **Downsample & retain**: the downsampler rolls sealed raw blocks into the 1m/1h tiers (1h + chained count-weighted from 1m); the enforcer expires each tier on its own TTL, raw first + (only once it has been rolled up). +6. **Recover**: on open, load blocks, then replay only the WAL beyond the blocks' maximum + low-water-mark → exactly-once reconstruction of the head. diff --git a/CHANGELOG.md b/CHANGELOG.md index 1e9e34a..548a0f0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,129 +3,120 @@ ## Unreleased ### Replication & clustering -- Rebalancing on membership change (ADR-031): adding or removing a storage node now +- **Rebalancing on membership change** (ADR-031): adding or removing a storage node now re-derives placement **and moves the data**. `cluster.PlacementDelta` diffs the owner - sets before/after the change (at the union of the rings' virtual-node boundaries) into the - arcs that moved; the coordinator **migrates** each moved arc to its new owners by reusing - the anti-entropy range read (`/api/internal/antientropy/range`) + the handoff backfill - apply (`/api/internal/backfill`), then **GCs** the data a node no longer owns via the new - `TSDB.DropSeriesInRanges` (`/api/internal/rebalance/drop`) — head flushed, then blocks - left, deleted whole, or rewritten to keep only owned series. GC runs only after the new - owners confirm receipt at quorum and never drops the last copy; a joining node catches up - out of routing before it serves, a leaving node re-homes its ranges before removal, and a - node returning from `dead` has the over-replication its absence created reclaimed off the - fallbacks. Reads stay complete throughout. Driven by `POST - /api/internal/cluster/{join,leave}`; `meridian_rebalance_*` exposes - migrations/bytes/GC/joins/leaves. On by default for the cluster tier - (`cluster.rebalance`); disabled, a membership change re-derives routing without moving - data (the pre-ADR-031 behaviour). -- Proactive anti-entropy (ADR-030): a rate-limited, jittered background sweep on the + sets before/after the change (at the union of the rings' virtual-node boundaries) into + the arcs that moved; the coordinator migrates each to its new owners — reusing the + anti-entropy range read and the handoff backfill apply — then GCs the data a node no + longer owns via the new `TSDB.DropSeriesInRanges` (head flushed, then blocks left, + deleted whole, or rewritten to keep only owned series). GC runs only after the new owners + confirm at quorum and never drops the last copy; a join catches up out of routing before + serving, a leave re-homes its ranges before removal, and a return from `dead` reclaims + over-replication off the fallbacks — reads stay complete throughout. Driven by `POST + /api/internal/cluster/{join,leave}`; `meridian_rebalance_*` exposes migrations/bytes/GC/ + joins/leaves. On by default (`cluster.rebalance`); off, routing re-derives without moving + data. +- **Proactive anti-entropy** (ADR-030): a rate-limited, jittered background sweep on the ingestor converges co-replicas that read-repair and hinted handoff never reach — cold data, an unobserved partial write, a hint dropped past the cap, a series no longer written. Each storage node summarises its data as **Merkle range digests** over - `(ring-range × time-window)` buckets (`/api/internal/antientropy/digest`); the sweep - round-robins the ring's replica groups, and where two replicas' digest roots differ it - reads only the divergent window (`/api/internal/antientropy/range`) and pushes each the - points it lacks back through the same out-of-order-tolerant `/api/internal/backfill` - apply (ADR-029) — bidirectional gap-fill to the window union. Agreement transfers - nothing. The ring hash is injected so the storage node stays ring-agnostic. - `meridian_anti_entropy_*` exposes rounds/divergence/repairs/bytes. On by default for the - cluster tier (`cluster.anti_entropy`); disabled it is exactly the ADR-029 behaviour. -- Hinted handoff (ADR-029): a write that can't reach a natural replica (it's down or - catching up) buffers a durable, bounded hint while quorum still succeeds on the - survivors, and replays it on the replica's return so it **fully converges** — including - an *interior* gap read-repair can't fill (read-repair only converges forward; storage - rejects out-of-order — ADR-015). Replay applies through a new out-of-order-tolerant - backfill path (`TSDB.Backfill` + a distinct WAL frame + `/api/internal/backfill`), and - a returning replica catches up in the reserved `joining` state — excluded from live - routing — before it is promoted to `active`, so it is made whole before it can strand a - gap. The hint store is per-target, capped in samples (drop-oldest past the cap), and - rebuilt from disk on restart; `meridian_handoff_*` exposes pending/replayed/dropped. - On by default for the cluster tier (`cluster.handoff`); disabled it is exactly ADR-022. + `(ring-range × time-window)` buckets; the sweep round-robins the ring's replica groups, + and where two replicas' roots differ it reads only the divergent window and pushes each + the points it lacks through the same out-of-order-tolerant backfill apply (ADR-029) — + bidirectional gap-fill to the window union. Agreement transfers nothing; the ring hash is + injected so storage stays ring-agnostic. `meridian_anti_entropy_*` exposes + rounds/divergence/repairs/bytes. On by default (`cluster.anti_entropy`); disabled it is + exactly ADR-029. +- **Hinted handoff** (ADR-029): a write that can't reach a natural replica (down or + catching up) buffers a durable, bounded hint while quorum still succeeds on the survivors, + and replays it on the replica's return so it **fully converges** — including an *interior* + gap read-repair can't fill (read-repair only converges forward; storage rejects + out-of-order, ADR-015). Replay applies through a new out-of-order-tolerant backfill path + (`TSDB.Backfill` + a distinct WAL frame), and a returning replica catches up in the + reserved `joining` state — excluded from live routing — before promotion to `active`, so + it is made whole before it can strand a gap. The hint store is per-target, capped in + samples (drop-oldest past the cap), and rebuilt from disk on restart; `meridian_handoff_*` + exposes pending/replayed/dropped. On by default (`cluster.handoff`); disabled it is + exactly ADR-022. ### Write path & flow control -- Per-series fair-share / priority-class load shedding (ADR-027): an opt-in admission - shaper in front of the bounded ingest queue makes overload shedding selective - instead of uniform. Priority bands (a label or `__name__` match → a capacity - ceiling) shed low priority before high; per-series token buckets throttle a hot or - high-cardinality series so it can't starve well-behaved ones. Both engage only under - contention, per-series state is bounded (sharded — a cardinality flood can't grow - it), and order within a series is preserved. Wired into the monolith BatchWriter and - the service WritePool; drops fold into `meridian_dropped_samples_total` and are - broken down by class/reason/series-bucket via `meridian_admission_*`. Off by default - (`ingestion.admission`), leaving ADR-023's uniform shedding as the fallback. +- **Per-series fair-share / priority-class load shedding** (ADR-027): an opt-in admission + shaper in front of the bounded ingest queue makes overload shedding selective instead of + uniform. Priority bands (a label or `__name__` match → a capacity ceiling) shed low + priority before high; per-series token buckets throttle a hot or high-cardinality series + so it can't starve well-behaved ones. Both engage only under contention, per-series state + is bounded (sharded, so a cardinality flood can't grow it), and order within a series is + preserved. Wired into the monolith BatchWriter and the service WritePool; drops fold into + `meridian_dropped_samples_total` and break down by class/reason/series-bucket via + `meridian_admission_*`. Off by default (`ingestion.admission`), leaving ADR-023's uniform + shedding as the fallback. ### Storage & durability -- WAL group commit: a single committer goroutine coalesces concurrently-submitted - frames behind one fsync, so a write still returns only after its own frame is - durable while concurrent writers stop serializing one fsync at a time — measured - ~30–37× WAL write throughput at 64 concurrent writers and ~4× at 8 (ADR-026). - Default on with zero linger; the on-disk frame format is byte-for-byte unchanged - and the synchronous path remains available (`storage.wal_group_commit`). +- **WAL group commit** (ADR-026): a single committer goroutine coalesces + concurrently-submitted frames behind one fsync, so a write still returns only after its + own frame is durable while concurrent writers stop serializing one fsync at a time — + measured ~30–37× WAL write throughput at 64 concurrent writers, ~4× at 8. Default on with + zero linger; the on-disk frame format is byte-for-byte unchanged and the synchronous path + remains available (`storage.wal_group_commit`). ### Anomaly detection -- Seasonal Holt-Winters model (ADR-028): per-series scoring now sits behind a model - interface, with an additive level+trend+seasonal **Holt-Winters** model selectable - via `anomaly.mode: holt_winters` alongside the default EWMA. It derives the diurnal - phase from the sample timestamp, warms up over one full season, and scores each value - against the band for its own time of day — so it flags a value normal globally but - abnormal for that phase (a midday level in the small hours, a scheduled dip that fails - to happen), which EWMA cannot. The Huber clamp, scale floor, debounce/hysteresis and - event emission are shared by both models; EWMA's math and defaults are unchanged. - `anomaly.{mode,season_length,season_period}` are configurable; the active model - appears on `/api/v1/stats`, `/api/v1/anomalies`, `meridian_anomaly_model_info`, and the - dashboard's Anomalies panel. +- **Seasonal Holt-Winters model** (ADR-028): per-series scoring now sits behind a model + interface, with an additive level+trend+seasonal **Holt-Winters** model selectable via + `anomaly.mode: holt_winters` alongside the default EWMA. It derives the diurnal phase from + the sample timestamp, warms up over one full season, and scores each value against the + band for its own time of day — so it flags a value normal globally but abnormal for that + phase (a midday level in the small hours, a scheduled dip that fails to happen), which + EWMA cannot. The Huber clamp, scale floor, debounce/hysteresis and event emission are + shared by both models; EWMA's math and defaults are unchanged. The active model appears on + `/api/v1/stats`, `/api/v1/anomalies`, `meridian_anomaly_model_info`, and the dashboard. ## v0.2.0 — 2026-06-18 -Meridian as it stands at this release: a distributed time-series database in Go -with a canvas-rendered React dashboard. See [DECISIONS.md](DECISIONS.md) for the -full ADR set and [PERFORMANCE.md](PERFORMANCE.md) for measured numbers. +Meridian as it stands at this release: a distributed time-series database in Go with a +canvas-rendered React dashboard. See [DECISIONS.md](DECISIONS.md) for the full ADR set and +[PERFORMANCE.md](PERFORMANCE.md) for measured numbers. ### Storage & durability -- Gorilla compression (delta-of-delta timestamps + XOR floats): ~28× on regular +- **Gorilla compression** (delta-of-delta timestamps + XOR floats): ~28× on regular integer-like gauges, ~2× on continuously varying floats (ADR-002). -- CRC32-framed WAL with 128 MB segment rotation; a corrupt frame resyncs to the - next 8-byte boundary instead of discarding the rest of the segment (ADR-003). -- Crash-consistent flush: atomic head-swap + WAL-rotate cut, temp-dir → fsync → - rename block writes, and a per-block WAL low-water-mark for exactly-once replay - (ADR-016). -- Out-of-order samples are rejected and counted; exact duplicates are deduped - (ADR-015). +- **CRC32-framed WAL** with 128 MB segment rotation; a corrupt frame resyncs to the next + 8-byte boundary instead of discarding the rest of the segment (ADR-003). +- **Crash-consistent flush**: atomic head-swap + WAL-rotate cut, temp-dir → fsync → rename + block writes, and a per-block WAL low-water-mark for exactly-once replay (ADR-016). +- **Out-of-order handling**: out-of-order samples are rejected and counted; exact duplicates + are deduped (ADR-015). ### Query engine -- PromQL subset with stepped range/matrix evaluation, unary `+`/`-`, bare - selectors (`{job="x"}`), `=`/`!=`/`=~`/`!~`, and compound durations (`1h30m`). -- Correct `rate()` (per selector-range, counter-reset corrected and - extrapolated), `histogram_quantile()` (cumulative-`le` interpolation), - vector↔vector ops with label matching and IEEE-754 division, `topk`/`bottomk`, - and both `by`/`without` grouping (ADR-014). +- **PromQL subset** with stepped range/matrix evaluation, unary `+`/`-`, bare selectors + (`{job="x"}`), `=`/`!=`/`=~`/`!~`, and compound durations (`1h30m`). +- **Correct semantics**: `rate()` (per selector-range, counter-reset corrected and + extrapolated), `histogram_quantile()` (cumulative-`le` interpolation), vector↔vector ops + with label matching and IEEE-754 division, `topk`/`bottomk`, and both `by`/`without` + grouping (ADR-014). ### Cluster (docker-compose tier) -- Quorum replication over a 64-bit consistent-hash ring (defaults N=3/W=2/R=2): +- **Quorum replication** over a 64-bit consistent-hash ring (defaults N=3/W=2/R=2): write-to-all at W acks, scatter reads at R with merge + async read-repair, and health-driven membership (ADR-022). -- Prometheus `/metrics` on every service — monolith, gateway, ingestor, storage, +- **Prometheus `/metrics`** on every service — monolith, gateway, ingestor, storage, querier, compactor (ADR-019). ### Flow control & detection -- Write-path backpressure: bounded block-then-shed ingest queues with HTTP 429 + - `Retry-After` / TCP NACK and cumulative drop/shed/backpressure metrics - (ADR-023). -- Streaming anomaly detection: per-series EWMA baseline + dispersion robust to a - moving diurnal baseline, with alerts over the WebSocket hub (ADR-024). +- **Write-path backpressure**: bounded block-then-shed ingest queues with HTTP 429 + + `Retry-After` / TCP NACK and cumulative drop/shed/backpressure metrics (ADR-023). +- **Streaming anomaly detection**: per-series EWMA baseline + dispersion robust to a moving + diurnal baseline, with alerts over the WebSocket hub (ADR-024). ### Operations -- Downsampling cascade raw → 1m → 1h with count-weighted chaining and query-time - resolution selection in the single binary; per-resolution retention (ADR-011). - Cluster query-time resolution selection is documented future work — storage - nodes generate rollups, but the querier still reads raw. -- Storage-service TSDB timings (`STORAGE_BLOCK_DURATION` / `_FLUSH_INTERVAL` / - `_RETENTION`) are configurable via the environment. +- **Downsampling cascade** raw → 1m → 1h with count-weighted chaining and query-time + resolution selection in the single binary; per-resolution retention (ADR-011). Cluster + query-time resolution selection is documented future work — storage nodes generate + rollups, but the querier still reads raw. +- **Configurable storage timings**: storage-service TSDB timings (`STORAGE_BLOCK_DURATION` / + `_FLUSH_INTERVAL` / `_RETENTION`) are set via the environment. ### Dashboard -- "Precision Instrument" design language (ADR-020, ADR-021): one accent, - self-hosted fonts, tabular-mono numerics, a signature strip-chart with a cursor - crosshair readout, dark/light themes, an accessibility floor, and reconnect - handling. WAL size is reported in bytes (previously mislabeled as a segment - count). +- **"Precision Instrument" design language** (ADR-020, ADR-021): one accent, self-hosted + fonts, tabular-mono numerics, a signature strip-chart with a cursor crosshair readout, + dark/light themes, an accessibility floor, and reconnect handling. WAL size is reported in + bytes (previously mislabeled as a segment count). diff --git a/DECISIONS.md b/DECISIONS.md index b07190c..c6cb91f 100644 --- a/DECISIONS.md +++ b/DECISIONS.md @@ -2,111 +2,142 @@ ## ADR-001: Go as Implementation Language -**Status**: Accepted +**Status**: Accepted + **Context**: Need a systems language with good concurrency, fast compilation, and -single-binary deployment. +single-binary deployment. + **Decision**: Go 1.25 with zero CGO dependencies. CI installs the toolchain from -`go.mod` (`go-version-file`) so the build version never drifts from the module. +`go.mod` (`go-version-file`) so the build version never drifts from the module. + **Consequences**: Simple cross-compilation, no shared library issues in containers, goroutine-per-connection model fits well. ## ADR-002: Gorilla Compression for Time-Series Data -**Status**: Accepted +**Status**: Accepted + **Context**: Time-series data exhibits high temporal locality and value similarity -that generic compression (gzip, lz4) cannot exploit effectively. +that generic compression (gzip, lz4) cannot exploit effectively. + **Decision**: Implement Facebook's Gorilla encoding with delta-of-delta timestamps and XOR float encoding. Extended to 64-bit millisecond timestamps (paper uses -32-bit seconds). 4-byte count header for decoder bootstrapping. +32-bit seconds). 4-byte count header for decoder bootstrapping. + **Consequences**: Achieves 20-30x compression on regular metrics vs. 3-5x with generic algorithms. Small decoder that can stream without seeking. ## ADR-003: CRC32-Framed WAL with Segment Rotation -**Status**: Accepted -**Context**: Must survive process crashes without losing acknowledged writes. +**Status**: Accepted + +**Context**: Must survive process crashes without losing acknowledged writes. + **Decision**: Write-ahead log with CRC32 checksums, 8-byte alignment for -efficient reads, and automatic rotation at 128 MB segments. +efficient reads, and automatic rotation at 128 MB segments. + **Consequences**: Crash recovery by replaying WAL. Segment rotation keeps individual files manageable and enables garbage collection. ## ADR-004: Inverted Index with Sorted Slices, Not Roaring Bitmaps -**Status**: Accepted +**Status**: Accepted + **Context**: Need an inverted index for label-based series lookup. Roaring bitmaps -are the standard choice but add a dependency. +are the standard choice but add a dependency. + **Decision**: Use `map[string]map[string][]uint64` with sorted slices and -set intersection/union via merge-join. +set intersection/union via merge-join. + **Consequences**: Zero external dependencies. Performance is adequate for the expected scale (< 100K series). Not optimal for millions of series. ## ADR-005: JSON-over-TCP Ingestion Protocol -**Status**: Accepted +**Status**: Accepted + **Context**: Protobuf would be the natural choice for ingestion, but protoc is not -available in the build environment. +available in the build environment. + **Decision**: JSON-over-TCP with newline framing. Same message structure as the -proto definition for future migration. +proto definition for future migration. + **Consequences**: ~3x larger on the wire than protobuf. Simpler debugging with netcat/telnet. Easy to switch to protobuf later since struct shapes match. ## ADR-006: PromQL Subset via Recursive Descent Parser -**Status**: Accepted -**Context**: Users expect a familiar query language for time-series databases. +**Status**: Accepted + +**Context**: Users expect a familiar query language for time-series databases. + **Decision**: Implement a PromQL subset with recursive descent parsing. Supports vector/range selectors, label matchers (=, !=, =~, !~), aggregations (sum, avg, min, max, count, topk, bottomk), functions (rate, histogram_quantile), binary -operators, and group-by clauses. +operators, and group-by clauses. + **Consequences**: No parser generator dependency. Easy to extend. Covers the most common monitoring use cases. ## ADR-007: Consistent Hash Ring for Data Distribution -**Status**: Accepted +**Status**: Accepted + **Context**: Need to distribute series across cluster nodes with even load -balance and minimal disruption during scaling. +balance and minimal disruption during scaling. + **Decision**: SHA256-based consistent hash ring with configurable virtual nodes -(default 64 per node). Series assigned by MetricKey = hash(sorted labels). +(default 64 per node). Series assigned by MetricKey = hash(sorted labels). + **Consequences**: Adding/removing a node only redistributes ~1/N of data. Virtual nodes smooth out hash distribution. ## ADR-008: requestAnimationFrame Batching for WebSocket Messages -**Status**: Accepted +**Status**: Accepted + **Context**: WebSocket messages arrive faster than the display refresh rate. Processing each message individually causes excessive React re-renders and -dropped frames. +dropped frames. + **Decision**: Buffer incoming WebSocket messages and flush them in a single batch -on each requestAnimationFrame callback. +on each requestAnimationFrame callback. + **Consequences**: Dashboard maintains 60fps even at high ingestion rates. Slight increase in perceived latency (up to 16ms) which is imperceptible. ## ADR-009: Canvas-Based Chart Rendering, No Chart Library -**Status**: Accepted +**Status**: Accepted + **Context**: The spec requires zero chart dependencies (no D3, Chart.js, -Recharts). Charts must render at 60fps for live streaming data. +Recharts). Charts must render at 60fps for live streaming data. + **Decision**: Direct Canvas 2D API rendering with custom TimeSeriesChart component. Features: area fills, a one-shot load transition, auto-scaling axes, and multi-series support. (The original revision also drew decorative glow; that -was removed under ADR-020.) +was removed under ADR-020.) + **Consequences**: Full control over rendering pipeline. No dependency bloat. Requires manual hit-testing for interactivity (tooltips, zoom). ## ADR-010: React Context + useReducer for State Management -**Status**: Accepted +**Status**: Accepted + **Context**: Dashboard state (theme, time range, query results, live metrics, -cluster nodes) needs to be shared across many components. +cluster nodes) needs to be shared across many components. + **Decision**: Single DashboardContext with useReducer pattern. No external state -library (Redux, Zustand, etc.). +library (Redux, Zustand, etc.). + **Consequences**: Zero dependencies for state management. Action-based updates are predictable and debuggable. Adequate for the component count. ## ADR-011: Three-Tier Downsampling Cascade **Status**: Accepted (implemented) + **Context**: High-resolution data is expensive to store and slow to scan over long ranges: a 30-day view at 5-second resolution is ~518k points per series, far more than a screen can show or a user needs. The rollup math (`Rollup`) existed but was @@ -220,33 +251,41 @@ pass; rollups are regenerable, so the extra on-disk state is never authoritative ## ADR-012: Single-Binary Architecture -**Status**: Accepted +**Status**: Accepted + **Context**: Deployment simplicity is a core design goal. Users should be able to -run `./meridian serve` and have a complete system. +run `./meridian serve` and have a complete system. + **Decision**: Single Go binary bundles server, ingestion, query engine, simulator, -CLI tools, and dashboard static files. +CLI tools, and dashboard static files. + **Consequences**: No orchestration required for single-node deployment. Dashboard assets are embedded or served from a directory. Trade-off: binary size is larger. ## ADR-013: Diurnal Simulation with Spike Injection -**Status**: Accepted +**Status**: Accepted + **Context**: Testing and demos require realistic-looking metric data, not random -noise. Real infrastructure exhibits predictable daily patterns. +noise. Real infrastructure exhibits predictable daily patterns. + **Decision**: Simulator generates diurnal curves (peak at 14:00 local time) with random spike injection (10% probability per host per cycle) and memory drift -(slow monotonic increase with periodic resets). +(slow monotonic increase with periodic resets). + **Consequences**: Dashboard screenshots and demos look realistic. Compression benchmarks reflect real-world data patterns. ## ADR-014: PromQL Evaluation Semantics **Status**: Accepted + **Context**: The query engine advertised a PromQL subset, but several pieces were incorrect rather than merely incomplete: a range window was subtracted twice, `rate()` divided by the sample span, `histogram_quantile()` ignored `le` buckets, vector÷vector returned the left operand, and `/0` returned `0`. Making the engine *correct* required committing to specific semantics. + **Decision**: - **Stepped range evaluation → matrix.** `Execute(start, end, step)` evaluates the expression as an instant query at each `t` in `{start, start+step, …, end}` and @@ -283,6 +322,7 @@ vector÷vector returned the left operand, and `/0` returned `0`. Making the engi points (floored at 1s); `start > end` is rejected; the step count is capped at 11000. The cap bounds output size and pre-empts a denial-of-service via attacker-controlled `start`/`end`/`step`. + **Consequences**: Evaluation is stepped/range — a query returns a matrix with one point per step per series — so `rate(x[5m])`, `sum(...) by (...)`, and `a/b` render as smooth multi-point lines rather than single values. The planner's `TimeRange` @@ -294,6 +334,7 @@ matching, counter resets across windows, and the step-count guard. ## ADR-015: Reject Out-of-Order Samples **Status**: Accepted + **Context**: Samples were appended to a series in arrival order without any ordering check, and `WriteBlock` took `Timestamps[0]`/`Timestamps[last]` as a block's min/max assuming sorted input. Ingesting `100, 50, 200, 10` therefore produced an *inverted* @@ -301,7 +342,9 @@ block (`minTime=200, maxTime=10`), which silently dropped overlapping queries (`Overlaps` compares against these bounds) and poisoned retention (the enforcer deletes by `MaxTime`). A policy was required; the options were reject, sort-on-flush, or full out-of-order support. + **Decision**: **Reject**, the Prometheus-classic model. + - A sample with a timestamp strictly older than the series' last is dropped and counted in `meridian_out_of_order_samples_total`. - A sample whose timestamp equals the series' last is **deduplicated** if its value @@ -320,6 +363,7 @@ or full out-of-order support. accepted into the new head. Blocks may then overlap in time, which the read path already handles (it merges and sorts). Cross-flush ordering is intentionally not enforced, matching the head-relative out-of-order window of mainstream TSDBs. + **Consequences**: Series stay sorted, so `Timestamps[0]`/`[last]` range checks and non-inverted block bounds hold; retention expires correctly. Out-of-order data is visibly dropped and counted rather than silently corrupting bounds. Full @@ -329,6 +373,7 @@ future work. ## ADR-016: Crash-Consistent Flush with a Per-Block WAL Low-Water-Mark **Status**: Accepted + **Context**: `Flush()` ran `WriteBlock(head)` → `head.Reset()` → `wal.Truncate()` with no lock spanning the three steps. A sample ingested after the block snapshot but before `Reset()` was discarded from the head *and* erased from the WAL by `Truncate()` @@ -336,6 +381,7 @@ before `Reset()` was discarded from the head *and* erased from the WAL by `Trunc next `Open()` to load the block *and* replay the whole WAL — a double-count. Block writes were also non-atomic (no temp/rename, nothing fsynced, several ignored I/O errors). + **Decision**: A three-phase flush with an atomic in-memory cut and a durable low-water-mark. 1. **Cut (under `db.mu` as writer).** Capture the old head, install a fresh head, and @@ -360,6 +406,7 @@ low-water-mark. cannot record a higher mark that would skip the failed flush's still-uncovered segments; that data remains in the WAL and is recovered on the next open. Leftover temp block dirs from an interrupted write are removed on open. + **Consequences**: Exactly-once recovery. A crash **before the block is durable** leaves no committed block and an un-truncated WAL, so replay rebuilds the data once from the WAL (the in-memory cut is lost with the crash, so there is no persistent @@ -376,13 +423,16 @@ low-water-mark cut — is implemented in ADR-026. ## ADR-017: Ingestion Rate as a Windowed Rate, Cumulative Count for the Counter **Status**: Accepted + **Context**: `TSDB.IngestionRate()` returned the cumulative `ingested` counter. The dashboard samples that value once per second and charts it as a rate, so it drew a monotonically rising line instead of throughput. Separately, the Prometheus exposition needs a *cumulative* `meridian_samples_ingested_total` — a `..._total` counter is correct by Prometheus convention (the scraper computes `rate()`). One method could not honestly serve both roles. + **Decision**: Split the two concerns. + - `IngestionRate()` returns a **windowed** samples/sec rate: a moving average over `RateWindow` (default 5s), fed by a background sampler that records the cumulative count every `RateSampleInterval` (default 1s). Idle intervals contribute the same @@ -394,6 +444,7 @@ method could not honestly serve both roles. windowed rate. Wire types stay `int64` (the rate is rounded to whole samples/sec), so neither the dashboard nor the inter-service protocol changes. Per-node rates are additive, so the gateway sums them into the cluster rate. + **Consequences**: The dashboard charts a true rate that tracks load and falls back to ~0 when idle, with no dashboard change. The Prometheus counter remains a proper cumulative counter. The reported rate lags a step change by up to the window, and @@ -404,10 +455,12 @@ rate read ~140/s and fell to 0 within a window of going idle, while ## ADR-018: CORS Restricted to Configured Origins (Default Localhost) **Status**: Accepted + **Context**: Both the monolith and the gateway returned `Access-Control-Allow-Origin: *` for every method, POST included. Any web page the operator visited could therefore script cross-origin reads and writes against a Meridian instance reachable from that browser (e.g. on a private network). + **Decision**: Replace the blanket wildcard with an origin-checked middleware that echoes only permitted origins. - Default (unconfigured): allow only `localhost` / `127.0.0.1` / `[::1]` origins — @@ -418,6 +471,7 @@ echoes only permitted origins. - The matched origin is reflected back (with `Vary: Origin`) rather than `"*"`, so the policy is also correct for credentialed requests. A disallowed cross-origin request receives no CORS headers and is blocked by the browser. + **Consequences**: Same-origin dashboard use is unchanged (a same-origin request sends no `Origin` header and proceeds normally). Cross-origin browser access now requires explicit configuration. The internal service APIs (storage/querier/ingestor) are not @@ -427,10 +481,12 @@ browser-facing surfaces and keep their permissive CORS; the public entry points ## ADR-019: Prometheus /metrics on Every Service **Status**: Accepted + **Context**: Only the monolith exposed `/metrics`. In the docker-compose topology the gateway, querier, storage, ingestor, and compactor had no scrape endpoint, so the "cluster" the README advertises could not actually be observed by a Prometheus-compatible collector. + **Decision**: Register `/metrics` on every service through shared `WriteStorageMetrics`/`WriteServiceMetrics` helpers. Storage nodes expose the full storage metrics — the cumulative `meridian_samples_ingested_total`, @@ -439,6 +495,7 @@ storage bytes by layer, and compression ratio. Every service additionally expose `meridian_up` and `meridian_uptime_seconds`, and the gateway reports connected WebSocket clients. The monolith handler reuses the same storage helper so a storage node emits identical metrics whether run as the monolith or as the storage service. + **Consequences**: The cluster is scrapeable end-to-end and metric names/labels are consistent across deployment modes. The endpoints are unauthenticated and intended for an internal scrape network (the same trust boundary as the internal RPC APIs). @@ -446,18 +503,23 @@ for an internal scrape network (the same trust boundary as the internal RPC APIs ## ADR-020: Visual Design Language — "Precision Instrument" **Status**: Accepted -**Context**: The dashboard was functional but read as templated / AI-generated. -The concrete tells: the brand palette was verbatim Mantine indigo (`#4c6ef5`); -every surface used the same glassmorphism (backdrop-blur + lifted shadow + -`rounded-xl`); charts carried decorative glow and a rainbow compression gradient -that encoded no scale; `Inter` was named ~30 times but never actually loaded (a -silent system-font fallback, with the DOM and canvas free to disagree); two -styling systems coexisted (semantic CSS-var inline styles next to hardcoded -Tailwind grays monkey-patched per theme); three different `formatBytes` rendered -the same quantity differently; and the chrome carried a self-congratulatory, -partly false "60fps" boast backed by an always-on `requestAnimationFrame` meter. -A single committed point of view was needed — and consistency, not novelty, is -what removes the generated feel. + +**Context**: The dashboard was functional but read as templated / AI-generated. The +concrete tells: +- the brand palette was verbatim Mantine indigo (`#4c6ef5`); +- every surface used the same glassmorphism (backdrop-blur + lifted shadow + `rounded-xl`); +- charts carried decorative glow and a rainbow compression gradient that encoded no scale; +- `Inter` was named ~30 times but never actually loaded (a silent system-font fallback, + with the DOM and canvas free to disagree); +- two styling systems coexisted (semantic CSS-var inline styles next to hardcoded Tailwind + grays monkey-patched per theme); +- three different `formatBytes` rendered the same quantity differently; +- the chrome carried a self-congratulatory, partly false "60fps" boast backed by an + always-on `requestAnimationFrame` meter. + +A single committed point of view was needed — and consistency, not novelty, is what removes +the generated feel. + **Decision**: Adopt one design language — *Precision Instrument* (mission-console-meets-chronometer: calm, engraved, exact), derived from the name *Meridian* (navigation / transit instruments). Deliberately avoid the other two @@ -501,6 +563,7 @@ AI-default looks as well (near-black + acid accent; cream + serif + terracotta). onMouseEnter/Leave hover hacks are gone. `utils/format.ts` is the single source for byte / number / duration / time formatting, so identical quantities render identically everywhere. + **Consequences**: One coherent identity per theme, verified in both dark and light: fonts load over the network (woff2 under `/assets`) and the canvas shares the DOM family; there is no Mantine indigo, no glow, and no glassmorphism; @@ -513,14 +576,20 @@ pass — is delivered in **ADR-021**. ## ADR-021: Visual Hierarchy, the Signature Instrument Chart, and the A11y Floor **Status**: Accepted -**Context**: ADR-020 committed the "Precision Instrument" language but left the -identity work for a follow-up: every panel was a uniform `.card` in a flat grid -(the primary query action and a tertiary stat tile carried identical weight), -heights were eyeballed magic pixels (`h-[294px]`, canvases pinned to 140/160/180), -there was no memorable element, the four empty states were worded four different -ways and drawn at canvas center, the logo was a generic circle-and-rising-line, -cluster nodes were two-letter codes, and there was no keyboard/focus/reduced-motion -story. A consistent instrument *needs* a hierarchy and one bold focal point. + +**Context**: ADR-020 committed the "Precision Instrument" language but left the identity +work for a follow-up: +- every panel was a uniform `.card` in a flat grid (the primary query action and a tertiary + stat tile carried identical weight); +- heights were eyeballed magic pixels (`h-[294px]`, canvases pinned to 140/160/180); +- there was no memorable element; +- the four empty states were worded four different ways and drawn at canvas center; +- the logo was a generic circle-and-rising-line; +- cluster nodes were two-letter codes; +- there was no keyboard/focus/reduced-motion story. + +A consistent instrument *needs* a hierarchy and one bold focal point. + **Decision**: - **Hierarchy through size, not decoration.** A `Panel` component renders one surface in three tiers — `primary` (the query bar and result), `secondary` @@ -560,6 +629,7 @@ story. A consistent instrument *needs* a hierarchy and one bold focal point. per-component focus styles, a global `prefers-reduced-motion` reset backs the `motion-safe:` utilities, and the layout is responsive to a phone width (header wraps, grids stack, stat tiles reflow). + **Consequences**: One coherent, hierarchical identity per theme, verified in dark, light, and at a 390px width by driving the live dashboard in a headless browser: the crosshair readout tracks the cursor, the empty/loading/error/reconnect states @@ -569,6 +639,7 @@ one place; everything else stays quiet, which is what makes it land. ## ADR-022: Replication — Quorum Writes, Quorum Reads, and Read-Repair **Status**: Accepted + **Context**: The project claimed "consistent-hash clustering with configurable replication," but `internal/cluster` (the ring and coordinator) was dead code with zero non-test importers: the live write path sharded each series to a single node @@ -578,6 +649,7 @@ returned silent partial results. The ring also truncated its hash to 32 bits and ignored node state, so it could route to dead nodes. The goal: make replication real over the existing docker-compose tier (ingestor → storage ← querier) while leaving the single-binary monolith single-node. + **Decision**: Adopt a Dynamo-style quorum model with N, W, R configured in `cluster` (defaults N=3, W=2, R=2; `Validate` enforces 1≤W,R≤N and W+R>N for read-your-writes). The consistent-hash ring (`internal/cluster`, now widened to a @@ -639,6 +711,7 @@ smaller than the configured N degrades gracefully rather than rejecting every wr ## ADR-023: Write-Path Backpressure — Bounded Queues with Block-Then-Shed Load Shedding **Status**: Accepted + **Context**: The ingest buffers were effectively unbounded. The monolith `BatchWriter` appended into a slice that grew without limit; the ingestor fanned out an unbounded number of concurrent quorum writes (one goroutine per HTTP/TCP @@ -649,6 +722,7 @@ bottleneck — the arrival rate exceeds the service rate, the backlog grows, and process simply consumes memory until it OOMs. There was no flow control: no signal to the producer to slow down, and no bound on resident work. A policy was required: how much to buffer, when to push back, and what to do past the limit. + **Decision**: Put a **bounded queue between accept and the drain-to-storage worker** on every ingest path, with **block-then-shed** enqueue. This is flow control by a bounded queue: queue depth ≈ arrival_rate × service_time (Little's @@ -734,6 +808,7 @@ wiring. ## ADR-024: Streaming Anomaly Detection — EWMA Baseline + Dispersion, Robust to a Moving Baseline **Status**: Accepted + **Context**: The telemetry on the live path is not stationary. The simulator (and real infrastructure) produces a slow **diurnal swing** (a 24-hour load cycle) and slow **drift** (a leaking memory gauge that climbs for minutes then resets), with @@ -991,6 +1066,7 @@ the whole-batch error path, and that one fsync covers many frames. ## ADR-027: Per-Series Fair-Share and Priority-Class Load Shedding **Status**: Accepted + **Context**: The write-path backpressure of ADR-023 bounds memory by shedding past a cap, but it sheds **uniformly**: whatever arrives next when the queue is full is dropped, regardless of which series it belongs to or how important it is. Two @@ -1004,6 +1080,7 @@ spends the scarce queue on whatever shows up rather than on the traffic that matters. ADR-023 named both gaps as deferred. The constraint from ADR-015 is that order **within** a single series must be preserved (the in-order head depends on it), so any fairness scheme must act *across* series only. + **Decision**: Add an **admission shaper** (`internal/backpressure.Shaper`) that is consulted **before** the bounded queue on every ingest path. It does not hold samples and is not the memory bound — the queue's capacity remains the hard cap @@ -1204,6 +1281,7 @@ deployment is byte-for-byte unchanged until it opts in. ## ADR-029: Hinted Handoff — Buffer Writes for a Down Replica, Replay on Return **Status**: Accepted + **Context**: ADR-022 made replication real but explicitly deferred hinted handoff, and that gap has a sharp edge. Two mechanisms were supposed to keep a replica consistent through downtime, and neither closes the hole: @@ -1311,6 +1389,7 @@ byte-for-byte as ADR-022. ## ADR-030: Proactive Anti-Entropy — Merkle Range Digests for Background Replica Convergence **Status**: Accepted + **Context**: After ADR-022 (replication) and ADR-029 (hinted handoff), every convergence path is still *triggered* — and each leaves a gap the others do not close: @@ -1416,6 +1495,7 @@ behaviour. ## ADR-031: Rebalancing on Membership Change — Migrate to New Owners, GC the Un-owned **Status**: Accepted + **Context**: The consistent-hash ring (ADR-022) re-derives *placement* the instant membership changes — add a node and `GetNodes`/`PreferenceList`/`ReplicaGroups` immediately return the new owner set; remove one and routing flows to the survivors. But the **data** diff --git a/PERFORMANCE.md b/PERFORMANCE.md index e52237e..4d28c70 100644 --- a/PERFORMANCE.md +++ b/PERFORMANCE.md @@ -100,71 +100,66 @@ returns only after the fsync covering its frame. These paths have no repeatable benchmark in this repo, so rather than quote invented numbers, here are the cost characteristics and where to read the real signal live: -- **Ingestion**: TCP JSON decode → (optional per-series/priority admission, ADR-027) → - bounded block-then-shed queue (ADR-023) → WAL - append (group-commit fsync, ADR-026 — see below) → in-memory head append (with an - inverted-index update). The live - rate is on `/api/v1/stats` (`ingestion_rate`, a windowed samples/sec — ADR-017) and - the cumulative `meridian_samples_ingested_total` counter; `ingest_queue_depth`/ - `_capacity` and `meridian_dropped_samples_total` expose backpressure. The default - simulator (8 hosts × ~43 series at a 5 s cadence) is a light, steady load. -- **Admission shaping** (when enabled, ADR-027): one classify (a short scan over the - configured classes) plus an allocation-free order-independent hash of the series - identity per offered series, and — only above the contention threshold — one - sharded token-bucket consult. State is O(shards + metric-buckets), fixed at - construction, so cost and memory are independent of series cardinality. Off by - default, it adds nothing to the hot path; the selectivity it buys is visible in the - `meridian_admission_*` counters. -- **Query**: cost scales with the number of series matched (inverted-index - intersection) × steps — each leaf selector is fetched once over the whole range and - sliced per step (ADR-014) — plus the block scan/merge. Per-query latency is recorded - in the `meridian_query_latency_seconds` histogram and shown on the dashboard's latency - panel. -- **Memory**: dominated by the in-memory head (per-series label set + sample buffer) and - the inverted index. It is bounded over time by flush-to-block plus retention, and - bounded under overload by the ingest queue capacity (depth ≤ capacity — ADR-023). - Persisted blocks store ~4.5 bits/sample for regular gauges (see above). -- **Hinted handoff** (ADR-029): off the live write path. A normal write adds, per missed - replica, one ring `PreferenceList` walk plus one durable hint file (fsync + rename) on - the ingestor — and only while a replica is actually down. Replay and backfill are a - background recovery path (one target at a time, FIFO), not the hot path. The buffer is - bounded per target by `max_samples_per_node` (drop-oldest past it), so a long outage - caps hint disk/memory rather than growing without bound; `meridian_handoff_pending_*` - shows the live backlog and `_replayed_/_dropped_samples_total` the catch-up progress. -- **Anti-entropy** (ADR-030): a background sweep, never on the read or write path, and - bounded on both axes. *Spatially*, `groups_per_round` caps how many replica groups a - round touches (the round-robin cursor covers the rest over later rounds) and the number - of groups tracks the cluster's distinct replica sets, not the virtual-node count. - *Temporally*, a round's per-group cost is one match-all read over `[now-lookback, now]` - per replica to compute the digest; the agreement case ends at a single root comparison - and transfers nothing, and only a divergent `window` is re-read and gap-filled. Smaller - `window` re-transfers less per divergence but enlarges the digest; `lookback` bounds the - per-round read on large datasets (`0` re-digests all history). `interval`+`jitter` set - the cadence and de-sync coordinators. Progress shows in - `meridian_anti_entropy_repairs_total` / `_transferred_samples_total`; `_divergent_windows_total` - climbing while `_repairs_total` does not flags an unresolvable difference (a same-timestamp - value conflict, which gap-fill does not overwrite). -- **Rebalancing on membership change** (ADR-031): off the live read and write path — it runs - only when a node is explicitly joined or left, and the work is proportional to the data that - actually moved, not the dataset. The owner-set diff is `O(vnodes × nodes)` over two ring - snapshots (no I/O); migration cost is one range read from a current owner plus one backfill - push per new owner, per moved arc-group, processed **sequentially** (no thundering herd) with - an optional `max_bytes_per_round` to spread a large move across passes. GC is a per-node drop - of the shed arcs: the head is flushed once, then only the blocks holding un-owned series are - rewritten (decode + re-encode of the kept series) — fully-owned blocks are untouched and - fully-un-owned blocks are deleted, so a node that loses a fraction of the keyspace pays for - rewriting only the blocks that mix owned and un-owned series. A joining node stays out of - routing until its data has arrived, so the migration adds no read-path latency; reads stay - complete throughout because the old owners keep their copy until the new owners are confirmed - at quorum. `meridian_rebalance_*` shows migrations/bytes moved and GC series/samples - reclaimed; `_skipped_total` climbing without `_migrations_total` flags a move that cannot - reach a source or a quorum. -- **Anomaly detection** (ADR-024, ADR-028): on the broadcast tick, not the read or write - path — one map lookup and a handful of float ops per live series per tick (~1 Hz), under - a single lock. **EWMA** (default) keeps O(1) state per series (a few `float64`s). - **Holt-Winters** (`mode: holt_winters`) keeps O(`season_length`) per series (the seasonal - array, plus a one-season warmup accumulator that is released once seeded); its per-tick - work is still a constant handful of ops (one bucket index from the timestamp, a forecast, - three smoothing updates). Memory follows live cardinality because unseen series are - evicted. Activity shows in `meridian_anomalies_total` / `meridian_active_anomalies`, and +- **Ingestion** — TCP JSON decode → optional admission (ADR-027) → bounded block-then-shed + queue (ADR-023) → WAL append (group-commit fsync, ADR-026) → in-memory head append (with + an inverted-index update). Live signal: `ingestion_rate` (windowed samples/sec, ADR-017) + and the cumulative `meridian_samples_ingested_total` on `/api/v1/stats`; + `ingest_queue_depth`/`_capacity` and `meridian_dropped_samples_total` expose backpressure. + The default simulator (8 hosts × ~43 series at a 5 s cadence) is a light, steady load. +- **Admission shaping** (when enabled, ADR-027) — per offered series: one classify (a short + scan over the configured classes) plus an allocation-free order-independent identity hash, + and — only above the contention threshold — one sharded token-bucket consult. State is + O(shards + metric-buckets), fixed at construction, so cost and memory are independent of + cardinality. Off by default it adds nothing to the hot path; the selectivity shows in + `meridian_admission_*`. +- **Query** — cost scales with matched series (inverted-index intersection) × steps — each + leaf selector is fetched once over the whole range and sliced per step (ADR-014) — plus + the block scan/merge. Per-query latency is in the `meridian_query_latency_seconds` + histogram and on the dashboard's latency panel. +- **Memory** — dominated by the in-memory head (per-series label set + sample buffer) and + the inverted index. Bounded over time by flush-to-block plus retention, and under overload + by the ingest queue capacity (depth ≤ capacity, ADR-023). Persisted blocks store ~4.5 + bits/sample for regular gauges. +- **Hinted handoff** (ADR-029) — off the live write path. A normal write adds, per missed + replica and only while it is down, one ring `PreferenceList` walk plus one durable hint + file (fsync + rename) on the ingestor. Replay and backfill are a background recovery path + (one target at a time, FIFO). The buffer is bounded per target by `max_samples_per_node` + (drop-oldest past it), so a long outage caps hint disk/memory; `meridian_handoff_pending_*` + shows the backlog and `_replayed_`/`_dropped_samples_total` the catch-up progress. +- **Anti-entropy** (ADR-030) — a background sweep, never on the read or write path, bounded + on both axes: + - *Spatially* — `groups_per_round` caps how many replica groups a round touches (the + round-robin cursor covers the rest later), and the group count tracks the cluster's + distinct replica sets, not the virtual-node count. + - *Temporally* — a round's per-group cost is one match-all read over `[now-lookback, now]` + per replica for the digest; agreement ends at a single root comparison and transfers + nothing, and only a divergent `window` is re-read and gap-filled. + - Tuning: smaller `window` re-transfers less per divergence but enlarges the digest; + `lookback` bounds the per-round read (`0` re-digests all history); `interval`+`jitter` + set the cadence and de-sync coordinators. + - `_repairs_total`/`_transferred_samples_total` show progress; `_divergent_windows_total` + climbing while `_repairs_total` does not flags an unresolvable same-timestamp value + conflict (gap-fill does not overwrite). +- **Rebalancing on membership change** (ADR-031) — off the live read and write path; runs + only on an explicit join/leave, and the work is proportional to the data that moved, not + the dataset: + - The owner-set diff is `O(vnodes × nodes)` over two ring snapshots (no I/O); migration is + one range read from a current owner plus one backfill push per new owner, per moved + arc-group, processed **sequentially** (no thundering herd) with an optional + `max_bytes_per_round`. + - GC drops the shed arcs per node: the head is flushed once, then only blocks mixing owned + and un-owned series are rewritten (decode + re-encode of the kept series) — fully-owned + blocks untouched, fully-un-owned deleted. + - A joining node stays out of routing until its data arrives, so migration adds no + read-path latency; reads stay complete because old owners keep their copy until the new + owners are confirmed at quorum. + - `meridian_rebalance_*` shows migrations/bytes and GC series/samples; `_skipped_total` + climbing without `_migrations_total` flags a move that cannot reach a source or quorum. +- **Anomaly detection** (ADR-024, ADR-028) — on the broadcast tick, not the read or write + path: one map lookup and a handful of float ops per live series per tick (~1 Hz), under a + single lock. **EWMA** (default) keeps O(1) state per series; **Holt-Winters** + (`mode: holt_winters`) keeps O(`season_length`) (the seasonal array, plus a one-season + warmup accumulator released once seeded) with the same constant per-tick work (one bucket + index, a forecast, three smoothing updates). Memory follows live cardinality (unseen + series evicted). `meridian_anomalies_total`/`meridian_active_anomalies` show activity; `meridian_anomaly_model_info` names the active model. diff --git a/PROTOCOL.md b/PROTOCOL.md index a3ccc50..b8a8afc 100644 --- a/PROTOCOL.md +++ b/PROTOCOL.md @@ -67,19 +67,17 @@ GET /api/query?q=&start=&end=&format= } ``` -**Resolution selection** (ADR-011, ADR-025): in the single-binary `serve`, the planner -transparently serves wide spans from a coarse rollup tier instead of raw, picking the -resolution from the query span and step. The response carries `resolution_ms` (the -rollup window served, in ms; `0` = raw) and `points_read` (points fetched from storage) -so a caller can see that a wide query read far fewer points from a coarse tier. A wide -range query is also served coarse when its function maps to a stored column — the -`*_over_time` family reads the matching aggregate (`max_over_time`→max, …) and `rate()` -reads the counter-increase column as `Σincrease / range` (window-averaged). Short ranges, -`last_over_time`, and a bare range selector still read raw. This selection runs in the -docker-compose **cluster** too: the querier runs the same planner and requests the chosen -resolution + aggregate column from the storage nodes (see Internal Cluster API below), so a -wide cluster query reports a non-zero `resolution_ms` exactly like the single binary -(ADR-011). +**Resolution selection** (ADR-011, ADR-025): the planner transparently serves wide spans +from a coarse rollup tier instead of raw, picking the resolution from the span and step. The +response carries `resolution_ms` (the rollup window served in ms; `0` = raw) and `points_read` +so a caller can see a wide query read far fewer points. A range query is served coarse when +its function maps to a stored column — the `*_over_time` family reads the matching aggregate +(`max_over_time`→max, …) and `rate()` reads the counter-increase column as `Σincrease / range` +(window-averaged); short ranges, `last_over_time`, and a bare range selector still read raw. +This runs in the single binary and the docker-compose **cluster** alike — the querier runs the +same planner and requests the chosen resolution + aggregate column from the storage nodes (see +Internal Cluster API below), so a wide cluster query reports a non-zero `resolution_ms` exactly +like the single binary. ### Labels @@ -168,22 +166,20 @@ resolution every live node can serve — the intersection): POST /api/internal/backfill ``` -The catch-up path for hinted handoff (ADR-029). The ingestor replays the hints it -buffered for a replica that was down through this endpoint when the replica returns. The -request body is identical to a storage write (`{"time_series":[...]}`); the difference is -the apply semantics: backfill accepts samples **older** than a series' last (inserting -them in sorted position, filling only gaps) where `/api/internal/write` would reject them -as out-of-order (ADR-015). This is what fills an interior gap read-repair cannot. The -response reports how many samples were applied (an exact duplicate of an existing point is -skipped): +The catch-up path for hinted handoff (ADR-029): the ingestor replays hints buffered for a +down replica through this endpoint when it returns. The body is identical to a storage write +(`{"time_series":[...]}`); the difference is apply semantics — backfill accepts samples +**older** than a series' last (inserting them in sorted position, filling only gaps) where +`/api/internal/write` would reject them as out-of-order (ADR-015). This is what fills an +interior gap read-repair cannot. The response reports how many samples were applied (an exact +duplicate is skipped): ```json { "samples_ingested": 128 } ``` -Backfilled samples are logged under a distinct WAL frame, so they survive a storage -restart and replay through the same out-of-order-tolerant path; only a recovering replica -uses this endpoint, so the live in-order write path is unaffected. +Backfilled samples are logged under a distinct WAL frame, so they survive a storage restart; +only a recovering replica uses this endpoint, so the live in-order write path is unaffected. ### Anti-entropy digest & range