fix: audit batch 2 — semi-sync, DISTINCT, set ops, MERGE, storage hardening
CI / test (push) Has been cancelled
CI / raft-e2e (push) Has been cancelled
CI / verify (push) Has been cancelled
Clients CI / build-server (push) Has been cancelled
Clients CI / test-python (push) Has been cancelled
Clients CI / test-javascript (push) Has been cancelled
Clients CI / test-nim (push) Has been cancelled
Clients CI / test-rust (push) Has been cancelled

Address 12 deep-audit findings: semi-sync fail-closed on partial ack,
COUNT/SUM/AVG(DISTINCT), UNION/INTERSECT/EXCEPT dedup, MERGE THEN DELETE,
WAL torn-record recovery, MVCC/checkpoint/flush/compaction/WAL rewrite
safety, mmap overflow bounds; remove stray protocol/scram ELF.
This commit is contained in:
2026-08-02 23:12:47 +03:00
parent ccc54e8f18
commit e44341e47c
17 changed files with 384 additions and 106 deletions
+1
View File
@@ -71,5 +71,6 @@ src/barabadb/storage/lsm
src/barabadb/storage/wal src/barabadb/storage/wal
src/barabadb/storage/btree src/barabadb/storage/btree
src/barabadb/storage/gate src/barabadb/storage/gate
src/barabadb/protocol/scram
clients/nim/tests/test_pool clients/nim/tests/test_pool
clients/nim/tests/test_wire clients/nim/tests/test_wire
+26 -24
View File
@@ -3,7 +3,7 @@
> Дата: 2026-08-02 > Дата: 2026-08-02
> Метод: 4 паралелни одит-агента по слоеве (Storage / Query / Core / Protocol), всеки чете всички файлове в обхвата си и проверява находките срещу реалния код. > Метод: 4 паралелни одит-агента по слоеве (Storage / Query / Core / Protocol), всеки чете всички файлове в обхвата си и проверява находките срещу реалния код.
> Обхват: **само нови дефекти** — 80-те вече оправени в `BUGS.md` / `BUG_AUDIT.md` / `BARADB_CLIENT_BUGS.md` са изключени. > Обхват: **само нови дефекти** — 80-те вече оправени в `BUGS.md` / `BUG_AUDIT.md` / `BARADB_CLIENT_BUGS.md` са изключени.
> **Общо: ~28 находки | Поправени (батч 1): 5 | Остават: 23** > **Общо: ~28 находки | Поправени: 17 (батч 1: 5 + батч 2: 12, вкл. hygiene) | Остават: 12**
--- ---
@@ -17,54 +17,56 @@
| H4 | 🟠 HIGH | **`**` и `++` се lower-ваха към equality** — `bkPow`/`bkConcat` липсваха в op-mapping case-а и попадаха в `else: irEq` (`2 ** 3``false`, `'a' ++ 'b'``false`). | `query/exec/lower.nim:79` | `of bkPow: irOp = irPow`, `of bkConcat: irOp = irAdd`; 2 регресионни теста | | H4 | 🟠 HIGH | **`**` и `++` се lower-ваха към equality** — `bkPow`/`bkConcat` липсваха в op-mapping case-а и попадаха в `else: irEq` (`2 ** 3``false`, `'a' ++ 'b'``false`). | `query/exec/lower.nim:79` | `of bkPow: irOp = irPow`, `of bkConcat: irOp = irAdd`; 2 регресионни теста |
| H5 | 🟠 HIGH | **`!=` не е отрицание на `=`** — `irNeq` short-circuit-ваше на string inequality, така че `5 != 5.0` → true, но `5 = 5.0` → true. | `query/exec/eval.nim:438` | `irNeq` numeric-first (точно допълнение на `irEq`); регресионен тест | | H5 | 🟠 HIGH | **`!=` не е отрицание на `=`** — `irNeq` short-circuit-ваше на string inequality, така че `5 != 5.0` → true, но `5 = 5.0` → true. | `query/exec/eval.nim:438` | `irNeq` numeric-first (точно допълнение на `irEq`); регресионен тест |
**Верификация:** `baradadb` build чист; `tests/test_all.nim` (пълен suite) и `tests/bugfix_test.nim` минават без `[FAILED]`. ## Поправени — батч 2 (12)
| # | Severity | Проблем | Файл | Fix |
|---|----------|---------|------|-----|
| H3 | 🟠 HIGH | **Semi-sync partial/zero ack** — връщаше LSN дори при 0 acks | `core/replication.nim` | `return 0` когато connected replicas < `syncReplicaCount` acks; 0 connected → local-only (като sync) |
| H6 | 🟠 HIGH | **`COUNT/SUM/AVG(DISTINCT)` игнорира DISTINCT** | `query/exec/lower.nim`, `plan_exec.nim` | `aggDistinct = node.funcDistinct`; dedup с `HashSet` в agg пътищата |
| H7 | 🟠 HIGH | **`UNION/INTERSECT/EXCEPT` KeyError** | `query/executor.nim` | Dedup fingerprint от projected cols, не `row["$value"]` |
| H8 | 🟠 HIGH | **`MERGE … THEN DELETE` / matched condition no-op** | `query/executor.nim` | Honor `mergeMatchedDelete` + `mergeMatchedCondition` |
| H9 | 🟠 HIGH | **WAL recovery crash на torn record** | `storage/lsm.nim`, `wal.nim`, `recovery.nim` | Bound key/val ≤ 64 MB; validate kind преди enum cast |
| M1 | 🟡 MEDIUM | **MVCC `write` delete-during-iteration** | `core/mvcc.nim` | Collect-then-delete stale txn ids |
| M3 | 🟡 MEDIUM | **`checkpoint` lock leak** | `storage/lsm.nim` | try/finally около write lock + walLock |
| M4 | 🟡 MEDIUM | **`flushUnsafe` clear-before-write** | `storage/lsm.nim` | Clear memtable едва след успешен `writeSSTable` |
| M5 | 🟡 MEDIUM | **Compaction empty-key skip** | `storage/compaction.nim` | `haveLast` флаг вместо `lastKey = ""` sentinel |
| M6 | 🟡 MEDIUM | **`rewriteLive` remove-before-move** | `storage/wal.nim` | Само атомен `moveFile` (rename replace) |
| L3 | 🟢 LOW | **mmap `offset+size` overflow** | `storage/mmap.nim` | Overflow-safe: `offset > size - length` |
| — | hygiene | **Stray ELF `protocol/scram`** | `.gitignore` | Премахнат binary + ignore entry |
**Верификация (батч 2):** `baradadb` build чист; `tests/bugfix_test.nim` (вкл. batch-2 suite) и `tests/test_all.nim` (501 OK) минават без `[FAILED]`. `tests/prop_test.nim` B-Tree suite OK (H10 *не* е в този батч — naive left-max fix чупи interleaved remove).
--- ---
## Остават (23) ## Остават (12)
### 🟠 HIGH (8) ### 🟠 HIGH (2)
| # | Проблем | Файл | Предложен fix | | # | Проблем | Файл | Предложен fix |
|---|---------|------|---------------| |---|---------|------|---------------|
| H2 | **TLS client връзките между възли не верифицират сертификата**`forwardQueryToLeader` ползва `verifyMode = CVerifyNone` → MITM на клъстър линка. Raft client dials са със същия default (`raftTlsVerifyPeer: false`). | `core/server.nim:70` | Verify peer cert срещу CA при client handshake (fail-closed при enabled TLS) | | H2 | **TLS client връзките между възли не верифицират сертификата**`forwardQueryToLeader` ползва `verifyMode = CVerifyNone` → MITM на клъстър линка. Raft client dials са със същия default (`raftTlsVerifyPeer: false`). | `core/server.nim:70` | Verify peer cert срещу CA при client handshake (fail-closed при enabled TLS) |
| H3 | **Semi-sync репликация връща durable LSN при partial/zero ack**`rmSync` връща 0 при partial ack, но `rmSemiSync` връща LSN безусловно (само debug echo при 0 acks). | `core/replication.nim:217` | `if syncReplicaCount > 0 and ackCount < syncReplicaCount: return 0` | | H10 | **B-tree `remove` separator convention** — audit: `splitChild` left-max vs `removeRec` right-min. Naive left-max rewrite of separators/borrows **fails** `prop_test` interleaved insert/remove; needs careful multi-level fix + more targeted repro first. | `storage/btree.nim:377` | Repro + full-tree separator invariant; keep borrow/merge/search consistent |
| H6 | **`COUNT/SUM/AVG(DISTINCT)` игнорира DISTINCT** — `funcDistinct` се set-ва в парсера (BUG-017), но `aggDistinct` никога не се копира/чете; няма dedup в aggregate пътищата. | `query/exec/lower.nim:116`, `plan_exec.nim` | Копирай `aggDistinct = node.funcDistinct`; dedup чрез `HashSet[string]` преди count/sum/avg |
| H7 | **`UNION/INTERSECT/EXCEPT` (без ALL) чупят с KeyError** — dedup ключът чете `row["$value"]`, но projected редове нямат този ключ. Само `UNION ALL` работи. | `query/executor.nim:459` | Dedup ключ от projected колоните (join `valueToString` по ред на `cols`), не `row["$value"]` |
| H8 | **`MERGE ... WHEN MATCHED THEN DELETE` / `AND <cond>` не се изпълняват** — AST/parser полетата (BUG-032) съществуват, но executor-ът не ги реферира; DELETE е no-op, condition се игнорира. | `query/executor.nim:752` | В matched клона: провери `mergeMatchedCondition`, после honor-вай `mergeMatchedDelete` |
| H9 | **WAL recovery чупи процеса при torn record** — recovery parser-ът вярва на `keyLen`/`valLen` (до ~4 GiB alloc) и `kind` (out-of-range enum → `CaseStmtError` Defect, не се catch-ва). | `storage/lsm.nim:704` | Bound lengths + валидирай `kind` преди use; дългосрочно per-record CRC32 |
| H10 | **B-tree `remove` пише separator с грешна конвенция**`splitChild` ползва left child max, `removeRec` пише right child min (`child.keys[0]`) → ключове стават ненамираеми при internal nodes (silent data loss). | `storage/btree.nim:377` | Separator = max ключ на left child, преизчислен след rebalance |
### 🟡 MEDIUM (11) ### 🟡 MEDIUM (7)
| # | Проблем | Файл | Предложен fix | | # | Проблем | Файл | Предложен fix |
|---|---------|------|---------------| |---|---------|------|---------------|
| M1 | **MVCC `write` трие от `activeTxns` по време на итерация**`delete` proc-ът ползва collect-then-delete, но `write` трие inline (unsafe, пропуска timed-out txns). | `core/mvcc.nim:180` | Collect stale ids в seq, трий след loop-а |
| M2 | **disttxn `connectWithTimeout` без SO_ERROR + uncaught RPC** — refused connect е "writable" → връща true; `sendDistTxnRpc` няма try/except → OSError wedge-ва 2PC състояние. (BUG-042 fix-нат в replication, не тук.) | `core/disttxn.nim:88` | `getsockopt(SO_ERROR)` + try/except около per-participant RPC | | M2 | **disttxn `connectWithTimeout` без SO_ERROR + uncaught RPC** — refused connect е "writable" → връща true; `sendDistTxnRpc` няма try/except → OSError wedge-ва 2PC състояние. (BUG-042 fix-нат в replication, не тук.) | `core/disttxn.nim:88` | `getsockopt(SO_ERROR)` + try/except около per-participant RPC |
| M3 | **`checkpoint` leaking write lock при exception** — `acquireWrite` без try/finally; IOError от flush/rotate пропуска `releaseWrite` → постоянен hang. | `storage/lsm.nim:932` | try/finally около lock-а (и walLock) | | M7 | **Compaction unlink-ва input-ите преди output-ът да е loadable в каталога**verifySSTable вече е преди unlink; остава catalog re-load ordering в caller. | `storage/compaction.nim` / LSM apply | Load/verify output в каталога ПРЕДИ unlink на input-ите |
| M4 | **`flushUnsafe` празни memtable преди SSTable write** — при IOError на `writeSSTable` данните са загубени от memory (остават само в WAL, невидими за live reads). | `storage/lsm.nim:870` | Първо `writeSSTable`, после clear на memtable |
| M5 | **Compaction пропуска empty-string ключа** — dedup sentinel `lastKey = ""` skip-ва ключ `""` → data loss при compact на празен ключ. | `storage/compaction.nim:130` | `haveLast` флаг вместо sentinel стойност |
| M6 | **`rewriteLive` data-loss window** — `removeFile(wal.path)` преди `moveFile`; crash между тях губи unflushed записи. `rename(2)` и без това е атомен replace. | `storage/wal.nim:329` | Махни `removeFile`, остави атомния `moveFile` |
| M7 | **Compaction unlink-ва input-ите преди output-ът да е loadable в каталога**`applyCompactionResult` re-load-ва output в try/except (само warning); ако fail-не след unlink → загуба на ключове. | `storage/compaction.nim:150` | Load/verify output в каталога ПРЕДИ unlink на input-ите |
| M8 | **`OFFSET n` без `LIMIT` връща 0 реда; negative `LIMIT` чупи** — `limitCount = 0` е sentinel и за "няма limit", и за "LIMIT 0"; `sourceRows[start..<endIdx]` с endIdx<start → IndexDefect. | `query/exec/lower.nim:411`, `plan_exec.nim:219` | Отделен sentinel (-1 = unlimited); clamp negative | | M8 | **`OFFSET n` без `LIMIT` връща 0 реда; negative `LIMIT` чупи** — `limitCount = 0` е sentinel и за "няма limit", и за "LIMIT 0"; `sourceRows[start..<endIdx]` с endIdx<start → IndexDefect. | `query/exec/lower.nim:411`, `plan_exec.nim:219` | Отделен sentinel (-1 = unlimited); clamp negative |
| M9 | **Aggregate window функции връщат NULL**`SUM/AVG/COUNT/MIN/MAX OVER (...)` попадат в `else` клона (`"\N"`); само ranking/lead/lag се handle-ват. | `query/exec/window.nim` | `of "sum","avg","count","min","max"` с `resolveFrameBounds` | | M9 | **Aggregate window функции връщат NULL**`SUM/AVG/COUNT/MIN/MAX OVER (...)` попадат в `else` клона (`"\N"`); само ranking/lead/lag се handle-ват. | `query/exec/window.nim` | `of "sum","avg","count","min","max"` с `resolveFrameBounds` |
| M10 | **WebSocket приема unmasked client frames** — RFC 6455 §5.1 изисква server да затвори връзката при unmasked client frame (cache-poisoning защита). | `core/websocket.nim:85` | Затвори връзката при `masked == false` | | M10 | **WebSocket приема unmasked client frames** — RFC 6455 §5.1 изисква server да затвори връзката при unmasked client frame (cache-poisoning защита). | `core/websocket.nim:85` | Затвори връзката при `masked == false` |
| M11 | **WebSocket без frame/message size limit → DoS**`buf.add(chunk)` расте неограничено; няма 125-byte control-frame cap. | `core/websocket.nim:215` | Max frame/message size + control-frame cap | | M11 | **WebSocket без frame/message size limit → DoS**`buf.add(chunk)` расте неограничено; няма 125-byte control-frame cap. | `core/websocket.nim:215` | Max frame/message size + control-frame cap |
| M12 | **WebSocket SUBSCRIBE bypass-ва table-level auth** — всеки автентикиран клиент subscribe-ва към任意 таблица и получава всички insert/delete. | `core/websocket.nim:232` | Table read authorization при subscribe | | M12 | **WebSocket SUBSCRIBE bypass-ва table-level auth** — всеки автентикиран клиент subscribe-ва към произволна таблица и получава всички insert/delete. | `core/websocket.nim:232` | Table read authorization при subscribe |
### 🟢 LOW (4) ### 🟢 LOW (3)
| # | Проблем | Файл | Предложен fix | | # | Проблем | Файл | Предложен fix |
|---|---------|------|---------------| |---|---------|------|---------------|
| L1 | **SCRAM timing user enumeration** — unknown user връща веднага, known user прави urandom+HMAC/PBKDF2 работа → timing delta (BUG-049 fix-на съобщението, не timing-а). | `protocol/auth.nim:227` | Equivalent dummy work за unknown users | | L1 | **SCRAM timing user enumeration** — unknown user връща веднага, known user прави urandom+HMAC/PBKDF2 работа → timing delta (BUG-049 fix-на съобщението, не timing-а). | `protocol/auth.nim:227` | Equivalent dummy work за unknown users |
| L2 | **SCRAM channel-binding не се верифицира**`c=` се приема verbatim; RFC 5802 изисква валидация. Не е exploitable днес (няма TLS-CB). | `protocol/auth.nim:259` | Enforce expected `c=` (напр. `biws`) | | L2 | **SCRAM channel-binding не се верифицира**`c=` се приема verbatim; RFC 5802 изисква валидация. Не е exploitable днес (няма TLS-CB). | `protocol/auth.nim:259` | Enforce expected `c=` (напр. `biws`) |
| L3 | **mmap read `offset + size` overflow**`offset + size > region.size` с native int wrap-ва негативно при corrupt offset/size → OOB read. v3 SSTable-ите са CRC-защитени (reachable само през legacy v1/v2). | `storage/mmap.nim` | `offset > region.size - size` (без overflow) |
| L4 | **NULL equality semantics**`NULL = NULL` и `col = NULL` → true (string sentinel сравнение), не unknown/false. Системно за string-based value модела. | `query/exec/eval.nim:431` | Three-valued logic за NULL (по-голям рефакторинг) | | L4 | **NULL equality semantics**`NULL = NULL` и `col = NULL` → true (string sentinel сравнение), не unknown/false. Системно за string-based value модела. | `query/exec/eval.nim:431` | Three-valued logic за NULL (по-голям рефакторинг) |
### Хигиена
- **Stray 94 KB компилиран ELF binary** в `src/barabadb/protocol/scram` — случайно commit-нат в source tree-то; да се премахне (+ `.gitignore`).
--- ---
## Проверени и чисти (не са бъгове) ## Проверени и чисти (не са бъгове)
+16 -1
View File
@@ -14,13 +14,28 @@ All notable changes to BaraDB are documented in this file.
- **Raft commit quorum (CRITICAL)** — commit now requires a strict majority (`N div 2 + 1`), matching the election check; the previous `(N+1) div 2` formula committed at a minority for even-sized clusters (`core/raft.nim`) - **Raft commit quorum (CRITICAL)** — commit now requires a strict majority (`N div 2 + 1`), matching the election check; the previous `(N+1) div 2` formula committed at a minority for even-sized clusters (`core/raft.nim`)
- **`**` / `++` operators (HIGH)** — power and concat are no longer lowered to equality: `2 ** 3` → 8, `'a' ++ 'b'``'ab'` (`query/exec/lower.nim`) - **`**` / `++` operators (HIGH)** — power and concat are no longer lowered to equality: `2 ** 3` → 8, `'a' ++ 'b'``'ab'` (`query/exec/lower.nim`)
- **`!=` semantics (HIGH)** — `!=` is now the exact complement of `=` for numerically-equal values (`5 != 5.0` is false) (`query/exec/eval.nim`) - **`!=` semantics (HIGH)** — `!=` is now the exact complement of `=` for numerically-equal values (`5 != 5.0` is false) (`query/exec/eval.nim`)
- **Semi-sync partial ack (HIGH)** — `writeLsn` in `rmSemiSync` returns `0` when connected replicas fail to meet `syncReplicaCount`; zero connected peers still succeed local-only (like sync) (`core/replication.nim`)
- **`COUNT/SUM/AVG(DISTINCT …)` (HIGH)** — `funcDistinct` is copied to `aggDistinct` and applied via `HashSet` dedup in aggregate paths (`query/exec/lower.nim`, `plan_exec.nim`)
- **`UNION` / `INTERSECT` / `EXCEPT` (HIGH)** — set-op dedup fingerprints projected columns instead of missing `row["$value"]` (KeyError crash) (`query/executor.nim`)
- **`MERGE … WHEN MATCHED THEN DELETE` (HIGH)** — executor honors `mergeMatchedDelete` and optional `mergeMatchedCondition` (`query/executor.nim`)
- **WAL recovery torn records (HIGH)** — recovery bounds key/value to 64 MB and rejects out-of-range entry kinds before enum cast (avoids multi-GiB alloc / `CaseStmtError` Defect) (`storage/lsm.nim`, `wal.nim`, `recovery.nim`)
- **MVCC `write` timeout cleanup** — stale active transactions are collected then deleted (no mutation during `activeTxns` iteration) (`core/mvcc.nim`)
- **`checkpoint` lock leak** — write lock and `walLock` released in `try/finally` (`storage/lsm.nim`)
- **`flushUnsafe` data-loss window** — memtable is cleared only after a successful SSTable write (`storage/lsm.nim`)
- **Compaction empty-string key** — dedup uses a `haveLast` flag so key `""` is not skipped (`storage/compaction.nim`)
- **`rewriteLive` crash window** — atomic `moveFile` replace only (no `removeFile` before rename) (`storage/wal.nim`)
- **mmap OOB on overflow** — length checks use `offset > size - length` instead of wrapping `offset + size` (`storage/mmap.nim`)
- **REP replication put/delete encoding** — the legacy REP payload carries an explicit op tag so PK-only inserts (empty value) replicate as puts instead of vanishing as deletes (`core/replication.nim`, `core/server.nim`) - **REP replication put/delete encoding** — the legacy REP payload carries an explicit op tag so PK-only inserts (empty value) replicate as puts instead of vanishing as deletes (`core/replication.nim`, `core/server.nim`)
- **REP receiver secondary indexes** — the legacy REP receiver applies via `applyReplicatedPut/Delete` under the storage gate, keeping B-tree/FTS/HNSW/graph indexes consistent on the replica (`core/server.nim`) - **REP receiver secondary indexes** — the legacy REP receiver applies via `applyReplicatedPut/Delete` under the storage gate, keeping B-tree/FTS/HNSW/graph indexes consistent on the replica (`core/server.nim`)
- **Snapshot send stall (partial)** — the leader's snapshot send runs gzip off the event loop on a worker thread (`gzipFileAsync`), so heartbeats keep firing during compression; tar (send) and the restore path still run on the loop (`core/backup.nim`, `core/raft.nim`) - **Snapshot send stall (partial)** — the leader's snapshot send runs gzip off the event loop on a worker thread (`gzipFileAsync`), so heartbeats keep firing during compression; tar (send) and the restore path still run on the loop (`core/backup.nim`, `core/raft.nim`)
### Removed
- Stray compiled ELF `src/barabadb/protocol/scram` from the source tree (added to `.gitignore`)
### Added ### Added
- Deep audit report `BUG_AUDIT_2026-08.md` (~28 findings; 5 fixed in this batch, 23 tracked) - Deep audit report `BUG_AUDIT_2026-08.md` (~28 findings; 17 fixed across batches 12, ~12 tracked)
--- ---
+4 -2
View File
@@ -163,7 +163,9 @@
**Батч 1 — поправени (5):** MIGRATE auth bypass (CRITICAL), raft commit strict-majority за even-N (CRITICAL), pre-auth wire-length DoS (HIGH), `**`/`++` lowering към equality (HIGH), `!=` не е отрицание на `=` (HIGH). **Батч 1 — поправени (5):** MIGRATE auth bypass (CRITICAL), raft commit strict-majority за even-N (CRITICAL), pre-auth wire-length DoS (HIGH), `**`/`++` lowering към equality (HIGH), `!=` не е отрицание на `=` (HIGH).
**Остават (~23):** вж. `BUG_AUDIT_2026-08.md` — TLS peer verify, semi-sync partial-ack, COUNT(DISTINCT), UNION/INTERSECT/EXCEPT crash, MERGE THEN DELETE, WAL recovery crash, B-tree separator convention, MVCC delete-during-iteration, disttxn SO_ERROR, checkpoint lock leak, flushUnsafe data-loss, compaction empty-key, wal rewriteLive window, OFFSET-без-LIMIT, window агрегати, WebSocket (3), SCRAM (2), mmap overflow. **Батч 2 — поправени (12):** semi-sync partial-ack (H3), COUNT/SUM/AVG(DISTINCT) (H6), UNION/INTERSECT/EXCEPT (H7), MERGE THEN DELETE (H8), WAL torn-record recovery (H9), MVCC write iteration (M1), checkpoint lock leak (M3), flushUnsafe order (M4), compaction empty-key (M5), rewriteLive atomic replace (M6), mmap overflow (L3), stray `protocol/scram` ELF.
**Остават (~12):** вж. `BUG_AUDIT_2026-08.md` — TLS peer verify (H2), B-tree separator (H10, needs careful repro), disttxn SO_ERROR (M2), compaction catalog order (M7), OFFSET-без-LIMIT (M8), window агрегати (M9), WebSocket (M10M12), SCRAM (L1L2), NULL equality (L4).
--- ---
@@ -179,7 +181,7 @@
| **Този план** — Сесии 10, 11, 12 | ✅ Завършен | | **Този план** — Сесии 10, 11, 12 | ✅ Завършен |
| Raft C3a/C3b + DDL/forward/compact/metrics (2026-07-30) | ✅ Завършен на `main``docs/superpowers/specs/2026-07-30-raft-cluster-status.md` | | Raft C3a/C3b + DDL/forward/compact/metrics (2026-07-30) | ✅ Завършен на `main``docs/superpowers/specs/2026-07-30-raft-cluster-status.md` |
| **Production GA v1.2.0** (single-node) | ✅ `docs/superpowers/plans/2026-07-30-production-ga.md` | | **Production GA v1.2.0** (single-node) | ✅ `docs/superpowers/plans/2026-07-30-production-ga.md` |
| **Сесия 13** — Stabilization & Deep Audit (2026-08) | 🔄 В процес — батч 1 завършен (5 поправки); `BUG_AUDIT_2026-08.md` | | **Сесия 13** — Stabilization & Deep Audit (2026-08) | 🔄 В процес — батч 1+2 (17 поправки); остават ~12; `BUG_AUDIT_2026-08.md` |
--- ---
+5 -1
View File
@@ -178,12 +178,16 @@ proc write*(tm: TxnManager, txn: Transaction, key: string, value: seq[byte]): bo
return false return false
# Timeout-based deadlock detection: abort stale transactions # Timeout-based deadlock detection: abort stale transactions
# Collect then delete — never mutate activeTxns while iterating it.
let now = getMonoTime().ticks() let now = getMonoTime().ticks()
var staleIds: seq[TxnId] = @[]
for otherId, otherTxn in tm.activeTxns: for otherId, otherTxn in tm.activeTxns:
if otherId != txn.id and otherTxn.state == tsActive: if otherId != txn.id and otherTxn.state == tsActive:
if now - otherTxn.startTime > tm.txnTimeoutMs * 1_000_000: if now - otherTxn.startTime > tm.txnTimeoutMs * 1_000_000:
otherTxn.state = tsAborted otherTxn.state = tsAborted
tm.activeTxns.del(otherId) staleIds.add(otherId)
for id in staleIds:
tm.activeTxns.del(id)
# Check for write-write conflict against other active transactions' write sets # Check for write-write conflict against other active transactions' write sets
for otherId, otherTxn in tm.activeTxns: for otherId, otherTxn in tm.activeTxns:
+11 -3
View File
@@ -213,10 +213,18 @@ proc writeLsn*(rm: ReplicationManager, data: seq[byte]): uint64 =
rm.pendingAcks[lsn].excl(id) rm.pendingAcks[lsn].excl(id)
if rm.pendingAcks[lsn].len == 0: if rm.pendingAcks[lsn].len == 0:
rm.pendingAcks.del(lsn) rm.pendingAcks.del(lsn)
# Semi-sync requires at least syncReplicaCount acks when replicas are
# connected. With zero connected peers (nothing to ship) the write is
# local-only — same as sync mode with an empty replica set.
if rm.syncReplicaCount > 0 and replicasToShip.len > 0 and
ackCount < rm.syncReplicaCount:
# Drop the LSN from pendingAcks — write is not durable
rm.pendingAcks.del(lsn)
release(rm.lock)
echo "[ERROR] Semi-sync replication failed: only ", ackCount, "/",
rm.syncReplicaCount, " replicas acked for LSN ", lsn
return 0
release(rm.lock) release(rm.lock)
if replicasToShip.len > 0 and ackCount == 0 and rm.syncReplicaCount > 0:
when defined(debug):
echo "Replication semi-sync: no replicas acked for LSN ", lsn
return lsn return lsn
proc ackLsn*(rm: ReplicationManager, replicaId: string, lsn: uint64) = proc ackLsn*(rm: ReplicationManager, replicaId: string, lsn: uint64) =
Binary file not shown.
+1
View File
@@ -122,6 +122,7 @@ proc lowerExpr*(node: Node): IRExpr =
else: discard else: discard
result.aggArgs = @[] result.aggArgs = @[]
for arg in node.funcArgs: result.aggArgs.add(lowerExpr(arg)) for arg in node.funcArgs: result.aggArgs.add(lowerExpr(arg))
result.aggDistinct = node.funcDistinct
if node.funcFilter != nil: if node.funcFilter != nil:
result.aggFilter = lowerExpr(node.funcFilter) result.aggFilter = lowerExpr(node.funcFilter)
else: else:
+72 -14
View File
@@ -5,6 +5,7 @@
## executor split). Pure code motion — no behavior changes. ## executor split). Pure code motion — no behavior changes.
import std/strutils import std/strutils
import std/tables import std/tables
import std/sets
import std/sequtils import std/sequtils
import std/algorithm import std/algorithm
import ../ir import ../ir
@@ -19,6 +20,19 @@ import eval
import scan import scan
import window import window
# ----------------------------------------------------------------------
# Aggregate DISTINCT helpers
# ----------------------------------------------------------------------
proc shouldKeepDistinct(seen: var HashSet[string], s: string, doDistinct: bool): bool =
## Returns true if `s` should be counted/included (first occurrence when distinct).
if not doDistinct:
return true
if s in seen:
return false
seen.incl(s)
return true
# ---------------------------------------------------------------------- # ----------------------------------------------------------------------
# IR Plan Execution (with actual filter/sort/projection) # IR Plan Execution (with actual filter/sort/projection)
# ---------------------------------------------------------------------- # ----------------------------------------------------------------------
@@ -113,49 +127,71 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] =
newRow[alias] = $filteredRows.len newRow[alias] = $filteredRows.len
else: else:
var count = 0 var count = 0
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(expr.aggArgs[0], row, ctx) let v = evalExpr(expr.aggArgs[0], row, ctx)
if v.kind != vkNull: count += 1 if v.kind != vkNull:
let s = valueToString(v)
if shouldKeepDistinct(seen, s, expr.aggDistinct):
count += 1
newRow[alias] = $count newRow[alias] = $count
of irSum: of irSum:
var sum = 0.0 var sum = 0.0
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(expr.aggArgs[0], row, ctx) let v = evalExpr(expr.aggArgs[0], row, ctx)
try: sum += parseFloat(valueToString(v)) except CatchableError: discard let s = valueToString(v)
if shouldKeepDistinct(seen, s, expr.aggDistinct):
try: sum += parseFloat(s) except CatchableError: discard
newRow[alias] = $sum newRow[alias] = $sum
of irAvg: of irAvg:
var sum = 0.0 var sum = 0.0
var count = 0 var count = 0
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(expr.aggArgs[0], row, ctx) let v = evalExpr(expr.aggArgs[0], row, ctx)
try: sum += parseFloat(valueToString(v)); count += 1 except CatchableError: discard let s = valueToString(v)
if shouldKeepDistinct(seen, s, expr.aggDistinct):
try: sum += parseFloat(s); count += 1 except CatchableError: discard
newRow[alias] = if count > 0: $(sum / float(count)) else: "0" newRow[alias] = if count > 0: $(sum / float(count)) else: "0"
of irMin: of irMin:
var minVal = "" var minVal = ""
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(expr.aggArgs[0], row, ctx) let v = evalExpr(expr.aggArgs[0], row, ctx)
if v.kind == vkNull: continue if v.kind == vkNull: continue
if minVal == "" or cmpMin(valueToString(v), minVal): minVal = valueToString(v) let s = valueToString(v)
if shouldKeepDistinct(seen, s, expr.aggDistinct):
if minVal == "" or cmpMin(s, minVal): minVal = s
newRow[alias] = minVal newRow[alias] = minVal
of irMax: of irMax:
var maxVal = "" var maxVal = ""
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(expr.aggArgs[0], row, ctx) let v = evalExpr(expr.aggArgs[0], row, ctx)
if v.kind == vkNull: continue if v.kind == vkNull: continue
if maxVal == "" or cmpMax(valueToString(v), maxVal): maxVal = valueToString(v) let s = valueToString(v)
if shouldKeepDistinct(seen, s, expr.aggDistinct):
if maxVal == "" or cmpMax(s, maxVal): maxVal = s
newRow[alias] = maxVal newRow[alias] = maxVal
of irArrayAgg: of irArrayAgg:
var arr: seq[string] var arr: seq[string]
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
if expr.aggArgs.len > 0: if expr.aggArgs.len > 0:
arr.add(valueToString(evalExpr(expr.aggArgs[0], row, ctx))) let s = valueToString(evalExpr(expr.aggArgs[0], row, ctx))
if shouldKeepDistinct(seen, s, expr.aggDistinct):
arr.add(s)
newRow[alias] = "[" & arr.join(", ") & "]" newRow[alias] = "[" & arr.join(", ") & "]"
of irStringAgg: of irStringAgg:
var parts: seq[string] var parts: seq[string]
var seen: HashSet[string]
let delim = if expr.aggArgs.len > 1: evalExpr(expr.aggArgs[1], initTable[string, Value](), ctx) else: Value(kind: vkString, strVal: ",") let delim = if expr.aggArgs.len > 1: evalExpr(expr.aggArgs[1], initTable[string, Value](), ctx) else: Value(kind: vkString, strVal: ",")
for row in filteredRows: for row in filteredRows:
if expr.aggArgs.len > 0: if expr.aggArgs.len > 0:
parts.add(valueToString(evalExpr(expr.aggArgs[0], row, ctx))) let s = valueToString(evalExpr(expr.aggArgs[0], row, ctx))
if shouldKeepDistinct(seen, s, expr.aggDistinct):
parts.add(s)
newRow[alias] = parts.join(valueToString(delim)) newRow[alias] = parts.join(valueToString(delim))
else: else:
let val = evalExpr(expr, if sourceRows.len > 0: sourceRows[0] else: initTable[string, Value](), ctx) let val = evalExpr(expr, if sourceRows.len > 0: sourceRows[0] else: initTable[string, Value](), ctx)
@@ -292,49 +328,71 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] =
aggRow[aggKey] = $filteredRows.len aggRow[aggKey] = $filteredRows.len
else: else:
var count = 0 var count = 0
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(aggExpr.aggArgs[0], row, ctx) let v = evalExpr(aggExpr.aggArgs[0], row, ctx)
if v.kind != vkNull: count += 1 if v.kind != vkNull:
let s = valueToString(v)
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
count += 1
aggRow[aggKey] = $count aggRow[aggKey] = $count
of irSum: of irSum:
var sum = 0.0 var sum = 0.0
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(aggExpr.aggArgs[0], row, ctx) let v = evalExpr(aggExpr.aggArgs[0], row, ctx)
try: sum += parseFloat(valueToString(v)) except CatchableError: discard let s = valueToString(v)
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
try: sum += parseFloat(s) except CatchableError: discard
aggRow[aggKey] = $sum aggRow[aggKey] = $sum
of irAvg: of irAvg:
var sum = 0.0 var sum = 0.0
var count = 0 var count = 0
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(aggExpr.aggArgs[0], row, ctx) let v = evalExpr(aggExpr.aggArgs[0], row, ctx)
try: sum += parseFloat(valueToString(v)); count += 1 except CatchableError: discard let s = valueToString(v)
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
try: sum += parseFloat(s); count += 1 except CatchableError: discard
aggRow[aggKey] = if count > 0: $(sum / float(count)) else: "0" aggRow[aggKey] = if count > 0: $(sum / float(count)) else: "0"
of irMin: of irMin:
var minVal = "" var minVal = ""
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(aggExpr.aggArgs[0], row, ctx) let v = evalExpr(aggExpr.aggArgs[0], row, ctx)
if v.kind == vkNull: continue if v.kind == vkNull: continue
if minVal == "" or cmpMin(valueToString(v), minVal): minVal = valueToString(v) let s = valueToString(v)
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
if minVal == "" or cmpMin(s, minVal): minVal = s
aggRow[aggKey] = minVal aggRow[aggKey] = minVal
of irMax: of irMax:
var maxVal = "" var maxVal = ""
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
let v = evalExpr(aggExpr.aggArgs[0], row, ctx) let v = evalExpr(aggExpr.aggArgs[0], row, ctx)
if v.kind == vkNull: continue if v.kind == vkNull: continue
if maxVal == "" or cmpMax(valueToString(v), maxVal): maxVal = valueToString(v) let s = valueToString(v)
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
if maxVal == "" or cmpMax(s, maxVal): maxVal = s
aggRow[aggKey] = maxVal aggRow[aggKey] = maxVal
of irArrayAgg: of irArrayAgg:
var arr: seq[string] var arr: seq[string]
var seen: HashSet[string]
for row in filteredRows: for row in filteredRows:
if aggExpr.aggArgs.len > 0: if aggExpr.aggArgs.len > 0:
arr.add(valueToString(evalExpr(aggExpr.aggArgs[0], row, ctx))) let s = valueToString(evalExpr(aggExpr.aggArgs[0], row, ctx))
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
arr.add(s)
aggRow[aggKey] = "[" & arr.join(", ") & "]" aggRow[aggKey] = "[" & arr.join(", ") & "]"
of irStringAgg: of irStringAgg:
var parts: seq[string] var parts: seq[string]
var seen: HashSet[string]
let delim = if aggExpr.aggArgs.len > 1: evalExpr(aggExpr.aggArgs[1], initTable[string, Value](), ctx) else: Value(kind: vkString, strVal: ",") let delim = if aggExpr.aggArgs.len > 1: evalExpr(aggExpr.aggArgs[1], initTable[string, Value](), ctx) else: Value(kind: vkString, strVal: ",")
for row in filteredRows: for row in filteredRows:
if aggExpr.aggArgs.len > 0: if aggExpr.aggArgs.len > 0:
parts.add(valueToString(evalExpr(aggExpr.aggArgs[0], row, ctx))) let s = valueToString(evalExpr(aggExpr.aggArgs[0], row, ctx))
if shouldKeepDistinct(seen, s, aggExpr.aggDistinct):
parts.add(s)
aggRow[aggKey] = parts.join(valueToString(delim)) aggRow[aggKey] = parts.join(valueToString(delim))
# Apply HAVING filter # Apply HAVING filter
if plan.groupHaving != nil: if plan.groupHaving != nil:
+44 -9
View File
@@ -448,6 +448,26 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu
if cols.len == 0: if cols.len == 0:
cols = rightRes.columns cols = rightRes.columns
# Fingerprint a projected row for set-op dedup. Prefer declared columns;
# fall back to non-system keys so UNION/INTERSECT/EXCEPT work without `$value`.
proc setOpRowKey(row: Row, colNames: seq[string]): string =
var parts: seq[string] = @[]
if colNames.len > 0:
for c in colNames:
if c in row:
parts.add(valueToString(row[c]))
else:
parts.add("")
else:
var keys: seq[string] = @[]
for k, _ in row:
if not k.startsWith("$"):
keys.add(k)
keys.sort()
for k in keys:
parts.add(k & "=" & valueToString(row[k]))
return parts.join("\x1f")
var rows: seq[Row] = @[] var rows: seq[Row] = @[]
case stmt.setOpKind case stmt.setOpKind
of sdkUnion: of sdkUnion:
@@ -460,28 +480,30 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu
# UNION: deduplicate # UNION: deduplicate
var seen: Table[string, bool] var seen: Table[string, bool]
for row in leftRes.rows: for row in leftRes.rows:
seen[valueToString(row["$value"])] = true seen[setOpRowKey(row, cols)] = true
for row in rightRes.rows: for row in rightRes.rows:
if not seen.getOrDefault(valueToString(row["$value"]), false): let k = setOpRowKey(row, cols)
seen[valueToString(row["$value"])] = true if not seen.getOrDefault(k, false):
seen[k] = true
rows.add(row) rows.add(row)
of sdkIntersect: of sdkIntersect:
var leftSet: Table[string, bool] var leftSet: Table[string, bool]
for row in leftRes.rows: for row in leftRes.rows:
leftSet[valueToString(row["$value"])] = true leftSet[setOpRowKey(row, cols)] = true
for row in rightRes.rows: for row in rightRes.rows:
if leftSet.getOrDefault(valueToString(row["$value"]), false): let k = setOpRowKey(row, cols)
if leftSet.getOrDefault(k, false):
rows.add(row) rows.add(row)
if not stmt.setOpAll: if not stmt.setOpAll:
leftSet.del(valueToString(row["$value"])) # remove to prevent duplicates for INTERSECT (not ALL) leftSet.del(k) # remove to prevent duplicates for INTERSECT (not ALL)
of sdkExcept: of sdkExcept:
var rightSet: Table[string, bool] var rightSet: Table[string, bool]
for row in rightRes.rows: for row in rightRes.rows:
rightSet[valueToString(row["$value"])] = true rightSet[setOpRowKey(row, cols)] = true
for row in leftRes.rows: for row in leftRes.rows:
if not rightSet.getOrDefault(valueToString(row["$value"]), false): if not rightSet.getOrDefault(setOpRowKey(row, cols), false):
rows.add(row) rows.add(row)
return okResult(rows, cols) return okResult(rows, cols)
@@ -759,7 +781,20 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu
let onExpr = lowerExpr(stmt.mergeOn) let onExpr = lowerExpr(stmt.mergeOn)
if valueToString(evalExpr(onExpr, rowWithTarget, ctx)) == "true": if valueToString(evalExpr(onExpr, rowWithTarget, ctx)) == "true":
matched = true matched = true
if stmt.mergeMatchedUpdate.len > 0 and "$key" in tgtRow: # Optional AND <condition> after WHEN MATCHED
var applyMatched = true
if stmt.mergeMatchedCondition != nil:
let condExpr = lowerExpr(stmt.mergeMatchedCondition)
applyMatched = valueToString(evalExpr(condExpr, rowWithTarget, ctx)) == "true"
if applyMatched and "$key" in tgtRow:
if stmt.mergeMatchedDelete:
fireTriggers(ctx, stmt.mergeTarget, "before", "delete", tgtRow)
count += execDelete(ctx, stmt.mergeTarget, valueToString(tgtRow["$key"]), kvPairs)
fireTriggers(ctx, stmt.mergeTarget, "after", "delete", tgtRow)
if ctx.onChange != nil:
ctx.onChange(ChangeEvent(table: stmt.mergeTarget, kind: ckDelete,
key: valueToString(tgtRow["$key"]), data: ""))
elif stmt.mergeMatchedUpdate.len > 0:
var updateSets = initTable[string, string]() var updateSets = initTable[string, string]()
for s in stmt.mergeMatchedUpdate: for s in stmt.mergeMatchedUpdate:
if s.kind == nkBinOp and s.binOp == bkAssign: if s.kind == nkBinOp and s.binOp == bkAssign:
+5 -2
View File
@@ -125,13 +125,16 @@ proc compact*(cs: CompactionStrategy, level: int): CompactionResult =
return cmp(b.timestamp, a.timestamp) # newest first return cmp(b.timestamp, a.timestamp) # newest first
) )
# Deduplicate: keep only the newest version of each key # Deduplicate: keep only the newest version of each key.
# Use a haveLast flag — sentinel lastKey="" would skip the empty-string key.
var merged: seq[Entry] = @[] var merged: seq[Entry] = @[]
var lastKey = "" var lastKey = ""
var haveLast = false
for entry in allEntries: for entry in allEntries:
if entry.key != lastKey: if not haveLast or entry.key != lastKey:
merged.add(entry) merged.add(entry)
lastKey = entry.key lastKey = entry.key
haveLast = true
# Keep tombstones to prevent deleted keys from resurrecting in lower levels # Keep tombstones to prevent deleted keys from resurrecting in lower levels
var final: seq[Entry] = @[] var final: seq[Entry] = @[]
+35 -9
View File
@@ -696,6 +696,10 @@ proc newLSMTree*(
var version: uint32 = 0 var version: uint32 = 0
if stream.readData(addr magic, 4) == 4 and magic == WALMagic: if stream.readData(addr magic, 4) == 4 and magic == WALMagic:
if stream.readData(addr version, 4) == 4: if stream.readData(addr version, 4) == 4:
# Cap per-record sizes to avoid multi-GiB alloc on torn/corrupt WAL.
# Kind must be a known WalEntryKind value (1..4) — out-of-range casts
# raise CaseStmtError (Defect) and crash the process.
const MaxWalRecordField = 64 * 1024 * 1024 # 64 MB
while not stream.atEnd(): while not stream.atEnd():
var kind: uint8 = 0 var kind: uint8 = 0
var timestamp: uint64 = 0 var timestamp: uint64 = 0
@@ -704,10 +708,20 @@ proc newLSMTree*(
if stream.readData(addr kind, 1) != 1: break if stream.readData(addr kind, 1) != 1: break
if stream.readData(addr timestamp, 8) != 8: break if stream.readData(addr timestamp, 8) != 8: break
if stream.readData(addr keyLen, 4) != 4: break if stream.readData(addr keyLen, 4) != 4: break
if keyLen.int > MaxWalRecordField:
echo "[WARN] WAL recovery: torn/corrupt record (keyLen=", keyLen, ") — stopping replay"
break
# Validate kind before allocating or branching (avoids CaseStmtError Defect)
if kind < uint8(wekPut) or kind > uint8(wekCommit):
echo "[WARN] WAL recovery: invalid entry kind ", kind, " — stopping replay"
break
var key = newString(keyLen.int) var key = newString(keyLen.int)
if keyLen > 0: if keyLen > 0:
if stream.readData(addr key[0], keyLen.int) != keyLen.int: break if stream.readData(addr key[0], keyLen.int) != keyLen.int: break
if stream.readData(addr valLen, 4) != 4: break if stream.readData(addr valLen, 4) != 4: break
if valLen.int > MaxWalRecordField:
echo "[WARN] WAL recovery: torn/corrupt record (valLen=", valLen, ") — stopping replay"
break
var value = newSeq[byte](valLen.int) var value = newSeq[byte](valLen.int)
if valLen > 0: if valLen > 0:
if stream.readData(addr value[0], valLen.int) != valLen.int: break if stream.readData(addr value[0], valLen.int) != valLen.int: break
@@ -866,26 +880,36 @@ proc flushUnsafe(db: LSMTree) =
if db.immutableMem.len == 0 and db.memTable.len == 0: if db.immutableMem.len == 0 and db.memTable.len == 0:
return return
# Flush immutable memtable if present, otherwise flush current memtable # Flush immutable memtable if present, otherwise flush current memtable.
var toFlush = db.immutableMem # Do NOT clear the source memtable until the SSTable is written — an IOError
if toFlush.len == 0: # mid-write must leave the data still visible to live reads (WAL still has it).
toFlush = db.memTable var flushingImmutable = false
db.memTable = newMemTable(db.memMaxSize) var toFlush: MemTable
if db.immutableMem.len > 0:
toFlush = db.immutableMem
flushingImmutable = true
else: else:
db.immutableMem = newMemTable(0) toFlush = db.memTable
if toFlush.len == 0: if toFlush.len == 0:
return return
let path = db.dir / "sstables" / ($db.nextSSTableId & ".sst") let path = db.dir / "sstables" / ($db.nextSSTableId & ".sst")
let sstId = db.nextSSTableId
inc db.nextSSTableId inc db.nextSSTableId
# Sort once at flush time (O(n log n)) — put/get stay O(1) # Sort once at flush time (O(n log n)) — put/get stay O(1)
var sst = writeSSTable(toFlush.sortedEntries(), path, level = 0) var sst = writeSSTable(toFlush.sortedEntries(), path, level = 0)
sst.id = db.nextSSTableId - 1 sst.id = sstId
db.sstables.add(sst) db.sstables.add(sst)
# SSTables are kept in insertion order (newest last) so getUnsafe can search newest-first # SSTables are kept in insertion order (newest last) so getUnsafe can search newest-first
# Only now drop the in-memory copy — SSTable is durable on disk
if flushingImmutable:
db.immutableMem = newMemTable(0)
else:
db.memTable = newMemTable(db.memMaxSize)
# Update MANIFEST atomically # Update MANIFEST atomically
inc db.manifestSequence inc db.manifestSequence
try: try:
@@ -930,7 +954,7 @@ proc checkpoint*(db: LSMTree) =
## rotate WAL, and write MANIFEST. This provides a clean boundary ## rotate WAL, and write MANIFEST. This provides a clean boundary
## for online backup without stopping the server. ## for online backup without stopping the server.
acquireWrite(db.lock) acquireWrite(db.lock)
try:
# Flush any pending immutable memtable first # Flush any pending immutable memtable first
if db.immutableMem.len > 0: if db.immutableMem.len > 0:
flushUnsafe(db) flushUnsafe(db)
@@ -946,10 +970,12 @@ proc checkpoint*(db: LSMTree) =
# Rotate WAL for a clean backup boundary # Rotate WAL for a clean backup boundary
acquire(db.walLock) acquire(db.walLock)
try:
db.wal.maybeRotate() db.wal.maybeRotate()
db.wal.sync() db.wal.sync()
finally:
release(db.walLock) release(db.walLock)
finally:
releaseWrite(db.lock) releaseWrite(db.lock)
proc close*(db: LSMTree) = proc close*(db: LSMTree) =
+8 -4
View File
@@ -80,7 +80,8 @@ proc readAt*(mf: MmapFile, offset: int, size: int): seq[byte] =
if mf.regions.len == 0: if mf.regions.len == 0:
return @[] return @[]
let region = mf.regions[0] let region = mf.regions[0]
if offset < 0 or size < 0 or offset + size > region.size: # overflow-safe bound: offset > size - length (not offset + length > size)
if offset < 0 or size < 0 or size > region.size or offset > region.size - size:
return @[] return @[]
result = newSeq[byte](size) result = newSeq[byte](size)
copyMem(addr result[0], unsafeAddr region.data[offset], size) copyMem(addr result[0], unsafeAddr region.data[offset], size)
@@ -91,21 +92,24 @@ proc readByte*(mf: MmapFile, offset: int): byte =
return mf.regions[0].data[offset] return mf.regions[0].data[offset]
proc readUint32*(mf: MmapFile, offset: int): uint32 = proc readUint32*(mf: MmapFile, offset: int): uint32 =
if mf.regions.len == 0 or offset < 0 or offset + 4 > mf.regions[0].size: if mf.regions.len == 0 or offset < 0 or 4 > mf.regions[0].size or
offset > mf.regions[0].size - 4:
return 0 return 0
var val: uint32 var val: uint32
copyMem(addr val, unsafeAddr mf.regions[0].data[offset], 4) copyMem(addr val, unsafeAddr mf.regions[0].data[offset], 4)
return val return val
proc readUint64*(mf: MmapFile, offset: int): uint64 = proc readUint64*(mf: MmapFile, offset: int): uint64 =
if mf.regions.len == 0 or offset < 0 or offset + 8 > mf.regions[0].size: if mf.regions.len == 0 or offset < 0 or 8 > mf.regions[0].size or
offset > mf.regions[0].size - 8:
return 0 return 0
var val: uint64 var val: uint64
copyMem(addr val, unsafeAddr mf.regions[0].data[offset], 8) copyMem(addr val, unsafeAddr mf.regions[0].data[offset], 8)
return val return val
proc readString*(mf: MmapFile, offset: int, size: int): string = proc readString*(mf: MmapFile, offset: int, size: int): string =
if mf.regions.len == 0 or offset < 0 or size < 0 or offset + size > mf.regions[0].size: if mf.regions.len == 0 or offset < 0 or size < 0 or
size > mf.regions[0].size or offset > mf.regions[0].size - size:
return "" return ""
result = newString(size) result = newString(size)
copyMem(addr result[0], unsafeAddr mf.regions[0].data[offset], size) copyMem(addr result[0], unsafeAddr mf.regions[0].data[offset], size)
+4
View File
@@ -68,6 +68,7 @@ proc scanWAL*(rec: CrashRecovery): seq[RecoveredEntry] =
var txnId: uint64 = 0 var txnId: uint64 = 0
var entryCount = 0 var entryCount = 0
const MaxWalRecordField = 64 * 1024 * 1024 # 64 MB
while not stream.atEnd(): while not stream.atEnd():
var kind: uint8 = 0 var kind: uint8 = 0
var timestamp: uint64 = 0 var timestamp: uint64 = 0
@@ -77,12 +78,15 @@ proc scanWAL*(rec: CrashRecovery): seq[RecoveredEntry] =
if stream.readData(addr kind, 1) != 1: break if stream.readData(addr kind, 1) != 1: break
if stream.readData(addr timestamp, 8) != 8: break if stream.readData(addr timestamp, 8) != 8: break
if stream.readData(addr keyLen, 4) != 4: break if stream.readData(addr keyLen, 4) != 4: break
if keyLen.int > MaxWalRecordField: break
if kind < uint8(wekPut) or kind > uint8(wekCommit): break
var key = newString(keyLen.int) var key = newString(keyLen.int)
if keyLen > 0: if keyLen > 0:
if stream.readData(addr key[0], keyLen.int) != keyLen.int: break if stream.readData(addr key[0], keyLen.int) != keyLen.int: break
if stream.readData(addr valLen, 4) != 4: break if stream.readData(addr valLen, 4) != 4: break
if valLen.int > MaxWalRecordField: break
var value = newSeq[byte](valLen.int) var value = newSeq[byte](valLen.int)
if valLen > 0: if valLen > 0:
if stream.readData(addr value[0], valLen.int) != valLen.int: break if stream.readData(addr value[0], valLen.int) != valLen.int: break
+7 -2
View File
@@ -326,8 +326,9 @@ proc rewriteLive*(wal: var WriteAheadLog,
if wal.stream != nil: if wal.stream != nil:
wal.stream.close() wal.stream.close()
if fileExists(wal.path): wal.stream = nil
removeFile(wal.path) # Atomic replace: moveFile overwrites the destination on POSIX rename(2).
# Do not removeFile first — a crash between unlink and rename would lose the WAL.
moveFile(tmpPath, wal.path) moveFile(tmpPath, wal.path)
wal.stream = newFileStream(wal.path, fmAppend) wal.stream = newFileStream(wal.path, fmAppend)
if wal.stream == nil: if wal.stream == nil:
@@ -363,6 +364,7 @@ proc readEntries*(walPath: string, untilTimestamp: uint64 = 0): seq[WalEntry] =
if s.readData(addr magic, 4) != 4: return if s.readData(addr magic, 4) != 4: return
if s.readData(addr version, 4) != 4: return if s.readData(addr version, 4) != 4: return
if magic != WALMagic: return if magic != WALMagic: return
const MaxWalRecordField = 64 * 1024 * 1024 # 64 MB
while not s.atEnd: while not s.atEnd:
var kind: uint8 var kind: uint8
if s.readData(addr kind, 1) != 1: break if s.readData(addr kind, 1) != 1: break
@@ -372,11 +374,14 @@ proc readEntries*(walPath: string, untilTimestamp: uint64 = 0): seq[WalEntry] =
break break
var keyLen: uint32 var keyLen: uint32
if s.readData(addr keyLen, 4) != 4: break if s.readData(addr keyLen, 4) != 4: break
if keyLen.int > MaxWalRecordField: break
if kind < uint8(wekPut) or kind > uint8(wekCommit): break
var key = newSeq[byte](keyLen) var key = newSeq[byte](keyLen)
if keyLen > 0: if keyLen > 0:
if s.readData(addr key[0], int(keyLen)) != int(keyLen): break if s.readData(addr key[0], int(keyLen)) != int(keyLen): break
var valLen: uint32 var valLen: uint32
if s.readData(addr valLen, 4) != 4: break if s.readData(addr valLen, 4) != 4: break
if valLen.int > MaxWalRecordField: break
var value = newSeq[byte](valLen) var value = newSeq[byte](valLen)
if valLen > 0: if valLen > 0:
if s.readData(addr value[0], int(valLen)) != int(valLen): break if s.readData(addr value[0], int(valLen)) != int(valLen): break
+96
View File
@@ -625,3 +625,99 @@ suite "Query operator correctness — audit batch 1":
let neq = executeQuery(ctx, parse("SELECT * FROM users WHERE id != 1.0")) let neq = executeQuery(ctx, parse("SELECT * FROM users WHERE id != 1.0"))
check eq.rows.len == 1 # 1 = 1.0 -> true check eq.rows.len == 1 # 1 = 1.0 -> true
check neq.rows.len == 0 # 1 != 1.0 -> false (old bug returned the row) check neq.rows.len == 0 # 1 != 1.0 -> false (old bug returned the row)
suite "Query correctness — audit batch 2":
test "COUNT(DISTINCT) deduplicates values":
## Regression: funcDistinct was parsed but never copied to aggDistinct /
## never consulted during aggregation.
var ctx = setupCtx()
defer: teardown(ctx)
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (2, 'bob')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (3, 'alice')"))
let r = executeQuery(ctx, parse("SELECT COUNT(DISTINCT name) AS c FROM users"))
check r.success
check r.rows.len == 1
check valueToString(r.rows[0]["c"]) == "2"
test "SUM(DISTINCT) sums unique values only":
var ctx = setupCtx()
defer: teardown(ctx)
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'a')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (2, 'b')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (5, 'c')"))
# ids 1, 2, 5 — insert another row with id-like values via a number col
discard executeQuery(ctx, parse("CREATE TABLE nums (id INTEGER PRIMARY KEY, n INTEGER)"))
discard executeQuery(ctx, parse("INSERT INTO nums (id, n) VALUES (1, 10)"))
discard executeQuery(ctx, parse("INSERT INTO nums (id, n) VALUES (2, 10)"))
discard executeQuery(ctx, parse("INSERT INTO nums (id, n) VALUES (3, 20)"))
let r = executeQuery(ctx, parse("SELECT SUM(DISTINCT n) AS s FROM nums"))
check r.success
check r.rows.len == 1
check parseFloat(valueToString(r.rows[0]["s"])) == 30.0
test "UNION deduplicates without KeyError":
## Regression: set-op dedup used row["$value"] which projected rows lack.
var ctx = setupCtx()
defer: teardown(ctx)
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (2, 'bob')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (3, 'alice')"))
let r = executeQuery(ctx, parse(
"SELECT name FROM users WHERE id = 1 UNION SELECT name FROM users WHERE id = 3"))
check r.success
check r.rows.len == 1
check valueToString(r.rows[0]["name"]) == "alice"
test "INTERSECT returns common rows":
var ctx = setupCtx()
defer: teardown(ctx)
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (2, 'bob')"))
let r = executeQuery(ctx, parse(
"SELECT name FROM users WHERE id <= 2 INTERSECT SELECT name FROM users WHERE id = 1"))
check r.success
check r.rows.len == 1
check valueToString(r.rows[0]["name"]) == "alice"
test "EXCEPT removes right-side rows":
var ctx = setupCtx()
defer: teardown(ctx)
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')"))
discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (2, 'bob')"))
let r = executeQuery(ctx, parse(
"SELECT name FROM users EXCEPT SELECT name FROM users WHERE id = 1"))
check r.success
check r.rows.len == 1
check valueToString(r.rows[0]["name"]) == "bob"
test "MERGE WHEN MATCHED THEN DELETE removes the row":
## Regression: mergeMatchedDelete was parsed but never executed.
var ctx = setupCtx()
defer: teardown(ctx)
discard executeQuery(ctx, parse("CREATE TABLE inv (id INTEGER PRIMARY KEY, qty INTEGER)"))
discard executeQuery(ctx, parse("INSERT INTO inv (id, qty) VALUES (1, 10)"))
discard executeQuery(ctx, parse("INSERT INTO inv (id, qty) VALUES (2, 20)"))
discard executeQuery(ctx, parse("CREATE TABLE deltas (id INTEGER PRIMARY KEY, qty INTEGER)"))
discard executeQuery(ctx, parse("INSERT INTO deltas (id, qty) VALUES (1, 0)"))
let r = executeQuery(ctx, parse("""
MERGE INTO inv AS t
USING deltas AS s
ON t.id = s.id
WHEN MATCHED THEN DELETE
"""))
check r.success
check r.affectedRows >= 1
let left = executeQuery(ctx, parse("SELECT id FROM inv ORDER BY id"))
check left.success
check left.rows.len == 1
check valueToString(left.rows[0]["id"]) == "2"
test "semi-sync writeLsn returns 0 when replicas do not ack":
var rm = newReplicationManager(rmSemiSync, syncCount = 1)
rm.addReplica(newReplica("r1", "10.0.0.1", 9472))
rm.connectReplica("r1")
let lsn = rm.writeLsn(@[1'u8, 2, 3])
check lsn == 0
+19 -5
View File
@@ -1659,14 +1659,28 @@ suite "Replication":
rm.connectReplica("r2") rm.connectReplica("r2")
rm.connectReplica("r3") rm.connectReplica("r3")
# Unreachable replicas cannot ack — semi-sync must fail closed (return 0)
let lsn = rm.writeLsn(@[1'u8]) let lsn = rm.writeLsn(@[1'u8])
check not rm.isFullyAcked(lsn) # needs 2 acks check lsn == 0
rm.ackLsn("r1", lsn) # No connected replicas → nothing to wait for; write succeeds
check not rm.isFullyAcked(lsn) # still needs 1 more var rm2 = newReplicationManager(rmSemiSync, syncCount = 2)
rm2.addReplica(newReplica("r1", "10.0.0.1", 9472))
# not connected
let lsn2 = rm2.writeLsn(@[1'u8])
check lsn2 > 0
check rm2.isFullyAcked(lsn2)
rm.ackLsn("r2", lsn) # ackLsn bookkeeping still clears pendingAcks at the required quorum
check rm.isFullyAcked(lsn) # 2 acks received var rm3 = newReplicationManager(rmSemiSync, syncCount = 2)
rm3.pendingAcks[1'u64] = initHashSet[string]()
rm3.pendingAcks[1'u64].incl("r1")
rm3.pendingAcks[1'u64].incl("r2")
check not rm3.isFullyAcked(1)
rm3.ackLsn("r1", 1)
check not rm3.isFullyAcked(1)
rm3.ackLsn("r2", 1)
check rm3.isFullyAcked(1)
test "Replica status": test "Replica status":
var rm = newReplicationManager(rmAsync) var rm = newReplicationManager(rmAsync)