feat(raft): run election timer in production, reset on AppendEntries
This commit is contained in:
@@ -548,12 +548,14 @@ type
|
|||||||
socket*: AsyncSocket
|
socket*: AsyncSocket
|
||||||
running*: bool
|
running*: bool
|
||||||
peerSockets*: Table[string, AsyncSocket]
|
peerSockets*: Table[string, AsyncSocket]
|
||||||
|
timer*: ElectionTimer
|
||||||
|
|
||||||
proc newRaftNetwork*(node: RaftNode): RaftNetwork =
|
proc newRaftNetwork*(node: RaftNode): RaftNetwork =
|
||||||
RaftNetwork(
|
RaftNetwork(
|
||||||
node: node,
|
node: node,
|
||||||
running: false,
|
running: false,
|
||||||
peerSockets: initTable[string, AsyncSocket](),
|
peerSockets: initTable[string, AsyncSocket](),
|
||||||
|
timer: newElectionTimer(node, node.electionTimeout),
|
||||||
)
|
)
|
||||||
|
|
||||||
proc connectToPeer(net: RaftNetwork, peerId: string) {.async.} =
|
proc connectToPeer(net: RaftNetwork, peerId: string) {.async.} =
|
||||||
@@ -585,7 +587,7 @@ proc broadcast*(net: RaftNetwork, msgs: seq[RaftMessage]) {.async.} =
|
|||||||
if i < msgs.len:
|
if i < msgs.len:
|
||||||
await net.send(peer, msgs[i])
|
await net.send(peer, msgs[i])
|
||||||
|
|
||||||
proc processMessage(net: RaftNetwork, msg: RaftMessage) {.async.} =
|
proc processMessage*(net: RaftNetwork, msg: RaftMessage) {.async.} =
|
||||||
case msg.kind
|
case msg.kind
|
||||||
of rmkRequestVote:
|
of rmkRequestVote:
|
||||||
let reply = net.node.handleRequestVote(msg)
|
let reply = net.node.handleRequestVote(msg)
|
||||||
@@ -593,6 +595,11 @@ proc processMessage(net: RaftNetwork, msg: RaftMessage) {.async.} =
|
|||||||
of rmkRequestVoteReply:
|
of rmkRequestVoteReply:
|
||||||
net.node.handleVoteReply(msg)
|
net.node.handleVoteReply(msg)
|
||||||
of rmkAppendEntries:
|
of rmkAppendEntries:
|
||||||
|
# A plausible current leader (same acceptance condition as
|
||||||
|
# handleAppendEntries) resets the election timer; stale-term
|
||||||
|
# messages must not.
|
||||||
|
if msg.term >= net.node.currentTerm:
|
||||||
|
net.timer.resetTimeout()
|
||||||
let reply = net.node.handleAppendEntries(msg)
|
let reply = net.node.handleAppendEntries(msg)
|
||||||
await net.send(msg.senderId, reply)
|
await net.send(msg.senderId, reply)
|
||||||
of rmkAppendEntriesReply:
|
of rmkAppendEntriesReply:
|
||||||
@@ -630,13 +637,17 @@ proc heartbeatLoop(net: RaftNetwork) {.async.} =
|
|||||||
await net.send(peer, msg)
|
await net.send(peer, msg)
|
||||||
await sleepAsync(net.node.heartbeatTimeout)
|
await sleepAsync(net.node.heartbeatTimeout)
|
||||||
|
|
||||||
|
proc timerLoop*(net: RaftNetwork) {.async.}
|
||||||
|
|
||||||
proc run*(net: RaftNetwork) {.async.} =
|
proc run*(net: RaftNetwork) {.async.} =
|
||||||
net.socket = newAsyncSocket()
|
net.socket = newAsyncSocket()
|
||||||
net.socket.setSockOpt(OptReuseAddr, true)
|
net.socket.setSockOpt(OptReuseAddr, true)
|
||||||
net.socket.bindAddr(Port(net.node.raftPort))
|
net.socket.bindAddr(Port(net.node.raftPort))
|
||||||
net.socket.listen()
|
net.socket.listen()
|
||||||
net.running = true
|
net.running = true
|
||||||
|
net.timer.resetTimeout()
|
||||||
asyncCheck net.heartbeatLoop()
|
asyncCheck net.heartbeatLoop()
|
||||||
|
asyncCheck net.timerLoop()
|
||||||
while net.running:
|
while net.running:
|
||||||
try:
|
try:
|
||||||
let client = await net.socket.accept()
|
let client = await net.socket.accept()
|
||||||
@@ -646,6 +657,7 @@ proc run*(net: RaftNetwork) {.async.} =
|
|||||||
|
|
||||||
proc stop*(net: RaftNetwork) =
|
proc stop*(net: RaftNetwork) =
|
||||||
net.running = false
|
net.running = false
|
||||||
|
net.timer.stop()
|
||||||
if net.socket != nil:
|
if net.socket != nil:
|
||||||
net.socket.close()
|
net.socket.close()
|
||||||
for peerId, sock in net.peerSockets:
|
for peerId, sock in net.peerSockets:
|
||||||
@@ -683,3 +695,10 @@ proc tick*(timer: ElectionTimer, net: RaftNetwork = nil) =
|
|||||||
timer.resetTimeout()
|
timer.resetTimeout()
|
||||||
of rsLeader:
|
of rsLeader:
|
||||||
timer.resetTimeout() # Keep alive
|
timer.resetTimeout() # Keep alive
|
||||||
|
|
||||||
|
proc timerLoop*(net: RaftNetwork) {.async.} =
|
||||||
|
## Production election timer: ticks the node's ElectionTimer until the
|
||||||
|
## network transport is stopped.
|
||||||
|
while net.running:
|
||||||
|
tick(net.timer, net)
|
||||||
|
await sleepAsync(50)
|
||||||
|
|||||||
@@ -2390,6 +2390,75 @@ suite "Raft Network Transport":
|
|||||||
if n3.isLeader: inc leaderCount
|
if n3.isLeader: inc leaderCount
|
||||||
check leaderCount == 1
|
check leaderCount == 1
|
||||||
|
|
||||||
|
test "timerLoop elects a leader without manual ticks":
|
||||||
|
var n1 = newRaftNode("n1", @["n2", "n3"], raftPort = 29011)
|
||||||
|
var n2 = newRaftNode("n2", @["n1", "n3"], raftPort = 29012)
|
||||||
|
var n3 = newRaftNode("n3", @["n1", "n2"], raftPort = 29013)
|
||||||
|
|
||||||
|
# Distinct deterministic timeouts avoid split-vote livelock
|
||||||
|
n1.electionTimeout = 150
|
||||||
|
n2.electionTimeout = 250
|
||||||
|
n3.electionTimeout = 350
|
||||||
|
|
||||||
|
n1.peerAddrs["n2"] = ("127.0.0.1", 29012)
|
||||||
|
n1.peerAddrs["n3"] = ("127.0.0.1", 29013)
|
||||||
|
n2.peerAddrs["n1"] = ("127.0.0.1", 29011)
|
||||||
|
n2.peerAddrs["n3"] = ("127.0.0.1", 29013)
|
||||||
|
n3.peerAddrs["n1"] = ("127.0.0.1", 29011)
|
||||||
|
n3.peerAddrs["n2"] = ("127.0.0.1", 29012)
|
||||||
|
|
||||||
|
let net1 = newRaftNetwork(n1)
|
||||||
|
let net2 = newRaftNetwork(n2)
|
||||||
|
let net3 = newRaftNetwork(n3)
|
||||||
|
|
||||||
|
asyncCheck net1.run()
|
||||||
|
asyncCheck net2.run()
|
||||||
|
asyncCheck net3.run()
|
||||||
|
waitFor sleepAsync(50)
|
||||||
|
|
||||||
|
# No manual ticks — the production timerLoop must drive the election
|
||||||
|
var leaderCount = 0
|
||||||
|
var waited = 0
|
||||||
|
while waited < 3000:
|
||||||
|
leaderCount = 0
|
||||||
|
if n1.isLeader: inc leaderCount
|
||||||
|
if n2.isLeader: inc leaderCount
|
||||||
|
if n3.isLeader: inc leaderCount
|
||||||
|
if leaderCount == 1: break
|
||||||
|
waitFor sleepAsync(100)
|
||||||
|
waited += 100
|
||||||
|
|
||||||
|
net1.stop()
|
||||||
|
net2.stop()
|
||||||
|
net3.stop()
|
||||||
|
waitFor sleepAsync(50)
|
||||||
|
|
||||||
|
check leaderCount == 1
|
||||||
|
|
||||||
|
test "inbound AppendEntries resets the election timer":
|
||||||
|
var nodeA = newRaftNode("a", @["leader"], raftPort = 0)
|
||||||
|
nodeA.electionTimeout = 0 # any elapsed millisecond counts as a timeout
|
||||||
|
let netA = newRaftNetwork(nodeA)
|
||||||
|
let timer = netA.timer
|
||||||
|
|
||||||
|
waitFor sleepAsync(5)
|
||||||
|
check timer.checkTimeout() # currently timed out
|
||||||
|
|
||||||
|
# Valid AppendEntries from a plausible current leader (term >= currentTerm)
|
||||||
|
let msg = RaftMessage(kind: rmkAppendEntries, term: nodeA.currentTerm,
|
||||||
|
senderId: "leader")
|
||||||
|
waitFor netA.processMessage(msg)
|
||||||
|
check not timer.checkTimeout()
|
||||||
|
|
||||||
|
# Stale-term AppendEntries must NOT reset the timer
|
||||||
|
nodeA.currentTerm = 5
|
||||||
|
waitFor sleepAsync(5)
|
||||||
|
check timer.checkTimeout()
|
||||||
|
let staleMsg = RaftMessage(kind: rmkAppendEntries, term: 1,
|
||||||
|
senderId: "leader")
|
||||||
|
waitFor netA.processMessage(staleMsg)
|
||||||
|
check timer.checkTimeout()
|
||||||
|
|
||||||
suite "CLI Autocomplete":
|
suite "CLI Autocomplete":
|
||||||
test "Autocomplete commands":
|
test "Autocomplete commands":
|
||||||
let res = autocomplete("he")
|
let res = autocomplete("he")
|
||||||
|
|||||||
Reference in New Issue
Block a user