This is an automated email from the ASF dual-hosted git repository. anmolnar pushed a commit to branch branch-3.9 in repository https://gitbox.apache.org/repos/asf/zookeeper.git
commit 509ffd01e2804c3a43a8828b5c7e2df051e97079 Author: Andor Molnar <[email protected]> AuthorDate: Wed Sep 30 10:08:35 2026 -0500 Check valid voting member before applying (re)config during LE --- .../server/quorum/FastLeaderElection.java | 60 +++++--- .../server/quorum/FLEConfigFromNonVoterTest.java | 151 +++++++++++++++++++++ 2 files changed, 189 insertions(+), 22 deletions(-) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java index cee02972e..e0a79ab53 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java @@ -301,31 +301,47 @@ public void run() { byte[] b = new byte[configLength]; response.buffer.get(b); - synchronized (self) { - try { - rqv = self.configFromString(new String(b, UTF_8)); - QuorumVerifier curQV = self.getQuorumVerifier(); - if (rqv.getVersion() > curQV.getVersion()) { - LOG.info("{} Received version: {} my version: {}", - self.getMyId(), - Long.toHexString(rqv.getVersion()), - Long.toHexString(self.getQuorumVerifier().getVersion())); - if (self.getPeerState() == ServerState.LOOKING) { - LOG.debug("Invoking processReconfig(), state: {}", self.getServerState()); - self.processReconfig(rqv, null, null, false); - if (!rqv.equals(curQV)) { - LOG.info("restarting leader election"); - self.shuttingDownLE = true; - self.getElectionAlg().shutdown(); - - break; + /* + * Only adopt a QuorumVerifier carried in a notification if + * the sender is a voting member of our current (or next) + * configuration. The election port accepts connections from + * any sid, so without this gate an unauthenticated peer could + * push an arbitrary config into processReconfig(), which + * persists it to the dynamic config file and restarts + * leader election. Non-voters (observers, joining servers) + * still get the reply notification below, which carries our + * own config, so they can learn the current membership. + */ + if (!validVoter(response.sid)) { + LOG.info("Ignoring config section in notification from non-voter sid={} (config length: {})", + response.sid, configLength); + } else { + synchronized (self) { + try { + rqv = self.configFromString(new String(b, UTF_8)); + QuorumVerifier curQV = self.getQuorumVerifier(); + if (rqv.getVersion() > curQV.getVersion()) { + LOG.info("{} Received version: {} my version: {}", + self.getMyId(), + Long.toHexString(rqv.getVersion()), + Long.toHexString(self.getQuorumVerifier().getVersion())); + if (self.getPeerState() == ServerState.LOOKING) { + LOG.debug("Invoking processReconfig(), state: {}", self.getServerState()); + self.processReconfig(rqv, null, null, false); + if (!rqv.equals(curQV)) { + LOG.info("restarting leader election"); + self.shuttingDownLE = true; + self.getElectionAlg().shutdown(); + + break; + } + } else { + LOG.debug("Skip processReconfig(), state: {}", self.getServerState()); } - } else { - LOG.debug("Skip processReconfig(), state: {}", self.getServerState()); } + } catch (IOException | ConfigException e) { + LOG.error("Something went wrong while processing config received from {}", response.sid); } - } catch (IOException | ConfigException e) { - LOG.error("Something went wrong while processing config received from {}", response.sid); } } } else { diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/FLEConfigFromNonVoterTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/FLEConfigFromNonVoterTest.java new file mode 100644 index 000000000..67ab4b083 --- /dev/null +++ b/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/FLEConfigFromNonVoterTest.java @@ -0,0 +1,151 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zookeeper.server.quorum; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import java.io.File; +import java.net.InetSocketAddress; +import java.util.HashMap; +import java.util.Map; +import org.apache.zookeeper.PortAssignment; +import org.apache.zookeeper.ZKTestCase; +import org.apache.zookeeper.server.quorum.QuorumPeer.QuorumServer; +import org.apache.zookeeper.server.quorum.QuorumPeer.ServerState; +import org.apache.zookeeper.server.quorum.flexible.QuorumVerifier; +import org.apache.zookeeper.test.ClientBase; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * A notification received on the election port may carry a QuorumVerifier + * (version >= 0x2). The election port accepts a connection from any sid, so + * the config section must only be adopted when the sender is a voting member + * of the receiver's current configuration. Otherwise an unauthenticated peer + * could inject an arbitrary quorum configuration into a LOOKING server. + */ +public class FLEConfigFromNonVoterTest extends ZKTestCase { + + protected static final Logger LOG = LoggerFactory.getLogger(FLEConfigFromNonVoterTest.class); + + private static final int COUNT = 3; + private static final long ROGUE_SID = 99L; + + private Map<Long, QuorumServer> peers; + private File[] tmpdir; + private int[] port; + private QuorumPeer victim; + private QuorumCnxManager[] cnxManagers; + private boolean savedReconfigEnabled; + + @BeforeEach + public void setUp() throws Exception { + savedReconfigEnabled = QuorumPeerConfig.isReconfigEnabled(); + QuorumPeerConfig.setReconfigEnabled(true); + + peers = new HashMap<>(); + tmpdir = new File[COUNT + 1]; + port = new int[COUNT + 1]; + cnxManagers = new QuorumCnxManager[2]; + + for (int i = 0; i < COUNT; i++) { + int clientport = PortAssignment.unique(); + peers.put((long) i, new QuorumServer(i, + new InetSocketAddress("127.0.0.1", PortAssignment.unique()), + new InetSocketAddress("127.0.0.1", PortAssignment.unique()), + new InetSocketAddress("127.0.0.1", clientport))); + tmpdir[i] = ClientBase.createTmpDir(); + port[i] = clientport; + } + tmpdir[COUNT] = ClientBase.createTmpDir(); + port[COUNT] = PortAssignment.unique(); + } + + @AfterEach + public void tearDown() throws Exception { + for (QuorumCnxManager m : cnxManagers) { + if (m != null) { + m.halt(); + } + } + if (victim != null) { + victim.shutdown(); + } + QuorumPeerConfig.setReconfigEnabled(savedReconfigEnabled); + } + + @Test + public void testConfigFromNonVoterIsIgnored() throws Exception { + victim = new QuorumPeer(peers, tmpdir[0], tmpdir[0], port[0], 3, 0, 1000, 2, 2, 2); + victim.startLeaderElection(); + QuorumVerifier originalQV = victim.getQuorumVerifier(); + long originalVersion = originalQV.getVersion(); + + FLETestUtils.LEThread thread = new FLETestUtils.LEThread(victim, 0); + thread.start(); + + // Rogue peer: not a member of the victim's view, but knows the + // victim's election address. Its sid is larger than the victim's so + // the victim's QuorumCnxManager keeps the inbound connection. + Map<Long, QuorumServer> rogueView = new HashMap<>(peers); + rogueView.put(ROGUE_SID, new QuorumServer(ROGUE_SID, + new InetSocketAddress("127.0.0.1", PortAssignment.unique()), + new InetSocketAddress("127.0.0.1", PortAssignment.unique()), + new InetSocketAddress("127.0.0.1", port[COUNT]))); + QuorumPeer rogue = new QuorumPeer(rogueView, tmpdir[COUNT], tmpdir[COUNT], port[COUNT], 3, ROGUE_SID, 1000, 2, 2, 2); + cnxManagers[0] = rogue.createCnxnManager(); + cnxManagers[0].listener.start(); + + String injected = originalQV.toString().replaceAll("version=.*", "") + + "server." + ROGUE_SID + "=127.0.0.1:12345:12346:participant;12347\n" + + "version=" + Long.toHexString(originalVersion + 1); + cnxManagers[0].toSend(0L, FastLeaderElection.buildMsg( + ServerState.LOOKING.ordinal(), 0, 0, 1, 1, injected.getBytes(UTF_8))); + + // Give the victim's WorkerReceiver time to process the message. + Thread.sleep(2000); + + assertEquals(originalVersion, victim.getQuorumVerifier().getVersion(), + "config from non-voter must not be adopted"); + assertFalse(victim.getQuorumVerifier().getAllMembers().containsKey(ROGUE_SID), + "non-voter must not be able to add itself to the quorum"); + assertTrue(thread.isAlive(), "leader election must not have been restarted by a non-voter"); + + // Positive control: the same config from a valid voter (sid 1) is adopted, + // which proves the delivery path in this test actually works. + QuorumPeer voter = new QuorumPeer(peers, tmpdir[1], tmpdir[1], port[1], 3, 1, 1000, 2, 2, 2); + cnxManagers[1] = voter.createCnxnManager(); + cnxManagers[1].listener.start(); + cnxManagers[1].toSend(0L, FastLeaderElection.buildMsg( + ServerState.LOOKING.ordinal(), 0, 0, 1, 1, injected.getBytes(UTF_8))); + + long deadline = System.currentTimeMillis() + 10000; + while (victim.getQuorumVerifier().getVersion() == originalVersion && System.currentTimeMillis() < deadline) { + Thread.sleep(100); + } + assertEquals(originalVersion + 1, victim.getQuorumVerifier().getVersion(), + "config from a valid voter should still be adopted"); + } + +}
