fix(raft): distinguish put-with-empty-value from delete in write path

This commit is contained in:
2026-07-31 00:58:51 +03:00
parent 63cb05afe2
commit 431334b70a
9 changed files with 119 additions and 51 deletions
+17 -12
View File
@@ -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
+6 -6
View File
@@ -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 & "."):
-1
View File
@@ -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:
+4 -4
View File
@@ -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)
+2 -2
View File
@@ -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)
+11 -10
View File
@@ -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)
+73
View File
@@ -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":
+3 -13
View File
@@ -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()
+3 -3
View File
@@ -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