From 853ec7dd3ba6887f85f1dd18b4097eb1b515277b Mon Sep 17 00:00:00 2001 From: dimgigov Date: Thu, 30 Jul 2026 17:34:39 +0300 Subject: [PATCH] feat(raft): parse id@host:port peers, enable raft state persistence --- src/barabadb/core/config.nim | 29 ++++++++++++++++++++++++++++- src/baradadb.nim | 6 +++++- tests/bugfix_test.nim | 26 ++++++++++++++++++++++++++ 3 files changed, 59 insertions(+), 2 deletions(-) diff --git a/src/barabadb/core/config.nim b/src/barabadb/core/config.nim index e3bbd46..0438f2d 100644 --- a/src/barabadb/core/config.nim +++ b/src/barabadb/core/config.nim @@ -1,6 +1,7 @@ import std/os import std/strutils import std/json +import std/tables type BaraConfig* = object @@ -38,6 +39,7 @@ type raftPort*: int raftPeers*: seq[string] raftNodeId*: string + raftPeerAddrs*: Table[string, tuple[host: string, port: int]] CompactionStrategy* = enum csSizeTiered = "size_tiered" @@ -76,6 +78,7 @@ proc defaultConfig*(): BaraConfig = raftPort: 9473, raftPeers: @[], raftNodeId: "", + raftPeerAddrs: initTable[string, tuple[host: string, port: int]](), ) # ---------------------------------------------------------------------- @@ -175,7 +178,31 @@ proc loadConfigFromEnv*(cfg: var BaraConfig) = cfg.raftPort = parseEnvInt(getEnv("BARADB_RAFT_PORT", ""), cfg.raftPort) let peersEnv = getEnv("BARADB_RAFT_PEERS", "") if peersEnv.len > 0: - cfg.raftPeers = peersEnv.split(",") + cfg.raftPeers = @[] + cfg.raftPeerAddrs = initTable[string, tuple[host: string, port: int]]() + for raw in peersEnv.split(","): + let entry = raw.strip() + if entry.len == 0: continue + let atPos = entry.rfind('@') + if atPos < 0: + # bare id — no network address + cfg.raftPeers.add(entry) + else: + let id = entry[0 ..< atPos] + let hostPort = entry[atPos + 1 .. ^1] + let colonPos = hostPort.rfind(':') + let host = if colonPos >= 0: hostPort[0 ..< colonPos] else: "" + let portStr = if colonPos >= 0: hostPort[colonPos + 1 .. ^1] else: "" + var port = 0 + try: + port = parseInt(portStr) + except ValueError: + discard + if id.len == 0 or host.len == 0 or port < 1 or port > 65535: + raise newException(ValueError, + "Invalid BARADB_RAFT_PEERS entry '" & entry & "': expected id@host:port with port 1-65535") + cfg.raftPeers.add(id) + cfg.raftPeerAddrs[id] = (host, port) cfg.raftNodeId = getEnv("BARADB_RAFT_NODE_ID", cfg.raftNodeId) # ---------------------------------------------------------------------- diff --git a/src/baradadb.nim b/src/baradadb.nim index 5eb32ba..6b0a952 100644 --- a/src/baradadb.nim +++ b/src/baradadb.nim @@ -330,7 +330,11 @@ proc main() = # Start Raft cluster if enabled if config.raftEnabled: info("Starting Raft node " & config.raftNodeId & " on port " & $config.raftPort) - var raftNode = newRaftNode(config.raftNodeId, config.raftPeers, config.raftPort) + let raftDataDir = config.dataDir / "raft" + createDir(raftDataDir) # idempotent; loadState reads from it, saveState writes + var raftNode = newRaftNode(config.raftNodeId, config.raftPeers, config.raftPort, + dataDir = raftDataDir) + raftNode.peerAddrs = config.raftPeerAddrs # Wire state machine to apply committed entries to the default database let defaultDbInfo = getDatabaseInfo(registry, "default") raftNode.applyCommand = proc(cmd: string, data: seq[byte]) {.gcsafe.} = diff --git a/tests/bugfix_test.nim b/tests/bugfix_test.nim index f84885b..1f2d46d 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/core/types +import ../src/barabadb/core/config import ../src/barabadb/storage/lsm const testDir = "/tmp/baradb_bugfix_test" @@ -357,3 +358,28 @@ suite "Bug fixes — UNIQUE index enforcement": discard executeQuery(ctx, parse("INSERT INTO accts (id, email) VALUES (2, 'a@b.c')")) let c = executeQuery(ctx, parse("CREATE UNIQUE INDEX accts_email ON accts (email)")) check not c.success + +suite "Raft peer address parsing": + + test "id@host:port entries populate raftPeerAddrs": + putEnv("BARADB_RAFT_PEERS", "n1@127.0.0.1:9473,n2@10.0.0.5:9474,n3") + defer: delEnv("BARADB_RAFT_PEERS") + var cfg = defaultConfig() + loadConfigFromEnv(cfg) + check cfg.raftPeers == @["n1", "n2", "n3"] + check cfg.raftPeerAddrs["n1"] == ("127.0.0.1", 9473) + check cfg.raftPeerAddrs["n2"] == ("10.0.0.5", 9474) + check "n3" notin cfg.raftPeerAddrs + + test "malformed peer entries raise with the entry in the message": + for bad in ["n1@:9473", "n1@host:notaport", "@host:9473", "n1@host:0", "n1@host:70000"]: + putEnv("BARADB_RAFT_PEERS", bad) + var cfg = defaultConfig() + var msg = "" + try: + loadConfigFromEnv(cfg) + except ValueError as e: + msg = e.msg + delEnv("BARADB_RAFT_PEERS") + check msg.len > 0 + check bad in msg