From ff330eed7e48ecdd287e248cacce2aa4d8e07e28 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Fri, 16 Sep 2022 10:30:13 +0200 Subject: [PATCH 1/3] [fix][functions] Ensure InternalConfigurationData data model is compatible across different versions --- .../apache/pulsar/broker/admin/AdminTest.java | 48 +++++++++++++++++++ .../conf/InternalConfigurationData.java | 36 +++++++++++++- 2 files changed, 82 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index 9bbcf18b9e954..c47e4d46e78d5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -51,6 +51,9 @@ import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.StreamingOutput; import javax.ws.rs.core.UriInfo; + +import lombok.AllArgsConstructor; +import lombok.Data; import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.mledger.proto.PendingBookieOpsStats; import org.apache.pulsar.broker.ServiceConfiguration; @@ -87,9 +90,11 @@ import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.apache.pulsar.common.stats.AllocatorStats; import org.apache.pulsar.common.stats.Metrics; +import org.apache.pulsar.common.util.ObjectMapperFactory; import org.apache.pulsar.functions.worker.WorkerConfig; import org.apache.pulsar.metadata.cache.impl.MetadataCacheImpl; import org.apache.pulsar.metadata.impl.AbstractMetadataStore; +import org.apache.pulsar.metadata.impl.MetadataStoreFactoryImpl; import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; import org.apache.zookeeper.KeeperException.Code; import org.apache.zookeeper.MockZooKeeper; @@ -205,6 +210,49 @@ public void internalConfiguration() throws Exception { assertEquals(response, expectedData); } + @Data + @AllArgsConstructor + /** + * Internal configuration data model before (before https://github.com/apache/pulsar/pull/14384). + */ + private static class OldInternalConfigurationData { + private String zookeeperServers; + private String configurationStoreServers; + @Deprecated + private String ledgersRootPath; + private String bookkeeperMetadataServiceUri; + private String stateStorageServiceUrl; + } + + /** + * This test verifies that the model data changes in InternalConfigurationData are retro-compatible. + * InternalConfigurationData is downloaded from the Function worker from a broker. + * The broker may be still serve an "old" version of InternalConfigurationData + * (before https://github.com/apache/pulsar/pull/14384) while the Worker already uses the new one. + * @throws Exception + */ + @Test + public void internalConfigurationRetroCompatibility() throws Exception { + OldInternalConfigurationData oldDataModel = new OldInternalConfigurationData( + MetadataStoreFactoryImpl.removeIdentifierFromMetadataURL(conf.getMetadataStoreUrl()), + conf.getConfigurationMetadataStoreUrl(), + new ClientConfiguration().getZkLedgersRootPath(), + conf.isBookkeeperMetadataStoreSeparated() ? conf.getBookkeeperMetadataStoreUrl() : null, + pulsar.getWorkerConfig().map(WorkerConfig::getStateStorageServiceUrl).orElse(null)); + + final Map oldDataJson = ObjectMapperFactory + .getThreadLocal().convertValue(oldDataModel, Map.class); + + final InternalConfigurationData newData = ObjectMapperFactory.getThreadLocal() + .convertValue(oldDataJson, InternalConfigurationData.class); + + assertEquals(newData.getMetadataStoreUrl(), conf.getMetadataStoreUrl()); + assertEquals(newData.getConfigurationMetadataStoreUrl(), oldDataModel.getConfigurationStoreServers()); + assertEquals(newData.getLedgersRootPath(), oldDataModel.getLedgersRootPath()); + assertEquals(newData.getBookkeeperMetadataServiceUri(), oldDataModel.getBookkeeperMetadataServiceUri()); + assertEquals(newData.getStateStorageServiceUrl(), oldDataModel.getStateStorageServiceUrl()); + } + @Test @SuppressWarnings("unchecked") public void clusters() throws Exception { diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java index 2927aad8afe41..e9990799363dd 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java @@ -27,6 +27,10 @@ @ToString public class InternalConfigurationData { + @Deprecated + private String zookeeperServers; + @Deprecated + private String configurationStoreServers; private String metadataStoreUrl; private String configurationMetadataStoreUrl; @Deprecated @@ -49,12 +53,40 @@ public InternalConfigurationData(String zookeeperServers, this.stateStorageServiceUrl = stateStorageServiceUrl; } + @Deprecated + public String getZookeeperServers() { + return zookeeperServers; + } + + @Deprecated + public void setZookeeperServers(String zookeeperServers) { + this.zookeeperServers = zookeeperServers; + } + + @Deprecated + public String getConfigurationStoreServers() { + return configurationStoreServers; + } + + @Deprecated + public void setConfigurationStoreServers(String configurationStoreServers) { + this.configurationStoreServers = configurationStoreServers; + } + public String getMetadataStoreUrl() { - return metadataStoreUrl; + if (metadataStoreUrl != null) { + return metadataStoreUrl; + } else if (zookeeperServers != null) { + return "zk:" + zookeeperServers; + } + return null; } public String getConfigurationMetadataStoreUrl() { - return configurationMetadataStoreUrl; + if (configurationMetadataStoreUrl != null) { + return configurationMetadataStoreUrl; + } + return configurationStoreServers; } /** @deprecated */ From 2b3190b9d543c3e2542f57abb1b54900d63ed021 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Fri, 16 Sep 2022 10:38:23 +0200 Subject: [PATCH 2/3] style --- .../src/test/java/org/apache/pulsar/broker/admin/AdminTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index c47e4d46e78d5..d035d7c4290e9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -51,7 +51,6 @@ import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.StreamingOutput; import javax.ws.rs.core.UriInfo; - import lombok.AllArgsConstructor; import lombok.Data; import org.apache.bookkeeper.conf.ClientConfiguration; From b98eb1ae02a953461330c58cfb0f04a73b6ae39b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Fri, 16 Sep 2022 10:39:31 +0200 Subject: [PATCH 3/3] fix the other way --- .../apache/pulsar/common/conf/InternalConfigurationData.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java index e9990799363dd..1099debfc178c 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java @@ -47,7 +47,9 @@ public InternalConfigurationData(String zookeeperServers, String bookkeeperMetadataServiceUri, String stateStorageServiceUrl) { this.metadataStoreUrl = zookeeperServers; + this.zookeeperServers = zookeeperServers; this.configurationMetadataStoreUrl = configurationMetadataStoreUrl; + this.configurationStoreServers = configurationMetadataStoreUrl; this.ledgersRootPath = ledgersRootPath; this.bookkeeperMetadataServiceUri = bookkeeperMetadataServiceUri; this.stateStorageServiceUrl = stateStorageServiceUrl;