diff --git a/src/barabadb/core/server.nim b/src/barabadb/core/server.nim index cc6f848..458718d 100644 --- a/src/barabadb/core/server.nim +++ b/src/barabadb/core/server.nim @@ -306,10 +306,12 @@ proc waitRaftCommit(node: RaftNode, lastIdx: uint64, timeoutMs: int): Future[(bo node.metrics.commitWaitMsTotal += waitedMs return (true, "") -proc appendWriteToRaft*(node: RaftNode, kvPairs: seq[(string, seq[byte])], +proc appendWriteToRaft*(node: RaftNode, + kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]], timeoutMs: int): Future[(bool, string)] {.async.} = ## C3b leader write path: append each written KV pair to the Raft log and - ## wait for majority commit. An empty value encodes a delete; the entry + ## wait for majority commit. The `deleted` flag encodes a delete — an empty + ## value alone is a put (PK-only tables store an empty LSM value); the entry ## format matches applyCommand ("put": key \x00 value, "delete": key). ## ## MUST be called from the async event-loop thread that owns `node` and @@ -317,11 +319,11 @@ proc appendWriteToRaft*(node: RaftNode, kvPairs: seq[(string, seq[byte])], ## handleAppendReply on the same loop, and applyCommand re-enters the ## (non-reentrant) gate — waiting under the gate would deadlock the loop. var lastIdx = 0'u64 - for (key, value) in kvPairs: - let entry = if value.len > 0: - node.appendLog("put", cast[seq[byte]](key & "\x00" & cast[string](value))) + for pair in kvPairs: + let entry = if pair.deleted: + node.appendLog("delete", cast[seq[byte]](pair.key)) else: - node.appendLog("delete", cast[seq[byte]](key)) + node.appendLog("put", cast[seq[byte]](pair.key & "\x00" & cast[string](pair.value))) if entry.index == 0: if node.metrics != nil: inc node.metrics.lostLeadershipTotal @@ -353,7 +355,7 @@ proc executeQuery(db: LSMTree, ctx: ExecutionContext, query: string, params: seq var ok = false var qr = QueryResult() var msg = "" - var kvPairs: seq[(string, seq[byte])] = @[] + var kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] = @[] var needsRaftDdl = false var needsForward = false var forwardHost = "" @@ -397,11 +399,14 @@ proc executeQuery(db: LSMTree, ctx: ExecutionContext, query: string, params: seq # Ship written key-value pairs to replicas (legacy path; skipped when # the raft path below handles the statement). if raftNode == nil and replication != nil and res.keyValuePairs.len > 0: - for (key, value) in res.keyValuePairs: - var data = newSeq[byte](key.len + 1 + value.len) - for i, c in key: data[i] = byte(c) - data[key.len] = byte(0) - for i, c in value: data[key.len + 1 + i] = c + for pair in res.keyValuePairs: + # Legacy REP wire format: key \x00 value, empty value = delete + # on the receiver. Deletes ship an empty value as before. + let value = if pair.deleted: @[] else: pair.value + var data = newSeq[byte](pair.key.len + 1 + value.len) + for i, c in pair.key: data[i] = byte(c) + data[pair.key.len] = byte(0) + for i, c in value: data[pair.key.len + 1 + i] = c discard replication.writeLsn(data) qr = QueryResult(affectedRows: res.affectedRows, rowCount: res.rows.len) qr.columns = res.columns diff --git a/src/barabadb/query/exec/dml.nim b/src/barabadb/query/exec/dml.nim index c465c3e..7fcd5ac 100644 --- a/src/barabadb/query/exec/dml.nim +++ b/src/barabadb/query/exec/dml.nim @@ -45,7 +45,7 @@ proc violatesUniqueIndex*(ctx: ExecutionContext, table: string, fields: seq[stri return "" proc execInsert*(ctx: ExecutionContext, table: string, fields: seq[string], values: seq[seq[string]], - kvPairs: var seq[(string, seq[byte])]): int = + kvPairs: var seq[tuple[key: string, value: seq[byte], deleted: bool]]): int = if not hasPrivilege(ctx, table, "INSERT"): return 0 let tblDef = if table in ctx.tables: ctx.tables[table] else: TableDef() @@ -87,7 +87,7 @@ proc execInsert*(ctx: ExecutionContext, table: string, fields: seq[string], valu discard ctx.txnManager.write(ctx.pendingTxn, fullKey, cast[seq[byte]](valStr)) else: ctx.db.put(fullKey, cast[seq[byte]](valStr)) - kvPairs.add((fullKey, cast[seq[byte]](valStr))) + kvPairs.add((fullKey, cast[seq[byte]](valStr), false)) for colName in ctx.btrees.keys.toSeq(): if colName.startsWith(table & "."): @@ -221,7 +221,7 @@ proc execInsert*(ctx: ExecutionContext, table: string, fields: seq[string], valu return count proc execDelete*(ctx: ExecutionContext, table: string, key: string, - kvPairs: var seq[(string, seq[byte])]): int = + kvPairs: var seq[tuple[key: string, value: seq[byte], deleted: bool]]): int = if not hasPrivilege(ctx, table, "DELETE"): return 0 let fullKey = table & "." & key @@ -238,7 +238,7 @@ proc execDelete*(ctx: ExecutionContext, table: string, key: string, discard ctx.txnManager.delete(ctx.pendingTxn, fullKey) else: ctx.db.delete(fullKey) - kvPairs.add((fullKey, @[])) + kvPairs.add((fullKey, @[], true)) # Update BTree indexes for colName in ctx.btrees.keys.toSeq(): if colName.startsWith(table & "."): @@ -264,7 +264,7 @@ proc execDelete*(ctx: ExecutionContext, table: string, key: string, return 0 proc execUpdateRow*(ctx: ExecutionContext, table: string, key: string, sets: Table[string, string], - kvPairs: var seq[(string, seq[byte])]): int = + kvPairs: var seq[tuple[key: string, value: seq[byte], deleted: bool]]): int = if not hasPrivilege(ctx, table, "UPDATE"): return 0 let fullKey = table & "." & key @@ -313,7 +313,7 @@ proc execUpdateRow*(ctx: ExecutionContext, table: string, key: string, sets: Tab discard ctx.txnManager.write(ctx.pendingTxn, fullKey, cast[seq[byte]](newVal)) else: ctx.db.put(fullKey, cast[seq[byte]](newVal)) - kvPairs.add((fullKey, cast[seq[byte]](newVal))) + kvPairs.add((fullKey, cast[seq[byte]](newVal), false)) # Update FTS indexes: remove old doc, add new for ftsKey, ftsIdx in ctx.ftsIndexes: if ftsKey.startsWith(table & "."): diff --git a/src/barabadb/query/exec/eval.nim b/src/barabadb/query/exec/eval.nim index a948371..4f15125 100644 --- a/src/barabadb/query/exec/eval.nim +++ b/src/barabadb/query/exec/eval.nim @@ -923,7 +923,6 @@ proc evalExprOld*(expr: IRExpr, row: Table[string, string], ctx: ExecutionContex ddl.add("\n") # Sample data - var kvPairs: seq[(string, seq[byte])] = @[] let rows = requireExecScanHook()(ctx, table) let sampleLimit = min(5, rows.len) if sampleLimit > 0: diff --git a/src/barabadb/query/exec/fk.nim b/src/barabadb/query/exec/fk.nim index 7b2d9d3..a4a3f4b 100644 --- a/src/barabadb/query/exec/fk.nim +++ b/src/barabadb/query/exec/fk.nim @@ -30,14 +30,14 @@ proc enforceFkOnDelete*(ctx: ExecutionContext, parentTable: string, parentCol: s of "CASCADE": for refRow in refs: if "$key" in refRow: - var dummy: seq[(string, seq[byte])] = @[] + var dummy: seq[tuple[key: string, value: seq[byte], deleted: bool]] = @[] discard execDelete(ctx, childTblName, valueToString(refRow["$key"]), dummy) of "SET NULL": for refRow in refs: if "$key" in refRow: var sets = initTable[string, string]() sets[col.name] = "\\N" - var dummy: seq[(string, seq[byte])] = @[] + var dummy: seq[tuple[key: string, value: seq[byte], deleted: bool]] = @[] discard execUpdateRow(ctx, childTblName, valueToString(refRow["$key"]), sets, dummy) of "RESTRICT", "NO ACTION": return (false, "FOREIGN KEY violation: row is referenced by " & childTblName & "." & col.name) @@ -56,14 +56,14 @@ proc enforceFkOnUpdate*(ctx: ExecutionContext, parentTable: string, parentCol: s if "$key" in refRow: var sets = initTable[string, string]() sets[col.name] = newVal - var dummy: seq[(string, seq[byte])] = @[] + var dummy: seq[tuple[key: string, value: seq[byte], deleted: bool]] = @[] discard execUpdateRow(ctx, childTblName, valueToString(refRow["$key"]), sets, dummy) of "SET NULL": for refRow in refs: if "$key" in refRow: var sets = initTable[string, string]() sets[col.name] = "\\N" - var dummy: seq[(string, seq[byte])] = @[] + var dummy: seq[tuple[key: string, value: seq[byte], deleted: bool]] = @[] discard execUpdateRow(ctx, childTblName, valueToString(refRow["$key"]), sets, dummy) of "RESTRICT", "NO ACTION": return (false, "FOREIGN KEY violation: row is referenced by " & childTblName & "." & col.name) diff --git a/src/barabadb/query/exec/types.nim b/src/barabadb/query/exec/types.nim index 61de048..f02c039 100644 --- a/src/barabadb/query/exec/types.nim +++ b/src/barabadb/query/exec/types.nim @@ -136,13 +136,13 @@ type rows*: seq[Row] affectedRows*: int message*: string - keyValuePairs*: seq[(string, seq[byte])] + keyValuePairs*: seq[tuple[key: string, value: seq[byte], deleted: bool]] proc `==`*(a, b: IndexEntry): bool = a.lsmKey == b.lsmKey and a.rowValue == b.rowValue proc okResult*(rows: seq[Row] = @[], cols: seq[string] = @[], affected: int = 0, msg: string = "", - kvPairs: seq[(string, seq[byte])] = @[]): ExecResult = + kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] = @[]): ExecResult = ExecResult(success: true, columns: cols, rows: rows, affectedRows: affected, message: msg, keyValuePairs: kvPairs) diff --git a/src/barabadb/query/executor.nim b/src/barabadb/query/executor.nim index 43dca17..4f1dcdb 100644 --- a/src/barabadb/query/executor.nim +++ b/src/barabadb/query/executor.nim @@ -574,7 +574,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu row[f] = mutableValues[0][i] fireTriggers(ctx, stmt.insTarget, "before", "insert", row) - var kvPairs: seq[(string, seq[byte])] + var kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] let count = execInsert(ctx, stmt.insTarget, mutableFields, mutableValues, kvPairs) # Fire AFTER INSERT triggers @@ -626,7 +626,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu # Scan and apply let rows = execScan(ctx, stmt.updTarget) var count = 0 - var kvPairs: seq[(string, seq[byte])] + var kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] for row in rows: # Compute sets for this row (expressions may reference columns) var sets = initTable[string, string]() @@ -701,7 +701,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu # Delete all rows matching WHERE let rows = execScan(ctx, stmt.delTarget) var count = 0 - var kvPairs: seq[(string, seq[byte])] + var kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] for row in rows: if stmt.delWhere != nil and stmt.delWhere.whereExpr != nil: let whereExpr = lowerExpr(stmt.delWhere.whereExpr) @@ -744,7 +744,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu let targetRows = execScan(ctx, stmt.mergeTarget) var count = 0 - var kvPairs: seq[(string, seq[byte])] + var kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] for srcRow in sourceRows: var matched = false @@ -793,7 +793,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu for i, f in fields: if i < values.len: row[f] = Value(kind: vkString, strVal: values[i]) fireTriggers(ctx, stmt.mergeTarget, "before", "insert", row) - var insKvPairs: seq[(string, seq[byte])] + var insKvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] count += execInsert(ctx, stmt.mergeTarget, fields, @[values], insKvPairs) for kv in insKvPairs: kvPairs.add(kv) fireTriggers(ctx, stmt.mergeTarget, "after", "insert", row) @@ -982,16 +982,17 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu of nkCommitTxn: if ctx.pendingTxn != nil and ctx.pendingTxn.state == tsActive: - var kvPairs: seq[(string, seq[byte])] + var kvPairs: seq[tuple[key: string, value: seq[byte], deleted: bool]] for key, version in ctx.pendingTxn.writeSet: if version.isDelete: ctx.db.delete(key) - # Empty value is the raft/replication "delete" convention — never - # ship a non-empty body for isDelete or followers will resurrect. - kvPairs.add((key, @[])) + # Empty value + deleted=true is the raft/replication "delete" + # convention — never ship a non-empty body for isDelete or + # followers will resurrect. + kvPairs.add((key, @[], true)) else: ctx.db.put(key, version.value) - kvPairs.add((key, version.value)) + kvPairs.add((key, version.value, false)) discard ctx.txnManager.commit(ctx.pendingTxn) ctx.pendingTxn = nil return okResult(msg="Transaction committed", kvPairs=kvPairs) diff --git a/tests/bugfix_test.nim b/tests/bugfix_test.nim index b8df821..3cfb2a0 100644 --- a/tests/bugfix_test.nim +++ b/tests/bugfix_test.nim @@ -4,6 +4,7 @@ import std/os import std/tables import ../src/barabadb/query/[parser, executor, lexer, ast] import ../src/barabadb/query/exec/params +import ../src/barabadb/query/exec/dml import ../src/barabadb/core/types import ../src/barabadb/core/config import ../src/barabadb/storage/lsm @@ -417,6 +418,78 @@ suite "Raft peer address parsing": check msg.len > 0 check bad in msg +suite "Raft put/delete encoding — empty value is not a delete": + + test "PK-only INSERT yields a put pair (deleted == false, empty value)": + var ctx = setupCtx() + defer: teardown(ctx) + discard executeQuery(ctx, parse("CREATE TABLE pkonly (id INTEGER PRIMARY KEY)")) + let r = executeQuery(ctx, parse("INSERT INTO pkonly (id) VALUES (1)")) + check r.success + check r.keyValuePairs.len == 1 + check r.keyValuePairs[0].value.len == 0 + check r.keyValuePairs[0].deleted == false + + test "DELETE yields a delete pair (deleted == true)": + var ctx = setupCtx() + defer: teardown(ctx) + discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')")) + let r = executeQuery(ctx, parse("DELETE FROM users WHERE id = 1")) + check r.success + check r.keyValuePairs.len == 1 + check r.keyValuePairs[0].deleted == true + check r.keyValuePairs[0].value.len == 0 + + test "UPDATE yields a put pair (deleted == false)": + var ctx = setupCtx() + defer: teardown(ctx) + discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')")) + let r = executeQuery(ctx, parse("UPDATE users SET name = 'bob' WHERE id = 1")) + check r.success + check r.keyValuePairs.len == 1 + check r.keyValuePairs[0].deleted == false + check r.keyValuePairs[0].value.len > 0 + + test "txn COMMIT pairs carry deleted flag for buffered writes": + var ctx = setupCtx() + defer: teardown(ctx) + discard executeQuery(ctx, parse("CREATE TABLE pkonly (id INTEGER PRIMARY KEY)")) + discard executeQuery(ctx, parse("INSERT INTO users (id, name) VALUES (1, 'alice')")) + discard executeQuery(ctx, parse("BEGIN")) + discard executeQuery(ctx, parse("INSERT INTO pkonly (id) VALUES (7)")) + discard executeQuery(ctx, parse("DELETE FROM users WHERE id = 1")) + let r = executeQuery(ctx, parse("COMMIT")) + check r.success + check r.keyValuePairs.len == 2 + var sawPut = false + var sawDelete = false + for pair in r.keyValuePairs: + if pair.deleted: + sawDelete = true + check pair.value.len == 0 + else: + sawPut = true + check pair.key == "pkonly.id=7" + check pair.value.len == 0 # empty value must still be a put + check sawPut and sawDelete + + test "apply of a put with empty value keeps the PK-only row": + var ctx = setupCtx() + defer: teardown(ctx) + discard executeQuery(ctx, parse("CREATE TABLE pkonly (id INTEGER PRIMARY KEY)")) + let r = executeQuery(ctx, parse("INSERT INTO pkonly (id) VALUES (3)")) + check r.success + check r.keyValuePairs.len == 1 + let pair = r.keyValuePairs[0] + # Same decode as applyCommand in src/baradadb.nim for a "put" entry. + let encoded = pair.key & "\x00" & cast[string](pair.value) + let parts = encoded.split("\x00") + check parts.len >= 2 + applyReplicatedPut(ctx, parts[0], cast[seq[byte]](parts[1])) + let sel = executeQuery(ctx, parse("SELECT * FROM pkonly WHERE id = 3")) + check sel.success + check sel.rows.len == 1 + suite "Raft write classification": test "isWrite classifies DML and COMMIT": diff --git a/tests/raft_failover_load_e2e_test.nim b/tests/raft_failover_load_e2e_test.nim index fc19299..fab8a7f 100644 --- a/tests/raft_failover_load_e2e_test.nim +++ b/tests/raft_failover_load_e2e_test.nim @@ -7,16 +7,6 @@ ## follower has caught up. ## Process-management conventions follow tests/raft_writes_e2e_test.nim; ## client access follows tests/nimforum_smoke_test.nim. -## -## NOTE: load_test deliberately has a non-PK column. A table whose only -## column is the PK stores an EMPTY LSM value per row (execInsert drops PK -## columns from the value), and the raft write path (appendWriteToRaft in -## src/barabadb/core/server.nim) encodes empty values as "delete" entries — -## so on commit every node applies a delete over the just-inserted row and -## it vanishes everywhere. That is a v1.2.0 product bug in the raft entry -## encoding (put/delete must not be inferred from value emptiness); until it -## is fixed in src/, this test exercises the two-column row shape that the -## current encoding handles correctly. import std/unittest import std/osproc import std/os @@ -176,7 +166,7 @@ proc writerLoop(args: WriterArgs) {.thread.} = sleep(50) continue try: - db.exec(sql("INSERT INTO load_test (id, n) VALUES (" & $n & ", " & $n & ")")) + db.exec(sql("INSERT INTO load_test (id) VALUES (" & $n & ")")) withLock args.lock[]: args.acked[].add n inc n @@ -290,7 +280,7 @@ proc runFailoverLoadScenario() = block: let db = openClient(nodes[leaderIdx].clientPort) try: - db.exec(sql"CREATE TABLE load_test (id INT PRIMARY KEY, n INT)") + db.exec(sql"CREATE TABLE load_test (id INT PRIMARY KEY)") except CatchableError as e: echo "leader CREATE TABLE failed: ", e.msg dumpAll(nodes) @@ -343,7 +333,7 @@ proc runFailoverLoadScenario() = try: let db = openClient(nodes[i].clientPort) try: - db.exec(sql("INSERT INTO load_test (id, n) VALUES (" & $probeId & ", " & $probeId & ")")) + db.exec(sql("INSERT INTO load_test (id) VALUES (" & $probeId & ")")) writerSurvivor = i finally: db.close() diff --git a/tests/test_all.nim b/tests/test_all.nim index 161eb90..b32e3eb 100644 --- a/tests/test_all.nim +++ b/tests/test_all.nim @@ -2670,7 +2670,7 @@ suite "Raft SQL Write Path": # Server-side leader write path: append + wait for majority commit let (ok, errMsg) = waitFor appendWriteToRaft(leader, - @[("users.1", cast[seq[byte]]("alice"))], timeoutMs = 3000) + @[("users.1", cast[seq[byte]]("alice"), false)], timeoutMs = 3000) check ok if not ok: echo "appendWriteToRaft failed: ", errMsg @@ -2782,7 +2782,7 @@ suite "Raft SQL Write Path": var n = newRaftNode("n1", @["n2"], raftPort = 29111) # Still a follower — appendLog returns index 0. let (ok, err) = waitFor appendWriteToRaft(n, - @[("k", cast[seq[byte]]("v"))], timeoutMs = 200) + @[("k", cast[seq[byte]]("v"), false)], timeoutMs = 200) check not ok check "lost leadership" in err @@ -2791,7 +2791,7 @@ suite "Raft SQL Write Path": var n = newRaftNode("n1", @["n2", "n3"], raftPort = 29112) n.becomeLeader() let (ok, err) = waitFor appendWriteToRaft(n, - @[("k", cast[seq[byte]]("v"))], timeoutMs = 300) + @[("k", cast[seq[byte]]("v"), false)], timeoutMs = 300) check not ok check "raft commit timeout" in err