diff --git a/pom.xml b/pom.xml
index 971de8db8d2ec..97b64ca0fc84b 100644
--- a/pom.xml
+++ b/pom.xml
@@ -240,7 +240,6 @@ flexible messaging model and an intuitive client API.
1.1.1
7.3.0
3.12.4
- 2.0.9
3.25.0-GA
1.5.0
3.1
@@ -339,12 +338,6 @@ flexible messaging model and an intuitive client API.
${mockito.version}
-
- org.powermock
- powermock-reflect
- ${powermock.version}
-
-
org.apache.zookeeper
zookeeper
@@ -1330,12 +1323,6 @@ flexible messaging model and an intuitive client API.
test
-
- org.powermock
- powermock-reflect
- test
-
-
org.assertj
assertj-core
diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnershipCache.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnershipCache.java
index 67e986b804cba..9c2030997fbdb 100644
--- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnershipCache.java
+++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnershipCache.java
@@ -106,7 +106,7 @@ public CompletableFuture asyncLoad(NamespaceBundle namespaceBundle,
.thenRun(() -> {
log.info("Resource lock for {} has expired", rl.getPath());
namespaceService.unloadNamespaceBundle(namespaceBundle);
- ownedBundlesCache.synchronous().invalidate(namespaceBundle);
+ invalidateLocalOwnerCache(namespaceBundle);
namespaceService.onNamespaceBundleUnload(namespaceBundle);
});
return new OwnedBundle(namespaceBundle);
@@ -330,6 +330,10 @@ public void invalidateLocalOwnerCache() {
this.ownedBundlesCache.synchronous().invalidateAll();
}
+ public void invalidateLocalOwnerCache(NamespaceBundle namespaceBundle) {
+ this.ownedBundlesCache.synchronous().invalidate(namespaceBundle);
+ }
+
public synchronized boolean refreshSelfOwnerInfo() {
this.selfOwnerInfo = new NamespaceEphemeralData(pulsar.getBrokerServiceUrl(),
pulsar.getBrokerServiceUrlTls(), pulsar.getSafeWebServiceAddress(),
diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
index fcb4135c97ece..e0327f5c29f8c 100644
--- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
+++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
@@ -1280,6 +1280,11 @@ public ManagedCursor getCursor() {
return cursor;
}
+ @VisibleForTesting
+ public PendingAckHandle getPendingAckHandle() {
+ return pendingAckHandle;
+ }
+
public void syncBatchPositionBitSetForPendingAck(PositionImpl position) {
this.pendingAckHandle.syncBatchPositionAckSetForTransaction(position);
}
diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
index 52bb165a61e91..4d6ea3a512079 100644
--- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
+++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
@@ -382,6 +382,11 @@ private void initializeDispatchRateLimiterIfNeeded() {
}
}
+ @VisibleForTesting
+ public AtomicLong getPendingWriteOps() {
+ return pendingWriteOps;
+ }
+
private PersistentSubscription createPersistentSubscription(String subscriptionName, ManagedCursor cursor,
boolean replicated, Map subscriptionProperties) {
checkNotNull(compactedTopic);
diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/PendingAckHandleImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/PendingAckHandleImpl.java
index 283bc038d764d..a565a20b6649d 100644
--- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/PendingAckHandleImpl.java
+++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/PendingAckHandleImpl.java
@@ -21,6 +21,7 @@
import static org.apache.bookkeeper.mledger.util.PositionAckSetUtil.andAckSet;
import static org.apache.bookkeeper.mledger.util.PositionAckSetUtil.compareToWithAckSet;
import static org.apache.bookkeeper.mledger.util.PositionAckSetUtil.isAckSetOverlap;
+import com.google.common.annotations.VisibleForTesting;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@@ -1027,6 +1028,11 @@ public PositionInPendingAckStats checkPositionInPendingAckState(PositionImpl pos
}
}
+ @VisibleForTesting
+ public Map> getIndividualAckPositions() {
+ return individualAckPositions;
+ }
+
@Override
public boolean checkIfPendingAckStoreInit() {
return this.pendingAckStoreFuture != null && this.pendingAckStoreFuture.isDone();
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/BookKeeperClientFactoryImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/BookKeeperClientFactoryImplTest.java
index e26b0aa756162..b9c32e91e4c18 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/BookKeeperClientFactoryImplTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/BookKeeperClientFactoryImplTest.java
@@ -37,10 +37,10 @@
import org.apache.bookkeeper.net.CachedDNSToSwitchMapping;
import org.apache.bookkeeper.stats.StatsLogger;
import org.apache.commons.configuration.ConfigurationException;
+import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.pulsar.bookie.rackawareness.BookieRackAffinityMapping;
import org.apache.pulsar.metadata.api.MetadataStore;
import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended;
-import org.powermock.reflect.Whitebox;
import org.testng.annotations.Test;
/**
@@ -281,7 +281,7 @@ public void testOpportunisticStripingConfiguration() {
}
@Test
- public void testBookKeeperIoThreadsConfiguration() {
+ public void testBookKeeperIoThreadsConfiguration() throws Exception {
BookKeeperClientFactoryImpl factory = new BookKeeperClientFactoryImpl();
ServiceConfiguration conf = new ServiceConfiguration();
assertEquals(factory.createBkClientConfiguration(mock(MetadataStoreExtended.class), conf)
@@ -292,11 +292,11 @@ public void testBookKeeperIoThreadsConfiguration() {
EventLoopGroup eventLoopGroup = mock(EventLoopGroup.class);
BookKeeper.Builder builder = factory.getBookKeeperBuilder(conf, eventLoopGroup,
mock(StatsLogger.class), mock(ClientConfiguration.class));
- assertEquals(Whitebox.getInternalState(builder, "eventLoopGroup"), eventLoopGroup);
+ assertEquals(FieldUtils.readField(builder, "eventLoopGroup", true), eventLoopGroup);
conf.setBookkeeperClientSeparatedIoThreadsEnabled(true);
builder = factory.getBookKeeperBuilder(conf, eventLoopGroup,
mock(StatsLogger.class), mock(ClientConfiguration.class));
- assertNull(Whitebox.getInternalState(builder, "eventLoopGroup"));
+ assertNull(FieldUtils.readField(builder, "eventLoopGroup", true));
}
}
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceCloseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceCloseTest.java
index c424132855b6d..1fbb40a6a5614 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceCloseTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceCloseTest.java
@@ -24,9 +24,9 @@
import java.util.concurrent.ScheduledFuture;
import lombok.extern.slf4j.Slf4j;
+import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.broker.loadbalance.LoadSheddingTask;
-import org.awaitility.reflect.WhiteboxImpl;
import org.testng.annotations.AfterClass;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.Test;
@@ -62,18 +62,21 @@ protected PulsarService startBrokerWithoutAuthorization(ServiceConfiguration con
@Test(timeOut = 30_000)
public void closeInTimeTest() throws Exception {
LoadSheddingTask task = pulsar.getLoadSheddingTask();
- boolean isCancel = WhiteboxImpl.getInternalState(task, "isCancel");
- assertFalse(isCancel);
- ScheduledFuture> loadSheddingFuture = WhiteboxImpl.getInternalState(task, "future");
- assertFalse(loadSheddingFuture.isCancelled());
+
+ {
+ assertFalse((boolean) FieldUtils.readField(task, "isCancel", true));
+ ScheduledFuture> loadSheddingFuture = (ScheduledFuture>) FieldUtils.readField(task, "future", true);
+ assertFalse(loadSheddingFuture.isCancelled());
+ }
// The pulsar service is not used, so it should be closed gracefully in short time.
pulsar.close();
- isCancel = WhiteboxImpl.getInternalState(task, "isCancel");
- assertTrue(isCancel);
- loadSheddingFuture = WhiteboxImpl.getInternalState(task, "future");
- assertTrue(loadSheddingFuture.isCancelled());
+ {
+ assertTrue((boolean) FieldUtils.readField(task, "isCancel", true));
+ ScheduledFuture> loadSheddingFuture = (ScheduledFuture>) FieldUtils.readField(task, "future", true);
+ assertTrue(loadSheddingFuture.isCancelled());
+ }
}
}
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiGetLastMessageIdTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiGetLastMessageIdTest.java
index 2a97bc4f8f2ea..d919958618185 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiGetLastMessageIdTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiGetLastMessageIdTest.java
@@ -22,7 +22,6 @@
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
-import java.lang.reflect.Field;
import java.util.Collection;
import java.util.Date;
import java.util.Map;
@@ -32,13 +31,11 @@
import java.util.concurrent.TimeUnit;
import javax.ws.rs.container.AsyncResponse;
import javax.ws.rs.container.TimeoutHandler;
-import javax.ws.rs.core.UriInfo;
import org.apache.pulsar.broker.admin.v2.PersistentTopics;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.broker.authentication.AuthenticationDataHttps;
import org.apache.pulsar.broker.service.Topic;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
-import org.apache.pulsar.broker.web.PulsarWebResource;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.MessageRoutingMode;
import org.apache.pulsar.client.api.Producer;
@@ -51,26 +48,16 @@
import org.awaitility.Awaitility;
import org.testng.Assert;
import org.testng.annotations.AfterMethod;
-import org.testng.annotations.BeforeClass;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
@Test(groups = "broker-admin")
public class AdminApiGetLastMessageIdTest extends MockedPulsarServiceBaseTest {
- private PersistentTopics persistentTopics;
- private final String testTenant = "my-tenant";
- private final String testLocalCluster = "use";
- private final String testNamespace = "my-namespace";
- protected Field uriField;
- protected UriInfo uriInfo;
+ private static final String testTenant = "my-tenant";
+ private static final String testNamespace = "my-namespace";
- @BeforeClass
- public void initPersistentTopics() throws Exception {
- uriField = PulsarWebResource.class.getDeclaredField("uri");
- uriField.setAccessible(true);
- uriInfo = mock(UriInfo.class);
- }
+ private PersistentTopics persistentTopics;
@Override
@BeforeMethod
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
index de555c1715d61..3a9bd21245bf0 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java
@@ -18,7 +18,6 @@
*/
package org.apache.pulsar.broker.admin;
-import static org.apache.pulsar.broker.BrokerTestUtil.spyWithClassAndConstructorArgs;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyString;
@@ -94,7 +93,6 @@
import org.apache.zookeeper.KeeperException;
import org.awaitility.Awaitility;
import org.mockito.ArgumentCaptor;
-import org.powermock.reflect.Whitebox;
import org.testng.Assert;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeClass;
@@ -149,9 +147,8 @@ protected void setup() throws Exception {
PulsarResources resources =
spy(new PulsarResources(pulsar.getLocalMetadataStore(), pulsar.getConfigurationMetadataStore()));
- doReturn(spyWithClassAndConstructorArgs(TopicResources.class, pulsar.getLocalMetadataStore())).when(resources)
- .getTopicResources();
- Whitebox.setInternalState(pulsar, "pulsarResources", resources);
+ doReturn(spy(new TopicResources(pulsar.getLocalMetadataStore()))).when(resources).getTopicResources();
+ doReturn(resources).when(pulsar).getPulsarResources();
admin.clusters().createCluster("use", ClusterData.builder().serviceUrl("http://broker-use.com:8080").build());
admin.clusters().createCluster("test", ClusterData.builder().serviceUrl("http://broker-use.com:8080").build());
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicAutoCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicAutoCreationTest.java
index 7bd15992f640d..e8c0683569f95 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicAutoCreationTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicAutoCreationTest.java
@@ -38,9 +38,9 @@
import org.apache.pulsar.client.api.ProducerConsumerBase;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.impl.LookupService;
+import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.common.naming.NamespaceName;
import org.apache.pulsar.common.partition.PartitionedTopicMetadata;
-import org.powermock.reflect.Whitebox;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
@@ -122,15 +122,14 @@ public void testPartitionedTopicAutoCreationForbiddenDuringNamespaceDeletion()
});
- LookupService original = Whitebox.getInternalState(pulsarClient, "lookup");
+ LookupService original = ((PulsarClientImpl) pulsarClient).getLookup();
try {
// we want to skip the "lookup" phase, because it is blocked by the HTTP API
LookupService mockLookup = mock(LookupService.class);
- Whitebox.setInternalState(pulsarClient, "lookup", mockLookup);
- when(mockLookup.getPartitionedTopicMetadata(any())).thenAnswer(i -> {
- return CompletableFuture.completedFuture(new PartitionedTopicMetadata(0));
- });
+ ((PulsarClientImpl) pulsarClient).setLookup(mockLookup);
+ when(mockLookup.getPartitionedTopicMetadata(any())).thenAnswer(
+ i -> CompletableFuture.completedFuture(new PartitionedTopicMetadata(0)));
when(mockLookup.getBroker(any())).thenAnswer(i -> {
InetSocketAddress brokerAddress =
new InetSocketAddress(pulsar.getAdvertisedAddress(), pulsar.getBrokerListenPort().get());
@@ -139,20 +138,20 @@ public void testPartitionedTopicAutoCreationForbiddenDuringNamespaceDeletion()
// Creating a producer and creating a Consumer may trigger automatic topic
// creation, let's try to create a Producer and a Consumer
- try (Producer producer = pulsarClient.newProducer()
+ try (Producer ignored = pulsarClient.newProducer()
.sendTimeout(1, TimeUnit.SECONDS)
.topic(topic)
- .create();) {
+ .create()) {
} catch (PulsarClientException.LookupException expected) {
String msg = "Namespace bundle for topic (%s) not served by this instance";
log.info("Expected error", expected);
assertTrue(expected.getMessage().contains(String.format(msg, topic)));
}
- try (Consumer consumer = pulsarClient.newConsumer()
+ try (Consumer ignored = pulsarClient.newConsumer()
.topic(topic)
.subscriptionName("test")
- .subscribe();) {
+ .subscribe()) {
} catch (PulsarClientException.LookupException expected) {
String msg = "Namespace bundle for topic (%s) not served by this instance";
log.info("Expected error", expected);
@@ -170,17 +169,16 @@ public void testPartitionedTopicAutoCreationForbiddenDuringNamespaceDeletion()
admin.topics().getList(namespaceName).isEmpty();
// create now the topic using auto creation
- Whitebox.setInternalState(pulsarClient, "lookup", original);
-
- try (Consumer consumer = pulsarClient.newConsumer()
+ ((PulsarClientImpl) pulsarClient).setLookup(original);
+ try (Consumer ignored = pulsarClient.newConsumer()
.topic(topic)
.subscriptionName("test")
- .subscribe();) {
+ .subscribe()) {
}
admin.topics().getList(namespaceName).contains(topic);
} finally {
- Whitebox.setInternalState(pulsarClient, "lookup", original);
+ ((PulsarClientImpl) pulsarClient).setLookup(original);
}
}
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicsTest.java
index e7dfff6de15df..49ead1c5fcbf2 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicsTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicsTest.java
@@ -49,6 +49,7 @@
import org.apache.avro.io.JsonEncoder;
import org.apache.avro.reflect.ReflectDatumWriter;
import org.apache.avro.util.Utf8;
+import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.broker.authentication.AuthenticationDataHttps;
@@ -88,7 +89,6 @@
import org.mockito.ArgumentCaptor;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
-import org.powermock.reflect.Whitebox;
import org.testng.Assert;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
@@ -320,7 +320,7 @@ public void testLookUpWithRedirect() throws Exception {
doReturn(false).when(topics).isRequestHttps();
UriInfo uriInfo = mock(UriInfo.class);
doReturn(requestPath).when(uriInfo).getRequestUri();
- Whitebox.setInternalState(topics, "uri", uriInfo);
+ FieldUtils.writeField(topics, "uri", uriInfo, true);
//do produce on another broker
topics.setPulsar(pulsar2);
AsyncResponse asyncResponse = mock(AsyncResponse.class);
diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/LoadBalancerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/LoadBalancerTest.java
index 8be88350ea653..5d62ec2c58c84 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/LoadBalancerTest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/LoadBalancerTest.java
@@ -18,7 +18,6 @@
*/
package org.apache.pulsar.broker.loadbalance;
-import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
@@ -26,7 +25,6 @@
import static org.testng.Assert.assertNotNull;
import static org.testng.Assert.assertTrue;
import java.lang.reflect.Field;
-import java.lang.reflect.Method;
import java.net.URL;
import java.util.ArrayList;
import java.util.Collections;
@@ -35,13 +33,17 @@
import java.util.Map;
import java.util.Optional;
import java.util.Set;
+import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
+import lombok.SneakyThrows;
import org.apache.bookkeeper.util.ZkUtils;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.commons.lang3.reflect.MethodUtils;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.loadbalance.impl.PulsarResourceDescription;
@@ -68,7 +70,8 @@
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.ZooDefs;
import org.awaitility.Awaitility;
-import org.powermock.reflect.Whitebox;
+import org.mockito.MockedConstruction;
+import org.mockito.Mockito;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.testng.Assert;
@@ -92,12 +95,12 @@ public class LoadBalancerTest {
private static final int MAX_RETRIES = 15;
private static final int BROKER_COUNT = 5;
- private int[] brokerWebServicePorts = new int[BROKER_COUNT];
- private int[] brokerNativeBrokerPorts = new int[BROKER_COUNT];
- private URL[] brokerUrls = new URL[BROKER_COUNT];
- private String[] lookupAddresses = new String[BROKER_COUNT];
- private PulsarService[] pulsarServices = new PulsarService[BROKER_COUNT];
- private PulsarAdmin[] pulsarAdmins = new PulsarAdmin[BROKER_COUNT];
+ private final int[] brokerWebServicePorts = new int[BROKER_COUNT];
+ private final int[] brokerNativeBrokerPorts = new int[BROKER_COUNT];
+ private final URL[] brokerUrls = new URL[BROKER_COUNT];
+ private final String[] lookupAddresses = new String[BROKER_COUNT];
+ private final PulsarService[] pulsarServices = new PulsarService[BROKER_COUNT];
+ private final PulsarAdmin[] pulsarAdmins = new PulsarAdmin[BROKER_COUNT];
@BeforeMethod
void setup() throws Exception {
@@ -222,7 +225,7 @@ public void testLoadReportsWrittenOnMetadataStore() throws Exception {
}
/*
- * tests rankings get updated when we write write the new load reports to the zookeeper on loadbalance root node
+ * tests rankings get updated when we write the new load reports to the zookeeper on load-balance root node
* tests writing pre-configured load report on the zookeeper translates the pre-calculated rankings
*/
@Test
@@ -237,18 +240,15 @@ public void testUpdateLoadReportAndCheckUpdatedRanking() throws Exception {
sru.setCpu(new ResourceUsage(5, 400));
lr.setSystemResourceUsage(sru);
- Whitebox.setInternalState(pulsarServices[0].getLoadManager().get(), "lastLoadReport", lr);
- ResourceLock lock = Whitebox.getInternalState(pulsarServices[i].getLoadManager().get(),
- "brokerLock");
- lock.updateValue(lr).join();
+ FieldUtils.writeField(pulsarServices[0].getLoadManager().get(), "lastLoadReport", lr, true);
+ updateLastReport(pulsarServices[i].getLoadManager().get(), lr);
}
for (int i = 0; i < BROKER_COUNT; i++) {
- Method updateRanking = Whitebox.getMethod(SimpleLoadManagerImpl.class, "updateRanking");
- updateRanking.invoke(pulsarServices[0].getLoadManager().get());
+ MethodUtils.invokeMethod(pulsarServices[0].getLoadManager().get(), true, "updateRanking");
}
- // do lookup for bunch of bundles
+ // do lookup for a bunch of bundles
int totalNamespaces = 200;
Map namespaceOwner = new HashMap<>();
for (int i = 0; i < totalNamespaces; i++) {
@@ -275,6 +275,13 @@ public void testUpdateLoadReportAndCheckUpdatedRanking() throws Exception {
}
}
+ @SuppressWarnings("unchecked")
+ @SneakyThrows
+ private void updateLastReport(LoadManager lm, LoadReport lr){
+ ResourceLock lock = (ResourceLock) FieldUtils.readField(lm, "brokerLock", true);
+ lock.updateValue(lr).join();
+ }
+
private AtomicReference