From 44060701b76b466831992c5c7e0dc5e593611cb8 Mon Sep 17 00:00:00 2001 From: dimgigov Date: Thu, 30 Jul 2026 21:06:49 +0300 Subject: [PATCH] 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. --- .gitignore | 1 + baradadb.nimble | 3 +- src/barabadb/core/raft.nim | 27 ++- src/baradadb.nim | 11 +- tests/raft_writes_e2e_test.nim | 340 +++++++++++++++++++++++++++++++++ 5 files changed, 373 insertions(+), 9 deletions(-) create mode 100644 tests/raft_writes_e2e_test.nim diff --git a/.gitignore b/.gitignore index 1a949a7..232eefc 100644 --- a/.gitignore +++ b/.gitignore @@ -15,6 +15,7 @@ tests/prop_test tests/bugfix_test tests/nimforum_smoke_test tests/raft_e2e_test +tests/raft_writes_e2e_test benchmarks/bench_all benchmarks/compare clients/nim/tests/test_client diff --git a/baradadb.nimble b/baradadb.nimble index e6f4a10..44d636d 100644 --- a/baradadb.nimble +++ b/baradadb.nimble @@ -29,7 +29,8 @@ task test, "Run all tests": # Quick embedded suites first, heavy fuzz/stress suites last. for t in ["test_minimal", "test_all", "bugfix_test", "join_tests", "test_lock", "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"]: exec "nim c -r tests/" & t & ".nim" diff --git a/src/barabadb/core/raft.nim b/src/barabadb/core/raft.nim index 55c54b3..2cc67ed 100644 --- a/src/barabadb/core/raft.nim +++ b/src/barabadb/core/raft.nim @@ -560,16 +560,26 @@ proc newRaftNetwork*(node: RaftNode): RaftNetwork = timer: newElectionTimer(node, node.electionTimeout), ) +const RaftConnectTimeoutMs = 200 + 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: return let (host, port) = net.node.peerAddrs[peerId] + var sock: AsyncSocket = nil try: - let sock = newAsyncSocket() - await sock.connect(host, Port(port)) + sock = newAsyncSocket() + let ok = await withTimeout(sock.connect(host, Port(port)), RaftConnectTimeoutMs) + if not ok: + sock.close() + return net.peerSockets[peerId] = sock except CatchableError: - discard + if sock != nil: + try: sock.close() except CatchableError: discard proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} = if peerId notin net.peerSockets: @@ -582,6 +592,7 @@ proc send*(net: RaftNetwork, peerId: string, msg: RaftMessage) {.async.} = try: await net.peerSockets[peerId].send(cast[string](header) & cast[string](data)) except CatchableError: + try: net.peerSockets[peerId].close() except CatchableError: discard net.peerSockets.del(peerId) proc broadcast*(net: RaftNetwork, msgs: seq[RaftMessage]) {.async.} = @@ -643,11 +654,19 @@ proc receiveLoop(net: RaftNetwork, client: AsyncSocket) {.async.} = client.close() 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: if net.node.state == rsLeader: + var futs: seq[Future[void]] = @[] for peer in net.node.peers: 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) proc timerLoop*(net: RaftNetwork) {.async.} diff --git a/src/baradadb.nim b/src/baradadb.nim index d198621..873288a 100644 --- a/src/baradadb.nim +++ b/src/baradadb.nim @@ -156,9 +156,11 @@ proc migrateLegacyData(registry: DatabaseRegistry, config: BaraConfig) = moveDir(legacyDir, legacyDir & ".migrated") info("Legacy data migration complete. Original renamed to " & legacyDir & ".migrated") -proc runTcpServer(config: BaraConfig) {.async.} = - info("BaraDB TCP listening on " & config.address & ":" & $config.port) - var server = newServer(config) +proc runTcpServer(server: Server) {.async.} = + ## Run the already-wired TCP Server. Must use the same instance that main + ## 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() proc wireRaftDistTxn(raftNode: RaftNode, tcpServer: Server) = @@ -381,7 +383,8 @@ proc main() = info("Joined gossip cluster via seed " & host & ":" & $port) # 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 httpServer.stop(closeStorage = false) diff --git a/tests/raft_writes_e2e_test.nim b/tests/raft_writes_e2e_test.nim new file mode 100644 index 0000000..c6886ff --- /dev/null +++ b/tests/raft_writes_e2e_test.nim @@ -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()