This is an automated email from the ASF dual-hosted git repository.

anmolnar pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/zookeeper.git

commit 22dad52e25d028c9bd5c2ea7b525b56ce1da1a53
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 61c0eb6010..8e238e5926 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 0000000000..67ab4b083e
--- /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");
+    }
+
+}

Reply via email to