test(raft): E2E replicated writes; wire SQL path to raft node

- Fix runTcpServer to run the already-wired Server (raftNode was assigned
  on a different instance that never accepted clients).
- Cap raft peer connect at 200ms and fan out heartbeats in parallel so a
  dead peer cannot stall AppendEntries to the live majority.
- Add raft_writes_e2e_test: 3-node write replication, follower rejection,
  and post-failover writes; wire into nimble test + gitignore.
This commit is contained in:
2026-07-30 21:06:49 +03:00
parent 333941ab65
commit 44060701b7
5 changed files with 373 additions and 9 deletions
+1
View File
@@ -15,6 +15,7 @@ tests/prop_test
tests/bugfix_test tests/bugfix_test
tests/nimforum_smoke_test tests/nimforum_smoke_test
tests/raft_e2e_test tests/raft_e2e_test
tests/raft_writes_e2e_test
benchmarks/bench_all benchmarks/bench_all
benchmarks/compare benchmarks/compare
clients/nim/tests/test_client clients/nim/tests/test_client
+2 -1
View File
@@ -29,7 +29,8 @@ task test, "Run all tests":
# Quick embedded suites first, heavy fuzz/stress suites last. # Quick embedded suites first, heavy fuzz/stress suites last.
for t in ["test_minimal", "test_all", "bugfix_test", "join_tests", "test_lock", for t in ["test_minimal", "test_all", "bugfix_test", "join_tests", "test_lock",
"test_schema_persist", "test_storage_hardening", "tla_faithfulness", "test_schema_persist", "test_storage_hardening", "tla_faithfulness",
"nimforum_smoke_test", "raft_e2e_test", "fuzz_test", "prop_test", "nimforum_smoke_test", "raft_e2e_test", "raft_writes_e2e_test",
"fuzz_test", "prop_test",
"test_wire_insert_stress", "stress_test"]: "test_wire_insert_stress", "stress_test"]:
exec "nim c -r tests/" & t & ".nim" exec "nim c -r tests/" & t & ".nim"
+23 -4
View File
@@ -560,16 +560,26 @@ proc newRaftNetwork*(node: RaftNode): RaftNetwork =
timer: newElectionTimer(node, node.electionTimeout), timer: newElectionTimer(node, node.electionTimeout),
) )
const RaftConnectTimeoutMs = 200
proc connectToPeer(net: RaftNetwork, peerId: string) {.async.} = proc connectToPeer(net: RaftNetwork, peerId: string) {.async.} =
## Dial a peer with a short timeout so a dead peer cannot stall the whole
## heartbeat / RequestVote fan-out (default TCP connect can hang for many
## seconds, which lets live followers trip their election timers).
if peerId notin net.node.peerAddrs: if peerId notin net.node.peerAddrs:
return return
let (host, port) = net.node.peerAddrs[peerId] let (host, port) = net.node.peerAddrs[peerId]
var sock: AsyncSocket = nil
try: try:
let sock = newAsyncSocket() sock = newAsyncSocket()
await sock.connect(host, Port(port)) let ok = await withTimeout(sock.connect(host, Port(port)), RaftConnectTimeoutMs)
if not ok:
sock.close()
return
net.peerSockets[peerId] = sock net.peerSockets[peerId] = sock
except CatchableError: except CatchableError:
discard if sock != nil:
try: sock.close() except CatchableError: discard
proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} = proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} =
if peerId notin net.peerSockets: if peerId notin net.peerSockets:
@@ -582,6 +592,7 @@ proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} =
try: try:
await net.peerSockets[peerId].send(cast[string](header) & cast[string](data)) await net.peerSockets[peerId].send(cast[string](header) & cast[string](data))
except CatchableError: except CatchableError:
try: net.peerSockets[peerId].close() except CatchableError: discard
net.peerSockets.del(peerId) net.peerSockets.del(peerId)
proc broadcast*(net: RaftNetwork, msgs: seq[RaftMessage]) {.async.} = proc broadcast*(net: RaftNetwork, msgs: seq[RaftMessage]) {.async.} =
@@ -643,11 +654,19 @@ proc receiveLoop(net: RaftNetwork, client: AsyncSocket) {.async.} =
client.close() client.close()
proc heartbeatLoop(net: RaftNetwork) {.async.} = proc heartbeatLoop(net: RaftNetwork) {.async.} =
## Fan out heartbeats in parallel so a slow/dead peer cannot delay
## AppendEntries to the rest of the cluster.
while net.running: while net.running:
if net.node.state == rsLeader: if net.node.state == rsLeader:
var futs: seq[Future[void]] = @[]
for peer in net.node.peers: for peer in net.node.peers:
let msg = net.node.appendEntries(peer) let msg = net.node.appendEntries(peer)
await net.send(peer, msg) futs.add(net.send(peer, msg))
for f in futs:
try:
await f
except CatchableError:
discard
await sleepAsync(net.node.heartbeatTimeout) await sleepAsync(net.node.heartbeatTimeout)
proc timerLoop*(net: RaftNetwork) {.async.} proc timerLoop*(net: RaftNetwork) {.async.}
+7 -4
View File
@@ -156,9 +156,11 @@ proc migrateLegacyData(registry: DatabaseRegistry, config: BaraConfig) =
moveDir(legacyDir, legacyDir & ".migrated") moveDir(legacyDir, legacyDir & ".migrated")
info("Legacy data migration complete. Original renamed to " & legacyDir & ".migrated") info("Legacy data migration complete. Original renamed to " & legacyDir & ".migrated")
proc runTcpServer(config: BaraConfig) {.async.} = proc runTcpServer(server: Server) {.async.} =
info("BaraDB TCP listening on " & config.address & ":" & $config.port) ## Run the already-wired TCP Server. Must use the same instance that main
var server = newServer(config) ## attaches raftNode / replication / gossip to — a fresh newServer(config)
## would leave those fields nil and open a second registry over the same data.
info("BaraDB TCP listening on " & server.config.address & ":" & $server.config.port)
await server.run() await server.run()
proc wireRaftDistTxn(raftNode: RaftNode, tcpServer: Server) = proc wireRaftDistTxn(raftNode: RaftNode, tcpServer: Server) =
@@ -381,7 +383,8 @@ proc main() =
info("Joined gossip cluster via seed " & host & ":" & $port) info("Joined gossip cluster via seed " & host & ":" & $port)
# Start TCP wire protocol server on main thread with async event loop # Start TCP wire protocol server on main thread with async event loop
waitFor runTcpServer(config) # (must be the wired tcpServer — not a fresh newServer(config)).
waitFor runTcpServer(tcpServer)
# Shutdown: stop listeners first, then close storage under the gate # Shutdown: stop listeners first, then close storage under the gate
httpServer.stop(closeStorage = false) httpServer.stop(closeStorage = false)
+340
View File
@@ -0,0 +1,340 @@
## Raft replicated writes E2E — real 3-node cluster over the TCP transport.
## Starts three actual build/baradadb processes, writes through the leader
## (exercising the raft commit wait), verifies followers apply committed
## entries, rejects follower writes, and writes again after failover.
## Process-management conventions follow tests/raft_e2e_test.nim; client
## access follows tests/nimforum_smoke_test.nim.
##
## NOTE: CREATE TABLE produces no keyValuePairs (schema DDL is out of scope
## for raft replication), so followers never learn the table from the log.
## The test therefore creates the table locally on ALL nodes — CREATE TABLE
## is not classified as a raft write (isWrite covers DML only), so followers
## accept it — and tests ROW replication only.
import std/unittest
import std/osproc
import std/os
import std/strtabs
import std/strutils
import std/sequtils
import std/times
import std/net
import std/posix
import ../adaptors/nim/baradb_sqlite as sqlite
const
BinaryPath = "./build/baradadb"
LeaderMarker = "became leader"
type
NodeProc = object
id: string
clientPort: int
p: Process
dataDir: string
output: string
alive: bool
proc drainOutput(n: var NodeProc) =
## Reads whatever the child has written so far. The pipe was set O_NONBLOCK
## at start, so this never blocks — a hung read is impossible here.
var tmp: array[8192, char]
while true:
let count = posix.read(n.p.outputHandle.cint, tmp[0].addr, tmp.len)
if count <= 0: break
for i in 0 ..< count: n.output.add tmp[i]
proc drainAll(nodes: var seq[NodeProc]) =
for n in nodes.mitems:
if n.p != nil: n.drainOutput()
proc dumpAll(nodes: var seq[NodeProc]) =
## Debuggability: on failure, everything the nodes said.
nodes.drainAll()
for n in nodes:
echo "===== output of ", n.id, " (port ", n.clientPort, ") ====="
echo n.output
proc portOpen(port: int): bool =
var s: Socket
try:
s = newSocket()
s.connect("127.0.0.1", Port(port), timeout = 250)
s.close()
result = true
except CatchableError:
if s != nil: s.close()
result = false
proc killNode(n: var NodeProc) =
if n.p != nil and n.alive:
try:
n.p.terminate()
discard n.p.waitForExit()
except CatchableError:
discard
n.alive = false
proc leaderTerms(output: string): seq[int] =
## All terms this node logged leadership for ("became leader for term T").
var pos = 0
while true:
let idx = output.find(LeaderMarker, pos)
if idx < 0: break
let tIdx = output.find("term ", idx)
if tIdx < 0: break
let numStart = tIdx + 5
var numEnd = numStart
while numEnd < output.len and output[numEnd] in Digits: inc numEnd
if numEnd > numStart:
result.add(parseInt(output[numStart ..< numEnd]))
pos = numEnd
proc maxLeader(nodes: seq[NodeProc]): tuple[idx, term: int] =
## Node that logged leadership for the highest term seen so far.
result = (-1, 0)
for i in 0 ..< nodes.len:
for t in leaderTerms(nodes[i].output):
if t > result.term: result = (i, t)
proc drainFor(nodes: var seq[NodeProc], ms: int) =
let start = getTime()
while getTime() - start < initDuration(milliseconds = ms):
nodes.drainAll()
sleep(50)
proc openClient(port: int): DbConn =
## Connect with retries — the port may accept TCP before the DB is usable.
for i in 0 ..< 50:
try:
return open("127.0.0.1:" & $port, "", "", "default")
except CatchableError:
sleep(100)
raise newException(IOError, "cannot connect to port " & $port)
proc waitForRow(port: int, name: string, deadlineSec: int): bool =
## Poll SELECT on `port` until a row with `name` appears. Tolerates errors
## (e.g. "unknown table" while schema has not been created yet) by retrying.
let db = openClient(port)
defer: db.close()
let start = getTime()
while getTime() - start < initDuration(seconds = deadlineSec):
try:
let rows = db.getAllRows(sql"SELECT * FROM rw_test")
for row in rows:
if row.len >= 2 and row[1] == name:
return true
except CatchableError:
discard
sleep(100)
return false
proc runWritesScenario() =
## Fatal phase failures dump all captured node output, record a test
## failure, and return; cleanup happens in the finally below either way.
let tstamp = getTime().toUnix.int
# Port bases: distinct from nimforum_smoke_test (35000+mod10000) and
# raft_e2e_test (41000+mod5000). Client ports are spaced by 10 because the
# server derives HTTP (port+440), WS (port+441) and gossip (raftPort+100)
# ports — consecutive client ports collide.
let cbase = 46000 + (tstamp mod 4000)
let rbase = cbase + 100
let peers = "n1@127.0.0.1:" & $(rbase + 1) &
",n2@127.0.0.1:" & $(rbase + 2) &
",n3@127.0.0.1:" & $(rbase + 3)
var nodes: seq[NodeProc]
for i in 1 .. 3:
let id = "n" & $i
let dataDir = getTempDir() / "baradb_raft_writes_e2e_" & $tstamp & "_" & id
createDir(dataDir)
var env = newStringTable()
for key, val in envPairs():
env[key] = val
env["BARADB_PORT"] = $(cbase + i * 10)
env["BARADB_RAFT_ENABLED"] = "true"
env["BARADB_RAFT_PORT"] = $(rbase + i)
env["BARADB_RAFT_NODE_ID"] = id
env["BARADB_RAFT_PEERS"] = peers
env["BARADB_DATA_DIR"] = dataDir
env["BARADB_LOG_LEVEL"] = "info"
let p = startProcess(BinaryPath, env = env,
options = {poStdErrToStdOut, poDaemon})
discard fcntl(p.outputHandle.cint, F_SETFL,
fcntl(p.outputHandle.cint, F_GETFL) or O_NONBLOCK)
nodes.add NodeProc(id: id, clientPort: cbase + i * 10, p: p,
dataDir: dataDir, alive: true)
try:
# Readiness: all three client ports accept TCP connections (10s each).
for i in 0 ..< nodes.len:
let readyStart = getTime()
var ok = false
while getTime() - readyStart < initDuration(seconds = 10):
if portOpen(nodes[i].clientPort):
ok = true
break
sleep(100)
if not ok:
echo "node ", nodes[i].id, " never became ready"
dumpAll(nodes)
fail()
return
# Election: timeouts are 150-300ms, heartbeat 50ms — a leader should
# emerge within ~2s; 10s deadline for margin.
var elected = false
let electStart = getTime()
while getTime() - electStart < initDuration(seconds = 10):
nodes.drainAll()
if maxLeader(nodes).idx >= 0:
elected = true
break
sleep(50)
if not elected:
echo "no leader elected within 10s"
dumpAll(nodes)
fail()
return
# Settle and require stability, same as raft_e2e_test.
nodes.drainFor(2000)
let (leaderIdx, leaderTerm) = maxLeader(nodes)
nodes.drainFor(1000)
let (stableIdx, stableTerm) = maxLeader(nodes)
if stableIdx != leaderIdx or stableTerm != leaderTerm:
echo "cluster unstable: leadership moved from ", nodes[leaderIdx].id,
" (term ", leaderTerm, ") to ", nodes[stableIdx].id,
" (term ", stableTerm, ")"
dumpAll(nodes)
fail()
return
echo "leader elected: ", nodes[leaderIdx].id, " (term ", leaderTerm, ")"
# Schema: CREATE TABLE is not a raft write (DML only is), and its _schema
# keys are not replicated — create the table locally on every node.
for i in 0 ..< nodes.len:
let db = openClient(nodes[i].clientPort)
try:
db.exec(sql"CREATE TABLE rw_test (id INT PRIMARY KEY, name STRING)")
except CatchableError as e:
echo "CREATE TABLE failed on ", nodes[i].id, ": ", e.msg
dumpAll(nodes)
fail()
return
db.close()
let followerIdx = (if leaderIdx == 0: 1 else: 0)
# Leader write: INSERT goes through the raft log and waits for majority
# commit before responding (Task 2) — expect success.
block:
let db = openClient(nodes[leaderIdx].clientPort)
try:
db.exec(sql"INSERT INTO rw_test (id, name) VALUES (1, 'raft-row')")
except CatchableError as e:
echo "leader INSERT failed: ", e.msg
dumpAll(nodes)
fail()
return
db.close()
echo "leader INSERT committed"
# Follower visibility: applyCommand puts committed entries into the
# follower's default DB — poll until the row shows up (5s deadline).
if not waitForRow(nodes[followerIdx].clientPort, "raft-row", 5):
echo "follower ", nodes[followerIdx].id,
" never saw the replicated row within 5s"
dumpAll(nodes)
fail()
return
echo "row replicated to follower ", nodes[followerIdx].id
# Follower rejection: DML on a follower must fail with "not leader".
block:
let db = openClient(nodes[followerIdx].clientPort)
var rejected = false
try:
db.exec(sql"INSERT INTO rw_test (id, name) VALUES (99, 'nope')")
except CatchableError as e:
rejected = "not leader" in e.msg
if not rejected:
echo "follower INSERT failed but without 'not leader': ", e.msg
db.close()
if not rejected:
echo "follower INSERT was not rejected with a 'not leader' error"
dumpAll(nodes)
fail()
return
echo "follower INSERT rejected with 'not leader'"
# Failover: kill the leader. A survivor must accept a write once it wins
# a new term (majority of the remaining 2-of-3). Log lines can thrash
# across terms, so discover the new leader by probing INSERT rather than
# relying solely on the first "became leader" line.
killNode(nodes[leaderIdx])
var newLeaderIdx = -1
var writeErr = ""
let foStart = getTime()
while getTime() - foStart < initDuration(seconds = 15):
nodes.drainAll()
# Prefer the highest-term survivor when choosing who to probe first.
var order: seq[int] = @[]
var bestTerm = 0
var bestIdx = -1
for i in 0 ..< nodes.len:
if i == leaderIdx: continue
order.add(i)
for t in leaderTerms(nodes[i].output):
if t > bestTerm: (bestIdx, bestTerm) = (i, t)
if bestIdx >= 0:
# Probe the current highest-term node first.
order = order.filterIt(it != bestIdx)
order.insert(bestIdx, 0)
for i in order:
try:
let db = openClient(nodes[i].clientPort)
try:
db.exec(sql"INSERT INTO rw_test (id, name) VALUES (2, 'after-failover')")
newLeaderIdx = i
finally:
db.close()
if newLeaderIdx >= 0: break
except CatchableError as e:
writeErr = e.msg
# "not leader" / commit timeout / connection blips — keep probing.
if newLeaderIdx >= 0: break
sleep(150)
if newLeaderIdx < 0:
echo "no survivor accepted a post-failover write within 15s",
(if writeErr.len > 0: " (last error: " & writeErr & ")" else: "")
dumpAll(nodes)
fail()
return
echo "failover complete: new leader ", nodes[newLeaderIdx].id,
" accepted post-failover write"
# The remaining follower (neither old nor new leader) must see the row.
let remainingIdx = 3 - leaderIdx - newLeaderIdx
if not waitForRow(nodes[remainingIdx].clientPort, "after-failover", 8):
echo "remaining follower ", nodes[remainingIdx].id,
" never saw the post-failover row within 8s"
dumpAll(nodes)
fail()
return
echo "post-failover row replicated to ", nodes[remainingIdx].id
check leaderIdx != newLeaderIdx
finally:
for n in nodes.mitems:
n.killNode()
if n.p != nil: n.p.close()
removeDir(n.dataDir)
suite "Raft replicated writes E2E":
test "writes replicate, followers reject, failover resumes writes":
if not fileExists(BinaryPath):
echo "[SKIP] ", BinaryPath, " missing — run `nimble test` (builds the server first)"
skip()
else:
runWritesScenario()