feat(raft): parse id@host:port peers, enable raft state persistence
This commit is contained in:
@@ -1,6 +1,7 @@
|
|||||||
import std/os
|
import std/os
|
||||||
import std/strutils
|
import std/strutils
|
||||||
import std/json
|
import std/json
|
||||||
|
import std/tables
|
||||||
|
|
||||||
type
|
type
|
||||||
BaraConfig* = object
|
BaraConfig* = object
|
||||||
@@ -38,6 +39,7 @@ type
|
|||||||
raftPort*: int
|
raftPort*: int
|
||||||
raftPeers*: seq[string]
|
raftPeers*: seq[string]
|
||||||
raftNodeId*: string
|
raftNodeId*: string
|
||||||
|
raftPeerAddrs*: Table[string, tuple[host: string, port: int]]
|
||||||
|
|
||||||
CompactionStrategy* = enum
|
CompactionStrategy* = enum
|
||||||
csSizeTiered = "size_tiered"
|
csSizeTiered = "size_tiered"
|
||||||
@@ -76,6 +78,7 @@ proc defaultConfig*(): BaraConfig =
|
|||||||
raftPort: 9473,
|
raftPort: 9473,
|
||||||
raftPeers: @[],
|
raftPeers: @[],
|
||||||
raftNodeId: "",
|
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)
|
cfg.raftPort = parseEnvInt(getEnv("BARADB_RAFT_PORT", ""), cfg.raftPort)
|
||||||
let peersEnv = getEnv("BARADB_RAFT_PEERS", "")
|
let peersEnv = getEnv("BARADB_RAFT_PEERS", "")
|
||||||
if peersEnv.len > 0:
|
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)
|
cfg.raftNodeId = getEnv("BARADB_RAFT_NODE_ID", cfg.raftNodeId)
|
||||||
|
|
||||||
# ----------------------------------------------------------------------
|
# ----------------------------------------------------------------------
|
||||||
|
|||||||
+5
-1
@@ -330,7 +330,11 @@ proc main() =
|
|||||||
# Start Raft cluster if enabled
|
# Start Raft cluster if enabled
|
||||||
if config.raftEnabled:
|
if config.raftEnabled:
|
||||||
info("Starting Raft node " & config.raftNodeId & " on port " & $config.raftPort)
|
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
|
# Wire state machine to apply committed entries to the default database
|
||||||
let defaultDbInfo = getDatabaseInfo(registry, "default")
|
let defaultDbInfo = getDatabaseInfo(registry, "default")
|
||||||
raftNode.applyCommand = proc(cmd: string, data: seq[byte]) {.gcsafe.} =
|
raftNode.applyCommand = proc(cmd: string, data: seq[byte]) {.gcsafe.} =
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import std/os
|
|||||||
import std/tables
|
import std/tables
|
||||||
import ../src/barabadb/query/[parser, executor, lexer, ast]
|
import ../src/barabadb/query/[parser, executor, lexer, ast]
|
||||||
import ../src/barabadb/core/types
|
import ../src/barabadb/core/types
|
||||||
|
import ../src/barabadb/core/config
|
||||||
import ../src/barabadb/storage/lsm
|
import ../src/barabadb/storage/lsm
|
||||||
|
|
||||||
const testDir = "/tmp/baradb_bugfix_test"
|
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')"))
|
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)"))
|
let c = executeQuery(ctx, parse("CREATE UNIQUE INDEX accts_email ON accts (email)"))
|
||||||
check not c.success
|
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
|
||||||
|
|||||||
Reference in New Issue
Block a user