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
The following commit(s) were added to refs/heads/master by this push:
new 5143ab88c Revert "Check valid voting member before applying (re)config
during LE"
5143ab88c is described below
commit 5143ab88c943675fe9c3e9c7a3c1049b615f6959
Author: Andor Molnar <[email protected]>
AuthorDate: Wed Sep 30 12:20:25 2026 -0500
Revert "Check valid voting member before applying (re)config during LE"
This reverts commit 22dad52e25d028c9bd5c2ea7b525b56ce1da1a53.
---
.../server/quorum/FastLeaderElection.java | 60 +++-----
.../server/quorum/FLEConfigFromNonVoterTest.java | 151 ---------------------
2 files changed, 22 insertions(+), 189 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 8e238e592..61c0eb601 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,47 +301,31 @@ public void run() {
byte[] b = new byte[configLength];
response.buffer.get(b);
- /*
- * 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());
+ 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());
}
- } 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
deleted file mode 100644
index 67ab4b083..000000000
---
a/zookeeper-server/src/test/java/org/apache/zookeeper/server/quorum/FLEConfigFromNonVoterTest.java
+++ /dev/null
@@ -1,151 +0,0 @@
-/*
- * 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");
- }
-
-}