From 0bebec08cd89039c32bd9b9e73d80d573b6bf0b3 Mon Sep 17 00:00:00 2001 From: Paul Buetow Date: Fri, 27 Mar 2026 13:30:14 +0200 Subject: sr: fix Raft replay leader election --- src/main/java/protocols/VSAbstractProtocol.java | 31 ++++++++++++++-- .../protocols/implementations/VSRaftProtocol.java | 14 ++++++++ .../java/simulator/builder/SimulationFactory.java | 8 +++-- .../implementations/VSRaftProtocolTest.java | 41 ++++++++++++++++++++++ 4 files changed, 89 insertions(+), 5 deletions(-) (limited to 'src') diff --git a/src/main/java/protocols/VSAbstractProtocol.java b/src/main/java/protocols/VSAbstractProtocol.java index da12d31..ee0d6c2 100644 --- a/src/main/java/protocols/VSAbstractProtocol.java +++ b/src/main/java/protocols/VSAbstractProtocol.java @@ -180,8 +180,8 @@ abstract public class VSAbstractProtocol extends VSAbstractEvent { * This method: * * * @param message the received message @@ -192,6 +192,21 @@ abstract public class VSAbstractProtocol extends VSAbstractEvent { if (isIncorrectProtocol(message)) return; + if (isServer && isClient) { + if (message.isServerMessage()) { + currentContextIsServer(false); + if (!isClientInitialized) + onInit(); + onClientRecv(message); + } else { + currentContextIsServer(true); + if (!isServerInitialized) + onInit(); + onServerRecv(message); + } + return; + } + if (isServer) { currentContextIsServer(true); if (!isServerInitialized) @@ -226,6 +241,18 @@ abstract public class VSAbstractProtocol extends VSAbstractEvent { if (isIncorrectProtocol(message)) return false; + return isRelevantMessageForContext(message); + } + + /** + * Checks whether a message is relevant for this protocol instance. + * Subclasses can relax or specialize the default server/client routing + * rules while keeping the protocol-name filter intact. + * + * @param message the message to check + * @return true if the message should be processed by this protocol instance + */ + protected boolean isRelevantMessageForContext(VSMessage message) { if (message.isServerMessage()) { if (!isClient) return false; diff --git a/src/main/java/protocols/implementations/VSRaftProtocol.java b/src/main/java/protocols/implementations/VSRaftProtocol.java index c75628d..bad893c 100644 --- a/src/main/java/protocols/implementations/VSRaftProtocol.java +++ b/src/main/java/protocols/implementations/VSRaftProtocol.java @@ -103,6 +103,11 @@ public class VSRaftProtocol extends VSAbstractProtocol { handleMessage(recvMessage); } + @Override + protected boolean isRelevantMessageForContext(VSMessage message) { + return isServer() || isClient(); + } + /* (non-Javadoc) * @see protocols.VSAbstractProtocol#onServerSchedule() */ @@ -366,7 +371,10 @@ public class VSRaftProtocol extends VSAbstractProtocol { heartbeatAck.setInteger("term", currentTerm); heartbeatAck.setInteger("pid", process.getProcessID()); heartbeatAck.setInteger("targetPid", messageLeaderId); + boolean previousContextIsServer = currentContextIsServer(); + currentContextIsServer(true); sendMessage(heartbeatAck); + currentContextIsServer(previousContextIsServer); } /** @@ -417,7 +425,10 @@ public class VSRaftProtocol extends VSAbstractProtocol { voteResponse.setInteger("pid", process.getProcessID()); voteResponse.setBoolean("voteGranted", voteGranted); voteResponse.setInteger("targetPid", candidateId); + boolean previousContextIsServer = currentContextIsServer(); + currentContextIsServer(true); sendMessage(voteResponse); + currentContextIsServer(previousContextIsServer); } /** @@ -487,7 +498,10 @@ public class VSRaftProtocol extends VSAbstractProtocol { appendAck.setInteger("pid", process.getProcessID()); appendAck.setInteger("logIndex", messageLogIndex); appendAck.setInteger("targetPid", messageLeaderId); + boolean previousContextIsServer = currentContextIsServer(); + currentContextIsServer(true); sendMessage(appendAck); + currentContextIsServer(previousContextIsServer); } /** diff --git a/src/main/java/simulator/builder/SimulationFactory.java b/src/main/java/simulator/builder/SimulationFactory.java index 48ec638..165e42d 100644 --- a/src/main/java/simulator/builder/SimulationFactory.java +++ b/src/main/java/simulator/builder/SimulationFactory.java @@ -96,10 +96,12 @@ public class SimulationFactory { .activateClientsAt(100, 1) .activateClientsAt(1700, 2) // Bias process 1 toward a fast, clean post-crash election while - // keeping process 2's timeout comfortably behind it. - .setProtocolLong(1, "electionTimeout", 4000) + // keeping process 2's timeout comfortably behind it. The shorter + // timeout also keeps the headless replay active long enough to + // reach the first post-crash election deterministically. + .setProtocolLong(1, "electionTimeout", 2500) .setProtocolLong(1, "electionJitter", 0) - .setProtocolLong(2, "electionTimeout", 9000) + .setProtocolLong(2, "electionTimeout", 12000) .setProtocolLong(2, "electionJitter", 0) .addCrashEvent(0, 3500); } diff --git a/src/test/java/protocols/implementations/VSRaftProtocolTest.java b/src/test/java/protocols/implementations/VSRaftProtocolTest.java index 40dfd10..802c80f 100644 --- a/src/test/java/protocols/implementations/VSRaftProtocolTest.java +++ b/src/test/java/protocols/implementations/VSRaftProtocolTest.java @@ -389,6 +389,33 @@ class VSRaftProtocolTest { assertEquals(4750L, taskCaptor.getValue().getTaskTime()); } + @Test + void testDualRoleVoteRequestIsHandledOnce() throws Exception { + protocol.currentContextIsServer(false); + protocol.isClient(true); + protocol.isServer(true); + protocol.onInit(); + clearInvocations(mockProcess, mockTaskManager); + when(mockProcess.getTime()).thenReturn(200L, 200L); + + VSMessage voteRequest = new VSMessage(); + setMessageProtocolClassname(voteRequest, VSRaftProtocol.class.getName()); + setMessageServerFlag(voteRequest, true); + voteRequest.setString("type", "voteRequest"); + voteRequest.setInteger("term", 2); + voteRequest.setInteger("candidateId", 11); + + ArgumentCaptor messageCaptor = + ArgumentCaptor.forClass(VSMessage.class); + + protocol.onMessageRecvStart(voteRequest); + + verify(mockProcess).sendMessage(messageCaptor.capture()); + assertEquals("voteResponse", messageCaptor.getValue().getString("type")); + assertEquals(2, messageCaptor.getValue().getInteger("term")); + assertTrue(messageCaptor.getValue().getBoolean("voteGranted")); + } + @Test void testClientReceiveHeartbeatBecomesFollowerResetsTimeoutAndSendsAck() throws Exception { @@ -764,6 +791,20 @@ class VSRaftProtocolTest { field.setBoolean(protocol, value); } + private void setMessageProtocolClassname(VSMessage message, String classname) + throws Exception { + Field field = VSMessage.class.getDeclaredField("protocolClassname"); + field.setAccessible(true); + field.set(message, classname); + } + + private void setMessageServerFlag(VSMessage message, boolean isServerMessage) + throws Exception { + Field field = VSMessage.class.getDeclaredField("isServerMessage"); + field.setAccessible(true); + field.setBoolean(message, isServerMessage); + } + private int getIntField(String fieldName) throws Exception { Field field = VSRaftProtocol.class.getDeclaredField(fieldName); field.setAccessible(true); -- cgit v1.2.3