From 78b848fbfed3827e84886cf4856b78e892086b1d Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 13 Jul 2023 15:08:22 +0300 Subject: [PATCH 1/4] [fix][test] Fix resource leak in PulsarTestContext - use fields in super class to hold resources instead of having duplicates - all fields in PulsarService have protected setters due to `@Setter(AccessLevel.PROTECTED)` - this fact was missed in the original implementation of PulsarTestContext --- .../AbstractTestPulsarService.java | 44 ++++------- .../NonStartableTestPulsarService.java | 79 +++---------------- 2 files changed, 28 insertions(+), 95 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java index a6861268b94d5..da6ce83906162 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java @@ -37,11 +37,6 @@ */ abstract class AbstractTestPulsarService extends PulsarService { protected final SpyConfig spyConfig; - protected final MetadataStoreExtended localMetadataStore; - protected final MetadataStoreExtended configurationMetadataStore; - protected final Compactor compactor; - protected final BrokerInterceptor brokerInterceptor; - protected final BookKeeperClientFactory bookKeeperClientFactory; public AbstractTestPulsarService(SpyConfig spyConfig, ServiceConfiguration config, MetadataStoreExtended localMetadataStore, @@ -50,53 +45,46 @@ public AbstractTestPulsarService(SpyConfig spyConfig, ServiceConfiguration confi BookKeeperClientFactory bookKeeperClientFactory) { super(config); this.spyConfig = spyConfig; - this.localMetadataStore = - NonClosingProxyHandler.createNonClosingProxy(localMetadataStore, MetadataStoreExtended.class); - this.configurationMetadataStore = - NonClosingProxyHandler.createNonClosingProxy(configurationMetadataStore, MetadataStoreExtended.class); - this.compactor = compactor; - this.brokerInterceptor = brokerInterceptor; - this.bookKeeperClientFactory = bookKeeperClientFactory; + setLocalMetadataStore( + NonClosingProxyHandler.createNonClosingProxy(localMetadataStore, MetadataStoreExtended.class)); + setConfigurationMetadataStore( + NonClosingProxyHandler.createNonClosingProxy(configurationMetadataStore, MetadataStoreExtended.class)); + setCompactor(compactor); + setBrokerInterceptor(brokerInterceptor); + setBkClientFactory(bookKeeperClientFactory); } @Override public MetadataStore createConfigurationMetadataStore(PulsarMetadataEventSynchronizer synchronizer) throws MetadataStoreException { if (synchronizer != null) { - synchronizer.registerSyncListener(configurationMetadataStore::handleMetadataEvent); + synchronizer.registerSyncListener( + ((MetadataStoreExtended) getConfigurationMetadataStore())::handleMetadataEvent); } - return configurationMetadataStore; + return getConfigurationMetadataStore(); } @Override public MetadataStoreExtended createLocalMetadataStore(PulsarMetadataEventSynchronizer synchronizer) throws MetadataStoreException, PulsarServerException { if (synchronizer != null) { - synchronizer.registerSyncListener(localMetadataStore::handleMetadataEvent); + synchronizer.registerSyncListener( + getLocalMetadataStore()::handleMetadataEvent); } - return localMetadataStore; + return getLocalMetadataStore(); } @Override public Compactor newCompactor() throws PulsarServerException { - if (compactor != null) { - return compactor; + if (getCompactor() != null) { + return getCompactor(); } else { return spyConfig.getCompactor().spy(super.newCompactor()); } } - @Override - public BrokerInterceptor getBrokerInterceptor() { - if (brokerInterceptor != null) { - return brokerInterceptor; - } else { - return super.getBrokerInterceptor(); - } - } - @Override public BookKeeperClientFactory newBookKeeperClientFactory() { - return bookKeeperClientFactory; + return getBkClientFactory(); } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/NonStartableTestPulsarService.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/NonStartableTestPulsarService.java index 13c4d7d72af2c..af365ed31934f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/NonStartableTestPulsarService.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/NonStartableTestPulsarService.java @@ -38,11 +38,9 @@ import org.apache.pulsar.broker.resources.TopicResources; import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.schema.DefaultSchemaRegistryService; -import org.apache.pulsar.broker.service.schema.SchemaRegistryService; import org.apache.pulsar.broker.storage.ManagedLedgerStorage; import org.apache.pulsar.broker.transaction.buffer.TransactionBufferProvider; import org.apache.pulsar.broker.transaction.pendingack.TransactionPendingAckStoreProvider; -import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; @@ -56,13 +54,6 @@ * for a "non-startable" PulsarService. Please see {@link PulsarTestContext} for more details. */ class NonStartableTestPulsarService extends AbstractTestPulsarService { - private final PulsarResources pulsarResources; - private final ManagedLedgerStorage managedLedgerClientFactory; - private final BrokerService brokerService; - - private final SchemaRegistryService schemaRegistryService; - - private final PulsarClientImpl pulsarClient; private final NamespaceService namespaceService; @@ -76,16 +67,16 @@ public NonStartableTestPulsarService(SpyConfig spyConfig, ServiceConfiguration c Function brokerServiceCustomizer) { super(spyConfig, config, localMetadataStore, configurationMetadataStore, compactor, brokerInterceptor, bookKeeperClientFactory); - this.pulsarResources = pulsarResources; - this.managedLedgerClientFactory = managedLedgerClientFactory; + setPulsarResources(pulsarResources); + setManagedLedgerClientFactory(managedLedgerClientFactory); try { - this.brokerService = brokerServiceCustomizer.apply( - spyConfig.getBrokerService().spy(TestBrokerService.class, this, getIoEventLoopGroup())); + setBrokerService(brokerServiceCustomizer.apply( + spyConfig.getBrokerService().spy(TestBrokerService.class, this, getIoEventLoopGroup()))); } catch (Exception e) { throw new RuntimeException(e); } - this.schemaRegistryService = spyWithClassAndConstructorArgs(DefaultSchemaRegistryService.class); - this.pulsarClient = mock(PulsarClientImpl.class); + setSchemaRegistryService(spyWithClassAndConstructorArgs(DefaultSchemaRegistryService.class)); + setClient(mock(PulsarClientImpl.class)); this.namespaceService = mock(NamespaceService.class); try { startNamespaceService(); @@ -118,63 +109,17 @@ public Supplier getNamespaceServiceProvider() throws PulsarSer return () -> namespaceService; } - @Override - public synchronized PulsarClient getClient() throws PulsarServerException { - return pulsarClient; - } - @Override public PulsarClientImpl createClientImpl(ClientConfigurationData clientConf) throws PulsarClientException { - return pulsarClient; - } - - @Override - public SchemaRegistryService getSchemaRegistryService() { - return schemaRegistryService; - } - - @Override - public PulsarResources getPulsarResources() { - return pulsarResources; - } - - public BrokerService getBrokerService() { - return brokerService; - } - - @Override - public MetadataStore getConfigurationMetadataStore() { - return configurationMetadataStore; - } - - @Override - public MetadataStoreExtended getLocalMetadataStore() { - return localMetadataStore; - } - - @Override - public ManagedLedgerStorage getManagedLedgerClientFactory() { - return managedLedgerClientFactory; - } - - @Override - protected PulsarResources newPulsarResources() { - return pulsarResources; - } - - @Override - protected ManagedLedgerStorage newManagedLedgerClientFactory() throws Exception { - return managedLedgerClientFactory; + try { + return (PulsarClientImpl) getClient(); + } catch (PulsarServerException e) { + throw new PulsarClientException(e); + } } - @Override protected BrokerService newBrokerService(PulsarService pulsar) throws Exception { - return brokerService; - } - - @Override - public BookKeeperClientFactory getBookKeeperClientFactory() { - return bookKeeperClientFactory; + return getBrokerService(); } static class TestBrokerService extends BrokerService { From eee7b19161638866fbc5a03ba59e9192ce394bca Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 13 Jul 2023 18:19:12 +0300 Subject: [PATCH 2/4] Fix StackOverflowError --- .../broker/testcontext/AbstractTestPulsarService.java | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java index da6ce83906162..d435b98281cc7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java @@ -37,6 +37,7 @@ */ abstract class AbstractTestPulsarService extends PulsarService { protected final SpyConfig spyConfig; + private boolean compactorExists; public AbstractTestPulsarService(SpyConfig spyConfig, ServiceConfiguration config, MetadataStoreExtended localMetadataStore, @@ -74,9 +75,17 @@ public MetadataStoreExtended createLocalMetadataStore(PulsarMetadataEventSynchro return getLocalMetadataStore(); } + @Override + protected void setCompactor(Compactor compactor) { + if (compactor != null) { + compactorExists = true; + } + super.setCompactor(compactor); + } + @Override public Compactor newCompactor() throws PulsarServerException { - if (getCompactor() != null) { + if (compactorExists) { return getCompactor(); } else { return spyConfig.getCompactor().spy(super.newCompactor()); From d1a901c64ade657484de9b1d9a743b2a3c2adb8e Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 13 Jul 2023 18:36:24 +0300 Subject: [PATCH 3/4] Handle broker interceptor for PulsarTestContext --- .../main/java/org/apache/pulsar/broker/PulsarService.java | 6 +++++- .../broker/testcontext/AbstractTestPulsarService.java | 6 ++++++ 2 files changed, 11 insertions(+), 1 deletion(-) 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 5cf3c47bcb047..40c5a2d652817 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 @@ -784,7 +784,7 @@ public void start() throws PulsarServerException { exposeTopicMetrics, offloaderScheduler, interval); this.defaultOffloader = createManagedLedgerOffloader(defaultOffloadPolicies); - this.brokerInterceptor = BrokerInterceptors.load(config); + setBrokerInterceptor(newBrokerInterceptor()); // use getter to support mocking getBrokerInterceptor method in tests BrokerInterceptor interceptor = getBrokerInterceptor(); if (interceptor != null) { @@ -927,6 +927,10 @@ public void start() throws PulsarServerException { } } + protected BrokerInterceptor newBrokerInterceptor() throws IOException { + return BrokerInterceptors.load(config); + } + @VisibleForTesting protected OrderedExecutor newOrderedExecutor() { return OrderedExecutor.newBuilder() diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java index d435b98281cc7..517d57d0042ae 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/AbstractTestPulsarService.java @@ -19,6 +19,7 @@ package org.apache.pulsar.broker.testcontext; +import java.io.IOException; import org.apache.pulsar.broker.BookKeeperClientFactory; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; @@ -96,4 +97,9 @@ public Compactor newCompactor() throws PulsarServerException { public BookKeeperClientFactory newBookKeeperClientFactory() { return getBkClientFactory(); } + + @Override + protected BrokerInterceptor newBrokerInterceptor() throws IOException { + return getBrokerInterceptor() != null ? getBrokerInterceptor() : super.newBrokerInterceptor(); + } } From 24f2528276e0c5fc625d8cb72456e9357615dc73 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 13 Jul 2023 22:22:43 +0300 Subject: [PATCH 4/4] Fix NPE in calling "compactor.getStats().removeTopic(topic)" --- .../apache/pulsar/broker/service/PersistentTopicTest.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index fefed1aaa0a6e..c49df3e85ce92 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -131,6 +131,7 @@ import org.apache.pulsar.compaction.CompactedTopic; import org.apache.pulsar.compaction.CompactedTopicContext; import org.apache.pulsar.compaction.Compactor; +import org.apache.pulsar.compaction.CompactorMXBean; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.impl.FaultInjectionMetadataStore; import org.awaitility.Awaitility; @@ -172,11 +173,13 @@ public void setup() throws Exception { svcConfig.setClusterName("pulsar-cluster"); svcConfig.setTopicLevelPoliciesEnabled(false); svcConfig.setSystemTopicEnabled(false); + Compactor compactor = mock(Compactor.class); + when(compactor.getStats()).thenReturn(mock(CompactorMXBean.class)); pulsarTestContext = PulsarTestContext.builderForNonStartableContext() .config(svcConfig) .spyByDefault() .useTestPulsarResources(metadataStore) - .compactor(mock(Compactor.class)) + .compactor(compactor) .build(); brokerService = pulsarTestContext.getBrokerService();