diff --git a/.gitignore b/.gitignore index 6dac124..a9d115b 100644 --- a/.gitignore +++ b/.gitignore @@ -71,5 +71,6 @@ src/barabadb/storage/lsm src/barabadb/storage/wal src/barabadb/storage/btree src/barabadb/storage/gate +src/barabadb/protocol/scram clients/nim/tests/test_pool clients/nim/tests/test_wire diff --git a/BUG_AUDIT_2026-08.md b/BUG_AUDIT_2026-08.md index eadbf3e..3f321cb 100644 --- a/BUG_AUDIT_2026-08.md +++ b/BUG_AUDIT_2026-08.md @@ -3,7 +3,7 @@ > Дата: 2026-08-02 > Метод: 4 паралелни одит-агента по слоеве (Storage / Query / Core / Protocol), всеки чете всички файлове в обхвата си и проверява находките срещу реалния код. > Обхват: **само нови дефекти** — 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 регресионни теста | | 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 | |---|---------|------|---------------| | 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` | -| 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 ` не се изпълняват** — 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 | +| 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 | -### 🟡 MEDIUM (11) +### 🟡 MEDIUM (7) | # | Проблем | Файл | Предложен 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 | -| M3 | **`checkpoint` leaking write lock при exception** — `acquireWrite` без try/finally; IOError от flush/rotate пропуска `releaseWrite` → постоянен hang. | `storage/lsm.nim:932` | try/finally около lock-а (и walLock) | -| 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-ите | +| M7 | **Compaction unlink-ва input-ите преди output-ът да е loadable в каталога** — verifySSTable вече е преди unlink; остава catalog re-load ordering в caller. | `storage/compaction.nim` / LSM apply | Load/verify output в каталога ПРЕДИ unlink на input-ите | | M8 | **`OFFSET n` без `LIMIT` връща 0 реда; negative `LIMIT` чупи** — `limitCount = 0` е sentinel и за "няма limit", и за "LIMIT 0"; `sourceRows[start.. 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 (по-голям рефакторинг) | -### Хигиена - -- **Stray 94 KB компилиран ELF binary** в `src/barabadb/protocol/scram` — случайно commit-нат в source tree-то; да се премахне (+ `.gitignore`). - --- ## Проверени и чисти (не са бъгове) diff --git a/CHANGELOG.md b/CHANGELOG.md index dafab0c..edfe6c0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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`) - **`**` / `++` 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`) +- **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 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`) +### Removed + +- Stray compiled ELF `src/barabadb/protocol/scram` from the source tree (added to `.gitignore`) + ### 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 1–2, ~12 tracked) --- diff --git a/PLAN.md b/PLAN.md index bf5d4d4..a0a7ea6 100644 --- a/PLAN.md +++ b/PLAN.md @@ -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). -**Остават (~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 (M10–M12), SCRAM (L1–L2), NULL equality (L4). --- @@ -179,7 +181,7 @@ | **Този план** — Сесии 10, 11, 12 | ✅ Завършен | | 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` | -| **Сесия 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` | --- diff --git a/src/barabadb/core/mvcc.nim b/src/barabadb/core/mvcc.nim index 839e830..a42f4db 100644 --- a/src/barabadb/core/mvcc.nim +++ b/src/barabadb/core/mvcc.nim @@ -178,12 +178,16 @@ proc write*(tm: TxnManager, txn: Transaction, key: string, value: seq[byte]): bo return false # Timeout-based deadlock detection: abort stale transactions + # Collect then delete — never mutate activeTxns while iterating it. let now = getMonoTime().ticks() + var staleIds: seq[TxnId] = @[] for otherId, otherTxn in tm.activeTxns: if otherId != txn.id and otherTxn.state == tsActive: if now - otherTxn.startTime > tm.txnTimeoutMs * 1_000_000: 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 for otherId, otherTxn in tm.activeTxns: diff --git a/src/barabadb/core/replication.nim b/src/barabadb/core/replication.nim index cf813ca..45ab7c5 100644 --- a/src/barabadb/core/replication.nim +++ b/src/barabadb/core/replication.nim @@ -213,10 +213,18 @@ proc writeLsn*(rm: ReplicationManager, data: seq[byte]): uint64 = rm.pendingAcks[lsn].excl(id) if rm.pendingAcks[lsn].len == 0: 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) - 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 proc ackLsn*(rm: ReplicationManager, replicaId: string, lsn: uint64) = diff --git a/src/barabadb/protocol/scram b/src/barabadb/protocol/scram deleted file mode 100755 index cec9518..0000000 Binary files a/src/barabadb/protocol/scram and /dev/null differ diff --git a/src/barabadb/query/exec/lower.nim b/src/barabadb/query/exec/lower.nim index 735a2a0..c4be989 100644 --- a/src/barabadb/query/exec/lower.nim +++ b/src/barabadb/query/exec/lower.nim @@ -122,6 +122,7 @@ proc lowerExpr*(node: Node): IRExpr = else: discard result.aggArgs = @[] for arg in node.funcArgs: result.aggArgs.add(lowerExpr(arg)) + result.aggDistinct = node.funcDistinct if node.funcFilter != nil: result.aggFilter = lowerExpr(node.funcFilter) else: diff --git a/src/barabadb/query/exec/plan_exec.nim b/src/barabadb/query/exec/plan_exec.nim index ef103e0..fffb3db 100644 --- a/src/barabadb/query/exec/plan_exec.nim +++ b/src/barabadb/query/exec/plan_exec.nim @@ -5,6 +5,7 @@ ## executor split). Pure code motion — no behavior changes. import std/strutils import std/tables +import std/sets import std/sequtils import std/algorithm import ../ir @@ -19,6 +20,19 @@ import eval import scan 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) # ---------------------------------------------------------------------- @@ -113,49 +127,71 @@ proc executePlan*(ctx: ExecutionContext, plan: IRPlan): seq[Row] = newRow[alias] = $filteredRows.len else: var count = 0 + var seen: HashSet[string] for row in filteredRows: 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 of irSum: var sum = 0.0 + var seen: HashSet[string] for row in filteredRows: 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 of irAvg: var sum = 0.0 var count = 0 + var seen: HashSet[string] for row in filteredRows: 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" of irMin: var minVal = "" + var seen: HashSet[string] for row in filteredRows: let v = evalExpr(expr.aggArgs[0], row, ctx) 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 of irMax: var maxVal = "" + var seen: HashSet[string] for row in filteredRows: let v = evalExpr(expr.aggArgs[0], row, ctx) 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 of irArrayAgg: var arr: seq[string] + var seen: HashSet[string] for row in filteredRows: 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(", ") & "]" of irStringAgg: 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: ",") for row in filteredRows: 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)) else: 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 else: var count = 0 + var seen: HashSet[string] for row in filteredRows: 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 of irSum: var sum = 0.0 + var seen: HashSet[string] for row in filteredRows: 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 of irAvg: var sum = 0.0 var count = 0 + var seen: HashSet[string] for row in filteredRows: 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" of irMin: var minVal = "" + var seen: HashSet[string] for row in filteredRows: let v = evalExpr(aggExpr.aggArgs[0], row, ctx) 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 of irMax: var maxVal = "" + var seen: HashSet[string] for row in filteredRows: let v = evalExpr(aggExpr.aggArgs[0], row, ctx) 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 of irArrayAgg: var arr: seq[string] + var seen: HashSet[string] for row in filteredRows: 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(", ") & "]" of irStringAgg: 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: ",") for row in filteredRows: 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)) # Apply HAVING filter if plan.groupHaving != nil: diff --git a/src/barabadb/query/executor.nim b/src/barabadb/query/executor.nim index 4f1dcdb..79cfc3d 100644 --- a/src/barabadb/query/executor.nim +++ b/src/barabadb/query/executor.nim @@ -448,6 +448,26 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu if cols.len == 0: 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] = @[] case stmt.setOpKind of sdkUnion: @@ -460,28 +480,30 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu # UNION: deduplicate var seen: Table[string, bool] for row in leftRes.rows: - seen[valueToString(row["$value"])] = true + seen[setOpRowKey(row, cols)] = true for row in rightRes.rows: - if not seen.getOrDefault(valueToString(row["$value"]), false): - seen[valueToString(row["$value"])] = true + let k = setOpRowKey(row, cols) + if not seen.getOrDefault(k, false): + seen[k] = true rows.add(row) of sdkIntersect: var leftSet: Table[string, bool] for row in leftRes.rows: - leftSet[valueToString(row["$value"])] = true + leftSet[setOpRowKey(row, cols)] = true 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) 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: var rightSet: Table[string, bool] for row in rightRes.rows: - rightSet[valueToString(row["$value"])] = true + rightSet[setOpRowKey(row, cols)] = true for row in leftRes.rows: - if not rightSet.getOrDefault(valueToString(row["$value"]), false): + if not rightSet.getOrDefault(setOpRowKey(row, cols), false): rows.add(row) return okResult(rows, cols) @@ -759,21 +781,34 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu let onExpr = lowerExpr(stmt.mergeOn) if valueToString(evalExpr(onExpr, rowWithTarget, ctx)) == "true": matched = true - if stmt.mergeMatchedUpdate.len > 0 and "$key" in tgtRow: - var updateSets = initTable[string, string]() - for s in stmt.mergeMatchedUpdate: - if s.kind == nkBinOp and s.binOp == bkAssign: - if s.binLeft.kind == nkIdent: - let valExpr = lowerExpr(s.binRight) - updateSets[s.binLeft.identName] = valueToString(evalExpr(valExpr, rowWithTarget, ctx)) - var newRow = tgtRow - for col, val in updateSets: - newRow[col] = Value(kind: vkString, strVal: val) - fireTriggers(ctx, stmt.mergeTarget, "before", "update", tgtRow) - count += execUpdateRow(ctx, stmt.mergeTarget, valueToString(tgtRow["$key"]), updateSets, kvPairs) - fireTriggers(ctx, stmt.mergeTarget, "after", "update", newRow) - if ctx.onChange != nil: - ctx.onChange(ChangeEvent(table: stmt.mergeTarget, kind: ckUpdate, key: valueToString(tgtRow["$key"]), data: "")) + # Optional AND 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]() + for s in stmt.mergeMatchedUpdate: + if s.kind == nkBinOp and s.binOp == bkAssign: + if s.binLeft.kind == nkIdent: + let valExpr = lowerExpr(s.binRight) + updateSets[s.binLeft.identName] = valueToString(evalExpr(valExpr, rowWithTarget, ctx)) + var newRow = tgtRow + for col, val in updateSets: + newRow[col] = Value(kind: vkString, strVal: val) + fireTriggers(ctx, stmt.mergeTarget, "before", "update", tgtRow) + count += execUpdateRow(ctx, stmt.mergeTarget, valueToString(tgtRow["$key"]), updateSets, kvPairs) + fireTriggers(ctx, stmt.mergeTarget, "after", "update", newRow) + if ctx.onChange != nil: + ctx.onChange(ChangeEvent(table: stmt.mergeTarget, kind: ckUpdate, key: valueToString(tgtRow["$key"]), data: "")) break if not matched and stmt.mergeNotMatchedInsert.len > 0: diff --git a/src/barabadb/storage/compaction.nim b/src/barabadb/storage/compaction.nim index 9404c1e..6ec15db 100644 --- a/src/barabadb/storage/compaction.nim +++ b/src/barabadb/storage/compaction.nim @@ -125,13 +125,16 @@ proc compact*(cs: CompactionStrategy, level: int): CompactionResult = 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 lastKey = "" + var haveLast = false for entry in allEntries: - if entry.key != lastKey: + if not haveLast or entry.key != lastKey: merged.add(entry) lastKey = entry.key + haveLast = true # Keep tombstones to prevent deleted keys from resurrecting in lower levels var final: seq[Entry] = @[] diff --git a/src/barabadb/storage/lsm.nim b/src/barabadb/storage/lsm.nim index c310a88..7316706 100644 --- a/src/barabadb/storage/lsm.nim +++ b/src/barabadb/storage/lsm.nim @@ -696,6 +696,10 @@ proc newLSMTree*( var version: uint32 = 0 if stream.readData(addr magic, 4) == 4 and magic == WALMagic: 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(): var kind: uint8 = 0 var timestamp: uint64 = 0 @@ -704,10 +708,20 @@ proc newLSMTree*( if stream.readData(addr kind, 1) != 1: break if stream.readData(addr timestamp, 8) != 8: 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) if keyLen > 0: if stream.readData(addr key[0], keyLen.int) != keyLen.int: 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) if valLen > 0: 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: return - # Flush immutable memtable if present, otherwise flush current memtable - var toFlush = db.immutableMem - if toFlush.len == 0: - toFlush = db.memTable - db.memTable = newMemTable(db.memMaxSize) + # Flush immutable memtable if present, otherwise flush current memtable. + # Do NOT clear the source memtable until the SSTable is written — an IOError + # mid-write must leave the data still visible to live reads (WAL still has it). + var flushingImmutable = false + var toFlush: MemTable + if db.immutableMem.len > 0: + toFlush = db.immutableMem + flushingImmutable = true else: - db.immutableMem = newMemTable(0) + toFlush = db.memTable if toFlush.len == 0: return let path = db.dir / "sstables" / ($db.nextSSTableId & ".sst") + let sstId = db.nextSSTableId inc db.nextSSTableId # Sort once at flush time (O(n log n)) — put/get stay O(1) var sst = writeSSTable(toFlush.sortedEntries(), path, level = 0) - sst.id = db.nextSSTableId - 1 + sst.id = sstId db.sstables.add(sst) # 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 inc db.manifestSequence try: @@ -930,27 +954,29 @@ proc checkpoint*(db: LSMTree) = ## rotate WAL, and write MANIFEST. This provides a clean boundary ## for online backup without stopping the server. acquireWrite(db.lock) + try: + # Flush any pending immutable memtable first + if db.immutableMem.len > 0: + flushUnsafe(db) - # Flush any pending immutable memtable first - if db.immutableMem.len > 0: - flushUnsafe(db) + # Freeze current memtable so writes can continue on a new one + if db.memTable.len > 0: + db.immutableMem = db.memTable + db.memTable = newMemTable(db.memMaxSize) - # Freeze current memtable so writes can continue on a new one - if db.memTable.len > 0: - db.immutableMem = db.memTable - db.memTable = newMemTable(db.memMaxSize) + # Flush the frozen memtable + if db.immutableMem.len > 0: + flushUnsafe(db) - # Flush the frozen memtable - if db.immutableMem.len > 0: - flushUnsafe(db) - - # Rotate WAL for a clean backup boundary - acquire(db.walLock) - db.wal.maybeRotate() - db.wal.sync() - release(db.walLock) - - releaseWrite(db.lock) + # Rotate WAL for a clean backup boundary + acquire(db.walLock) + try: + db.wal.maybeRotate() + db.wal.sync() + finally: + release(db.walLock) + finally: + releaseWrite(db.lock) proc close*(db: LSMTree) = acquireWrite(db.lock) diff --git a/src/barabadb/storage/mmap.nim b/src/barabadb/storage/mmap.nim index af3612d..c13bd2d 100644 --- a/src/barabadb/storage/mmap.nim +++ b/src/barabadb/storage/mmap.nim @@ -80,7 +80,8 @@ proc readAt*(mf: MmapFile, offset: int, size: int): seq[byte] = if mf.regions.len == 0: return @[] 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 @[] result = newSeq[byte](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] 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 var val: uint32 copyMem(addr val, unsafeAddr mf.regions[0].data[offset], 4) return val 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 var val: uint64 copyMem(addr val, unsafeAddr mf.regions[0].data[offset], 8) return val 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 "" result = newString(size) copyMem(addr result[0], unsafeAddr mf.regions[0].data[offset], size) diff --git a/src/barabadb/storage/recovery.nim b/src/barabadb/storage/recovery.nim index daed440..488c0d1 100644 --- a/src/barabadb/storage/recovery.nim +++ b/src/barabadb/storage/recovery.nim @@ -68,6 +68,7 @@ proc scanWAL*(rec: CrashRecovery): seq[RecoveredEntry] = var txnId: uint64 = 0 var entryCount = 0 + const MaxWalRecordField = 64 * 1024 * 1024 # 64 MB while not stream.atEnd(): var kind: uint8 = 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 timestamp, 8) != 8: 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) if keyLen > 0: if stream.readData(addr key[0], keyLen.int) != keyLen.int: break if stream.readData(addr valLen, 4) != 4: break + if valLen.int > MaxWalRecordField: break var value = newSeq[byte](valLen.int) if valLen > 0: if stream.readData(addr value[0], valLen.int) != valLen.int: break diff --git a/src/barabadb/storage/wal.nim b/src/barabadb/storage/wal.nim index 1d4cab2..8419e17 100644 --- a/src/barabadb/storage/wal.nim +++ b/src/barabadb/storage/wal.nim @@ -326,8 +326,9 @@ proc rewriteLive*(wal: var WriteAheadLog, if wal.stream != nil: wal.stream.close() - if fileExists(wal.path): - removeFile(wal.path) + wal.stream = nil + # 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) wal.stream = newFileStream(wal.path, fmAppend) 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 version, 4) != 4: return if magic != WALMagic: return + const MaxWalRecordField = 64 * 1024 * 1024 # 64 MB while not s.atEnd: var kind: uint8 if s.readData(addr kind, 1) != 1: break @@ -372,11 +374,14 @@ proc readEntries*(walPath: string, untilTimestamp: uint64 = 0): seq[WalEntry] = break var keyLen: uint32 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) if keyLen > 0: if s.readData(addr key[0], int(keyLen)) != int(keyLen): break var valLen: uint32 if s.readData(addr valLen, 4) != 4: break + if valLen.int > MaxWalRecordField: break var value = newSeq[byte](valLen) if valLen > 0: if s.readData(addr value[0], int(valLen)) != int(valLen): break diff --git a/tests/bugfix_test.nim b/tests/bugfix_test.nim index a18d27f..7204d72 100644 --- a/tests/bugfix_test.nim +++ b/tests/bugfix_test.nim @@ -625,3 +625,99 @@ suite "Query operator correctness — audit batch 1": let neq = executeQuery(ctx, parse("SELECT * FROM users WHERE id != 1.0")) check eq.rows.len == 1 # 1 = 1.0 -> true 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 diff --git a/tests/test_all.nim b/tests/test_all.nim index 0f829dc..3d4b040 100644 --- a/tests/test_all.nim +++ b/tests/test_all.nim @@ -1659,14 +1659,28 @@ suite "Replication": rm.connectReplica("r2") rm.connectReplica("r3") + # Unreachable replicas cannot ack — semi-sync must fail closed (return 0) let lsn = rm.writeLsn(@[1'u8]) - check not rm.isFullyAcked(lsn) # needs 2 acks + check lsn == 0 - rm.ackLsn("r1", lsn) - check not rm.isFullyAcked(lsn) # still needs 1 more + # No connected replicas → nothing to wait for; write succeeds + 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) - check rm.isFullyAcked(lsn) # 2 acks received + # ackLsn bookkeeping still clears pendingAcks at the required quorum + 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": var rm = newReplicationManager(rmAsync)