From 077ce3c02d198389de2683f089267c8399f974c4 Mon Sep 17 00:00:00 2001 From: Kezhu Wang Date: Fri, 5 May 2023 15:53:34 +0800 Subject: [PATCH 1/2] CURATOR-593: Append chroot for EnsembleProvider::setConnectionString in EnsembleTracker Curator uses `EnsembleProvider::getConnectionString` as connection string to `ZooKeeper`. `EnsembleTracker` subscribes to config node to construct up to date connection string for `EnsembleProvider`. This is great. But, currently, `EnsembleTracker` omits chroot part of connection string which could cause curator locating at ZooKeeper root after reconnection. This could damage clients' data hierarchy in unexpected manner. --- .../framework/imps/EnsembleTracker.java | 11 +++- .../framework/imps/TestReconfiguration.java | 54 ++++++++++++++++++- 2 files changed, 62 insertions(+), 3 deletions(-) diff --git a/curator-framework/src/main/java/org/apache/curator/framework/imps/EnsembleTracker.java b/curator-framework/src/main/java/org/apache/curator/framework/imps/EnsembleTracker.java index 82ba6e2fc6..f7b99f8f27 100644 --- a/curator-framework/src/main/java/org/apache/curator/framework/imps/EnsembleTracker.java +++ b/curator-framework/src/main/java/org/apache/curator/framework/imps/EnsembleTracker.java @@ -209,10 +209,17 @@ private void processConfigData(byte[] data) throws Exception if (!properties.isEmpty()) { QuorumMaj newConfig = new QuorumMaj(properties); - String connectionString = configToConnectionString(newConfig); - if (connectionString.trim().length() > 0) + String connectionString = configToConnectionString(newConfig).trim(); + if (!connectionString.isEmpty()) { currentConfig.set(newConfig); + String oldConnectionString = ensembleProvider.getConnectionString(); + int i = oldConnectionString.indexOf('/'); + if (i >= 0) + { + String chroot = oldConnectionString.substring(i); + connectionString += chroot; + } ensembleProvider.setConnectionString(connectionString); } else diff --git a/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java b/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java index 13c162e3b6..e7854b8557 100644 --- a/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java +++ b/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java @@ -51,6 +51,7 @@ import java.net.InetAddress; import java.net.InetSocketAddress; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.Iterator; @@ -60,6 +61,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; @@ -69,6 +71,7 @@ public class TestReconfiguration extends CuratorTestBase { private final Timing2 timing = new Timing2(); private TestingCluster cluster; + private CountDownLatch ensembleLatch; private EnsembleProvider ensembleProvider; private static final String superUserPasswordDigest = "curator-test:zghsj3JfJqK7DbWf0RQ1BgbJH9w="; // ran from DigestAuthenticationProvider.generateDigest(superUserPassword); @@ -436,6 +439,51 @@ public void testNewMembers() throws Exception } } + @Test + public void testRemoveWithChroot() throws Exception + { + // Use a long chroot path to circumvent ZOOKEEPER-4565 and ZOOKEEPER-4601 + String chroot = "/pretty-long-chroot"; + + try (CuratorFramework client = newClient(cluster.getConnectString() + chroot)) { + client.start(); + client.create().forPath("/", "deadbeef".getBytes()); + + QuorumVerifier oldConfig = toQuorumVerifier(client.getConfig().forEnsemble()); + assertConfig(oldConfig, cluster.getInstances()); + + CountDownLatch latch = setChangeWaiter(client); + + Collection oldInstances = cluster.getInstances(); + InstanceSpec us = cluster.findConnectionInstance(client.getZookeeperClient().getZooKeeper()); + InstanceSpec removeSpec = oldInstances.iterator().next(); + if ( us.equals(removeSpec) ) { + Iterator iterator = oldInstances.iterator(); + iterator.next(); + removeSpec = iterator.next(); + } + + client.reconfig().leaving(Integer.toString(removeSpec.getServerId())).fromConfig(oldConfig.getVersion()).forEnsemble(); + + assertTrue(timing.awaitLatch(latch)); + + byte[] newConfigData = client.getConfig().forEnsemble(); + QuorumVerifier newConfig = toQuorumVerifier(newConfigData); + List newInstances = Lists.newArrayList(cluster.getInstances()); + newInstances.remove(removeSpec); + assertConfig(newConfig, newInstances); + + assertTrue(timing.awaitLatch(ensembleLatch)); + String connectString = EnsembleTracker.configToConnectionString(newConfig) + chroot; + assertEquals(connectString, ensembleProvider.getConnectionString()); + + client.getZookeeperClient().reset(); + client.sync().forPath("/"); + byte[] data = client.getData().forPath("/"); + assertArrayEquals("deadbeef".getBytes(), data, () -> "expected \"deedbeef\", got data: " + Arrays.toString(data)); + } + } + @Test public void testConfigToConnectionStringIPv4Normal() throws Exception { @@ -558,6 +606,7 @@ private CuratorFramework newClient(String connectionString) { private CuratorFramework newClient(String connectionString, boolean withEnsembleProvider) { final AtomicReference connectString = new AtomicReference<>(connectionString); + ensembleLatch = new CountDownLatch(1); ensembleProvider = new EnsembleProvider() { @Override @@ -585,7 +634,10 @@ public void close() throws IOException @Override public void setConnectionString(String connectionString) { - connectString.set(connectionString); + if (!connectionString.equals(getConnectionString())) { + connectString.set(connectionString); + ensembleLatch.countDown(); + } } }; return CuratorFrameworkFactory.builder() From 28b27c68a8518cbc924aaaf0707816e16b8ed72d Mon Sep 17 00:00:00 2001 From: Kezhu Wang Date: Mon, 22 May 2023 16:01:57 +0800 Subject: [PATCH 2/2] fixup! CURATOR-593: Append chroot for EnsembleProvider::setConnectionString in EnsembleTracker --- .../framework/imps/TestReconfiguration.java | 28 +++++++++++++------ 1 file changed, 19 insertions(+), 9 deletions(-) diff --git a/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java b/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java index e7854b8557..082a948a1e 100644 --- a/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java +++ b/curator-framework/src/test/java/org/apache/curator/framework/imps/TestReconfiguration.java @@ -51,7 +51,6 @@ import java.net.InetAddress; import java.net.InetSocketAddress; import java.util.ArrayList; -import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.Iterator; @@ -61,7 +60,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; -import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; @@ -71,7 +70,6 @@ public class TestReconfiguration extends CuratorTestBase { private final Timing2 timing = new Timing2(); private TestingCluster cluster; - private CountDownLatch ensembleLatch; private EnsembleProvider ensembleProvider; private static final String superUserPasswordDigest = "curator-test:zghsj3JfJqK7DbWf0RQ1BgbJH9w="; // ran from DigestAuthenticationProvider.generateDigest(superUserPassword); @@ -444,8 +442,9 @@ public void testRemoveWithChroot() throws Exception { // Use a long chroot path to circumvent ZOOKEEPER-4565 and ZOOKEEPER-4601 String chroot = "/pretty-long-chroot"; + CountDownLatch ensembleLatch = new CountDownLatch(1); - try (CuratorFramework client = newClient(cluster.getConnectString() + chroot)) { + try (CuratorFramework client = newClient(cluster.getConnectString() + chroot, ensembleLatch)) { client.start(); client.create().forPath("/", "deadbeef".getBytes()); @@ -480,7 +479,7 @@ public void testRemoveWithChroot() throws Exception client.getZookeeperClient().reset(); client.sync().forPath("/"); byte[] data = client.getData().forPath("/"); - assertArrayEquals("deadbeef".getBytes(), data, () -> "expected \"deedbeef\", got data: " + Arrays.toString(data)); + assertThat(data).asString().isEqualTo("deadbeef"); } } @@ -603,10 +602,19 @@ private CuratorFramework newClient(String connectionString) { return newClient(connectionString, true); } - private CuratorFramework newClient(String connectionString, boolean withEnsembleProvider) + private CuratorFramework newClient(String connectionString, boolean withEnsembleTracker) + { + return newClient(connectionString, withEnsembleTracker, null); + } + + private CuratorFramework newClient(String connectionString, CountDownLatch ensembleLatch) + { + return newClient(connectionString, ensembleLatch != null, ensembleLatch); + } + + private CuratorFramework newClient(String connectionString, boolean withEnsembleTracker, CountDownLatch ensembleLatch) { final AtomicReference connectString = new AtomicReference<>(connectionString); - ensembleLatch = new CountDownLatch(1); ensembleProvider = new EnsembleProvider() { @Override @@ -636,13 +644,15 @@ public void setConnectionString(String connectionString) { if (!connectionString.equals(getConnectionString())) { connectString.set(connectionString); - ensembleLatch.countDown(); + if (ensembleLatch != null) { + ensembleLatch.countDown(); + } } } }; return CuratorFrameworkFactory.builder() .ensembleProvider(ensembleProvider) - .ensembleTracker(withEnsembleProvider) + .ensembleTracker(withEnsembleTracker) .sessionTimeoutMs(timing.session()) .connectionTimeoutMs(timing.connection()) .authorization("digest", superUserPassword.getBytes())