From b78f828e3a875d9ec5928df1ee539354a1664c8d Mon Sep 17 00:00:00 2001 From: Yan Zhao Date: Sat, 14 Jan 2023 18:56:05 +0800 Subject: [PATCH 1/3] [fix][broker] Support zookeeper read-only config. (#19156) (cherry picked from commit accae9f371f87dcf0b2a2f8d4ed277b742bfdc5f) --- .../apache/pulsar/broker/ServiceConfiguration.java | 6 ++++++ .../pulsar/broker/resources/PulsarResources.java | 12 ++++++++++-- .../java/org/apache/pulsar/broker/PulsarService.java | 4 ++-- .../websocket/proxy/ProxyAuthenticationTest.java | 4 +++- .../websocket/proxy/ProxyAuthorizationTest.java | 4 +++- .../websocket/proxy/ProxyConfigurationTest.java | 4 +++- .../websocket/proxy/ProxyPublishConsumeTest.java | 4 +++- .../websocket/proxy/ProxyPublishConsumeTlsTest.java | 4 +++- .../proxy/ProxyPublishConsumeWithoutZKTest.java | 4 +++- .../proxy/v1/V1_ProxyAuthenticationTest.java | 4 +++- .../apache/pulsar/functions/worker/WorkerConfig.java | 6 ++++++ .../org/apache/pulsar/functions/worker/Worker.java | 3 ++- .../pulsar/proxy/server/ProxyConfiguration.java | 6 ++++++ .../org/apache/pulsar/proxy/server/ProxyService.java | 8 +++++--- .../apache/pulsar/websocket/WebSocketService.java | 10 ++++++---- 15 files changed, 64 insertions(+), 19 deletions(-) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index b0a8581aac953..6337d3cf7e4a6 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -417,6 +417,12 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private int zooKeeperCacheExpirySeconds = -1; + @FieldContext( + category = CATEGORY_SERVER, + doc = "Is zookeeper allow read-only operations." + ) + private boolean zooKeeperAllowReadOnlyOperations; + @FieldContext( category = CATEGORY_SERVER, dynamic = true, diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java index a087d8090d3a4..3c1a8fb41dfe0 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java @@ -88,9 +88,17 @@ public PulsarResources(MetadataStore localMetadataStore, MetadataStore configura this.configurationMetadataStore = Optional.ofNullable(configurationMetadataStore); } - public static MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs) + public static MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs, + boolean allowReadOnlyOperations) throws MetadataStoreException { return MetadataStoreExtended.create(serverUrls, MetadataStoreConfig.builder() - .sessionTimeoutMillis(sessionTimeoutMs).allowReadOnlyOperations(false).build()); + .sessionTimeoutMillis(sessionTimeoutMs).allowReadOnlyOperations(allowReadOnlyOperations).build()); + } + + public static MetadataStoreExtended createConfigMetadataStore(String serverUrls, int sessionTimeoutMs, + boolean allowReadOnlyOperations) + throws MetadataStoreException { + return MetadataStoreExtended.create(serverUrls, MetadataStoreConfig.builder() + .sessionTimeoutMillis(sessionTimeoutMs).allowReadOnlyOperations(allowReadOnlyOperations).build()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 2bda5c5539883..e2034bdc60b62 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -347,7 +347,7 @@ public MetadataStore createConfigurationMetadataStore() throws MetadataStoreExce return MetadataStoreFactory.create(config.getConfigurationMetadataStoreUrl(), MetadataStoreConfig.builder() .sessionTimeoutMillis((int) config.getMetadataStoreSessionTimeoutMillis()) - .allowReadOnlyOperations(false) + .allowReadOnlyOperations(config.isZooKeeperAllowReadOnlyOperations()) .configFilePath(config.getMetadataStoreConfigPath()) .batchingEnabled(config.isMetadataStoreBatchingEnabled()) .batchingMaxDelayMillis(config.getMetadataStoreBatchingMaxDelayMillis()) @@ -980,7 +980,7 @@ public MetadataStoreExtended createLocalMetadataStore() throws MetadataStoreExce return MetadataStoreExtended.create(config.getMetadataStoreUrl(), MetadataStoreConfig.builder() .sessionTimeoutMillis((int) config.getMetadataStoreSessionTimeoutMillis()) - .allowReadOnlyOperations(false) + .allowReadOnlyOperations(config.isZooKeeperAllowReadOnlyOperations()) .configFilePath(config.getMetadataStoreConfigPath()) .batchingEnabled(config.isMetadataStoreBatchingEnabled()) .batchingMaxDelayMillis(config.getMetadataStoreBatchingMaxDelayMillis()) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java index a34ec879ba609..49ba21ce3882e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java @@ -20,6 +20,7 @@ import static java.util.concurrent.Executors.newFixedThreadPool; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -83,7 +84,8 @@ public void setup() throws Exception { } service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java index 7e0ee1bd4669e..a62f4751459c4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.websocket.proxy; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -64,7 +65,8 @@ protected void setup() throws Exception { config.setWebServicePort(Optional.of(0)); config.setConfigurationMetadataStoreUrl(GLOBAL_DUMMY_VALUE); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); service.start(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java index 26cf8a0e1549f..8ff00cf0ce9b3 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.websocket.proxy; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -67,7 +68,8 @@ public void configTest(int numIoThreads, int connectionsPerBroker) throws Except config.setServiceUrl("http://localhost:8080"); config.getProperties().setProperty("brokerClient_lookupTimeoutMs", "100"); WebSocketService service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); service.start(); PulsarClientImpl client = (PulsarClientImpl) service.getPulsarClient(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java index 918640642ecbd..02d9ceb1c1ad9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java @@ -20,6 +20,7 @@ import static java.util.concurrent.Executors.newFixedThreadPool; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -101,7 +102,8 @@ public void setup() throws Exception { config.setClusterName("test"); config.setConfigurationMetadataStoreUrl(GLOBAL_DUMMY_VALUE); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java index a8b67416107c7..b98414a927014 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java @@ -20,6 +20,7 @@ import static java.util.concurrent.Executors.newFixedThreadPool; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -74,7 +75,8 @@ public void setup() throws Exception { config.setBrokerClientAuthenticationPlugin(AuthenticationTls.class.getName()); config.setConfigurationMetadataStoreUrl(GLOBAL_DUMMY_VALUE); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java index 1fb12645e5e34..dadc246c98ebc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java @@ -20,6 +20,7 @@ import static java.util.concurrent.Executors.newFixedThreadPool; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -62,7 +63,8 @@ public void setup() throws Exception { config.setServiceUrl(pulsar.getSafeWebServiceAddress()); config.setServiceUrlTls(pulsar.getWebServiceAddressTls()); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeper)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeper)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java index b80c3fb07be5e..af9b00e59211f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java @@ -20,6 +20,7 @@ import static java.util.concurrent.Executors.newFixedThreadPool; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -85,7 +86,8 @@ public void setup() throws Exception { } service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java index a55f49bfc2031..27a6e8eea5112 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java @@ -205,6 +205,12 @@ public class WorkerConfig implements Serializable, PulsarConfiguration { ) private int zooKeeperCacheExpirySeconds = -1; + @FieldContext( + category = CATEGORY_WORKER, + doc = "Is zooKeeper allow read-only operations." + ) + private boolean zooKeeperAllowReadOnlyOperations; + @FieldContext( category = CATEGORY_CONNECTORS, doc = "The path to the location to locate builtin connectors" diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java index f8871cb86e792..189daae348f49 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java @@ -75,7 +75,8 @@ private AuthorizationService getAuthorizationService() throws PulsarServerExcept try { configMetadataStore = PulsarResources.createMetadataStore( workerConfig.getConfigurationMetadataStoreUrl(), - (int) workerConfig.getMetadataStoreSessionTimeoutMillis()); + (int) workerConfig.getMetadataStoreSessionTimeoutMillis(), + workerConfig.isZooKeeperAllowReadOnlyOperations()); } catch (IOException e) { throw new PulsarServerException(e); } diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java index d139bb83002c0..a3183eb749c8e 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConfiguration.java @@ -148,6 +148,12 @@ public class ProxyConfiguration implements PulsarConfiguration { ) private int zooKeeperCacheExpirySeconds = -1; + @FieldContext( + category = CATEGORY_SERVER, + doc = "Is zooKeeper allow read-only operations." + ) + private boolean zooKeeperAllowReadOnlyOperations; + @FieldContext( category = CATEGORY_BROKER_DISCOVERY, doc = "The service url points to the broker cluster" diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java index 1960b5143a038..c29e2ba169650 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyService.java @@ -427,12 +427,14 @@ public Optional getListenPortTls() { public MetadataStoreExtended createLocalMetadataStore() throws MetadataStoreException { return PulsarResources.createMetadataStore(proxyConfig.getMetadataStoreUrl(), - proxyConfig.getMetadataStoreSessionTimeoutMillis()); + proxyConfig.getMetadataStoreSessionTimeoutMillis(), + proxyConfig.isZooKeeperAllowReadOnlyOperations()); } public MetadataStoreExtended createConfigurationMetadataStore() throws MetadataStoreException { - return PulsarResources.createMetadataStore(proxyConfig.getConfigurationMetadataStoreUrl(), - proxyConfig.getMetadataStoreSessionTimeoutMillis()); + return PulsarResources.createConfigMetadataStore(proxyConfig.getConfigurationMetadataStoreUrl(), + proxyConfig.getMetadataStoreSessionTimeoutMillis(), + proxyConfig.isZooKeeperAllowReadOnlyOperations()); } public Authentication getProxyClientAuthenticationPlugin() { diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java index a57c6c491e78a..09be1b7026452 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java @@ -102,8 +102,9 @@ public void start() throws PulsarServerException, PulsarClientException, Malform if (isNotBlank(config.getConfigurationMetadataStoreUrl())) { try { - configMetadataStore = createMetadataStore(config.getConfigurationMetadataStoreUrl(), - (int) config.getMetadataStoreSessionTimeoutMillis()); + configMetadataStore = createConfigMetadataStore(config.getConfigurationMetadataStoreUrl(), + (int) config.getMetadataStoreSessionTimeoutMillis(), + config.isZooKeeperAllowReadOnlyOperations()); } catch (MetadataStoreException e) { throw new PulsarServerException(e); } @@ -123,9 +124,10 @@ public void start() throws PulsarServerException, PulsarClientException, Malform log.info("Pulsar WebSocket Service started"); } - public MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs) + public MetadataStoreExtended createConfigMetadataStore(String serverUrls, int sessionTimeoutMs, boolean + isAllowReadOnlyOperations) throws MetadataStoreException { - return PulsarResources.createMetadataStore(serverUrls, sessionTimeoutMs); + return PulsarResources.createConfigMetadataStore(serverUrls, sessionTimeoutMs, isAllowReadOnlyOperations); } @Override From d2a2ce959e9a5a74aa3a4b402283b3b21d20a98f Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Sun, 26 Feb 2023 20:16:14 +0800 Subject: [PATCH 2/3] fix --- .../apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java | 3 ++- .../java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java index 71f8c0b6d8602..32b02adc1b2cf 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java @@ -64,7 +64,8 @@ public void setup() throws Exception { config.setConfigurationStoreServers(GLOBAL_DUMMY_VALUE); config.setWebSocketSessionIdleTimeoutMillis(3 * 1000); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)) + .when(service).createConfigMetadataStore(anyString(), anyInt(), false); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java index 6e0dfa46cfdbc..ddc892fcb12fa 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java @@ -66,7 +66,8 @@ public void setup() throws Exception { config.setWebSocketSessionIdleTimeoutMillis(3 * 1000); config.setWebSocketPingDurationSeconds(2); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createConfigMetadataStore(anyString(), anyInt(), false); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); From 029393e30f019910471b65d520967034dc2ef6b0 Mon Sep 17 00:00:00 2001 From: xiangying <1984997880@qq.com> Date: Wed, 1 Mar 2023 17:49:06 +0800 Subject: [PATCH 3/3] fix --- .../apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java | 3 ++- .../java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java index 32b02adc1b2cf..18451d9e7510f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyIdleTimeoutTest.java @@ -21,6 +21,7 @@ import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -65,7 +66,7 @@ public void setup() throws Exception { config.setWebSocketSessionIdleTimeoutMillis(3 * 1000); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); doReturn(new ZKMetadataStore(mockZooKeeperGlobal)) - .when(service).createConfigMetadataStore(anyString(), anyInt(), false); + .when(service).createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java index ddc892fcb12fa..9ab2629fe019d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPingTest.java @@ -21,6 +21,7 @@ import static java.util.concurrent.TimeUnit.SECONDS; import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; @@ -67,7 +68,7 @@ public void setup() throws Exception { config.setWebSocketPingDurationSeconds(2); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) - .createConfigMetadataStore(anyString(), anyInt(), false); + .createConfigMetadataStore(anyString(), anyInt(), anyBoolean()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started");