diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryShutdownTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryShutdownTest.java index 9dc88efeb13d7..1b863a850cab1 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryShutdownTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryShutdownTest.java @@ -25,7 +25,8 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.doAnswer; -import static org.powermock.api.mockito.PowerMockito.mock; +import static org.mockito.Mockito.mock; + import java.util.Collections; import java.util.Optional; import java.util.UUID; diff --git a/pom.xml b/pom.xml index 136d37fc2a149..46880324d4cd7 100644 --- a/pom.xml +++ b/pom.xml @@ -324,13 +324,7 @@ flexible messaging model and an intuitive client API. org.powermock - powermock-api-mockito2 - ${powermock.version} - - - - org.powermock - powermock-module-testng + powermock-reflect ${powermock.version} @@ -1312,13 +1306,7 @@ flexible messaging model and an intuitive client API. org.powermock - powermock-api-mockito2 - test - - - - org.powermock - powermock-module-testng + powermock-reflect test 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 7b77b1a74f20c..4edfbfd17db02 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 @@ -87,12 +87,12 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; -import static org.powermock.api.mockito.PowerMockito.doAnswer; -import static org.powermock.api.mockito.PowerMockito.doReturn; -import static org.powermock.api.mockito.PowerMockito.mock; -import static org.powermock.api.mockito.PowerMockito.spy; public class TopicsTest extends MockedPulsarServiceBaseTest { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index 7c1cfd01d541d..e2652ba2fa94d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -66,14 +66,12 @@ import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.MockZooKeeper; import org.apache.zookeeper.data.ACL; -import org.powermock.core.classloader.annotations.PowerMockIgnore; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * Base class for all tests that need a Pulsar instance without a ZK and BK cluster. */ -@PowerMockIgnore(value = {"org.slf4j.*", "com.sun.org.apache.xerces.*" }) public abstract class MockedPulsarServiceBaseTest extends TestRetrySupport { protected final String DUMMY_VALUE = "DUMMY_VALUE"; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentMessageFinderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentMessageFinderTest.java index 8be6efff6cf56..e70d0dc2b5b6a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentMessageFinderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentMessageFinderTest.java @@ -21,12 +21,12 @@ import static org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest.retryStrategically; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.clearInvocations; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; -import static org.powermock.api.mockito.PowerMockito.doAnswer; -import static org.powermock.api.mockito.PowerMockito.mock; -import static org.powermock.api.mockito.PowerMockito.spy; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotEquals; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicOwnerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicOwnerTest.java index 2c3248d709a10..8c68c2de91a4f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicOwnerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/TopicOwnerTest.java @@ -20,8 +20,9 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.nullable; -import static org.powermock.api.mockito.PowerMockito.doAnswer; -import static org.powermock.api.mockito.PowerMockito.spy; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.spy; + import com.github.benmanes.caffeine.cache.AsyncLoadingCache; import com.google.common.collect.Sets; import java.util.LinkedHashMap; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReaderTests.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReaderTests.java index 217ac0f290306..52a684bbf27dc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReaderTests.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/streamingdispatch/StreamingEntryReaderTests.java @@ -54,11 +54,11 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.reset; -import static org.powermock.api.mockito.PowerMockito.doAnswer; -import static org.powermock.api.mockito.PowerMockito.mock; -import static org.powermock.api.mockito.PowerMockito.spy; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; /** diff --git a/pulsar-client-tools/src/test/java/org/apache/pulsar/admin/cli/TestCmdPackages.java b/pulsar-client-tools/src/test/java/org/apache/pulsar/admin/cli/TestCmdPackages.java index 6564d3f81c903..23daf81362ce6 100644 --- a/pulsar-client-tools/src/test/java/org/apache/pulsar/admin/cli/TestCmdPackages.java +++ b/pulsar-client-tools/src/test/java/org/apache/pulsar/admin/cli/TestCmdPackages.java @@ -35,7 +35,6 @@ import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.packages.management.core.common.PackageMetadata; -import org.powermock.core.classloader.annotations.PrepareForTest; import org.testng.annotations.BeforeMethod; import org.testng.annotations.DataProvider; import org.testng.annotations.Test; @@ -43,7 +42,6 @@ /** * Unit test for packages commands. */ -@PrepareForTest(CmdPackages.class) public class TestCmdPackages { private PulsarAdmin pulsarAdmin; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AutoClusterFailoverTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AutoClusterFailoverTest.java index 573ddd8086da7..6310d49da74c4 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AutoClusterFailoverTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/AutoClusterFailoverTest.java @@ -30,10 +30,11 @@ import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.awaitility.Awaitility; import org.mockito.Mockito; -import org.powermock.api.mockito.PowerMockito; import org.testng.Assert; import org.testng.annotations.Test; +import static org.mockito.Mockito.mock; + @Test(groups = "broker-impl") public class AutoClusterFailoverTest { @Test @@ -116,7 +117,7 @@ public void testAutoClusterFailoverSwitchWithoutAuthentication() { .build(); AutoClusterFailover autoClusterFailover = Mockito.spy((AutoClusterFailover) provider); - PulsarClientImpl pulsarClient = PowerMockito.mock(PulsarClientImpl.class); + PulsarClientImpl pulsarClient = mock(PulsarClientImpl.class); Mockito.doReturn(false).when(autoClusterFailover).probeAvailable(primary); Mockito.doReturn(true).when(autoClusterFailover).probeAvailable(secondary); Mockito.doReturn(configurationData).when(pulsarClient).getConfiguration(); @@ -172,7 +173,7 @@ public void testAutoClusterFailoverSwitchWithAuthentication() throws IOException .build(); AutoClusterFailover autoClusterFailover = Mockito.spy((AutoClusterFailover) provider); - PulsarClientImpl pulsarClient = PowerMockito.mock(PulsarClientImpl.class); + PulsarClientImpl pulsarClient = mock(PulsarClientImpl.class); Mockito.doReturn(false).when(autoClusterFailover).probeAvailable(primary); Mockito.doReturn(true).when(autoClusterFailover).probeAvailable(secondary); Mockito.doReturn(configurationData).when(pulsarClient).getConfiguration(); @@ -225,7 +226,7 @@ public void testAutoClusterFailoverSwitchTlsTrustStore() throws IOException { .build(); AutoClusterFailover autoClusterFailover = Mockito.spy((AutoClusterFailover) provider); - PulsarClientImpl pulsarClient = PowerMockito.mock(PulsarClientImpl.class); + PulsarClientImpl pulsarClient = mock(PulsarClientImpl.class); Mockito.doReturn(false).when(autoClusterFailover).probeAvailable(primary); Mockito.doReturn(true).when(autoClusterFailover).probeAvailable(secondary); Mockito.doReturn(configurationData).when(pulsarClient).getConfiguration(); diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchMessageContainerImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchMessageContainerImplTest.java index 3d554871141e3..54af93c87e393 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchMessageContainerImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchMessageContainerImplTest.java @@ -18,58 +18,53 @@ */ package org.apache.pulsar.client.impl; -import org.apache.bookkeeper.common.allocator.impl.ByteBufAllocatorBuilderImpl; import org.apache.bookkeeper.common.allocator.impl.ByteBufAllocatorImpl; import org.apache.pulsar.client.api.CompressionType; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.mockito.MockedConstruction; import org.mockito.Mockito; -import org.powermock.api.mockito.PowerMockito; -import org.powermock.core.classloader.annotations.PowerMockIgnore; -import org.powermock.core.classloader.annotations.PrepareForTest; -import org.testng.IObjectFactory; -import org.testng.annotations.ObjectFactory; import org.testng.annotations.Test; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; -@PrepareForTest({ByteBufAllocatorImpl.class, ByteBufAllocatorBuilderImpl.class}) -@PowerMockIgnore({"javax.management.*", "javax.ws.*", "org.apache.logging.log4j.*"}) -public class BatchMessageContainerImplTest { +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.when; - @ObjectFactory - public IObjectFactory getObjectFactory() { - return new org.powermock.modules.testng.PowerMockObjectFactory(); - } +public class BatchMessageContainerImplTest { @Test public void recoveryAfterOom() throws Exception { - final ByteBufAllocatorImpl mockAllocator = PowerMockito.mock(ByteBufAllocatorImpl.class); - PowerMockito.whenNew(ByteBufAllocatorImpl.class).withAnyArguments().thenReturn(mockAllocator); - PowerMockito.when(mockAllocator.buffer(Mockito.anyInt(), Mockito.anyInt())).thenThrow(new OutOfMemoryError("test")).thenReturn(null); - final ProducerImpl producer = Mockito.mock(ProducerImpl.class); - final ProducerConfigurationData producerConfigurationData = new ProducerConfigurationData(); - producerConfigurationData.setCompressionType(CompressionType.NONE); - Mockito.when(producer.getConfiguration()).thenReturn(producerConfigurationData); - final BatchMessageContainerImpl batchMessageContainer = new BatchMessageContainerImpl(); - batchMessageContainer.setProducer(producer); - MessageMetadata messageMetadata1 = new MessageMetadata(); - messageMetadata1.setSequenceId(1L); - messageMetadata1.setProducerName("producer1"); - messageMetadata1.setPublishTime(System.currentTimeMillis()); - ByteBuffer payload1 = ByteBuffer.wrap("payload1".getBytes(StandardCharsets.UTF_8)); - final MessageImpl message1 = MessageImpl.create(messageMetadata1, payload1, Schema.BYTES, null); - batchMessageContainer.add(message1, null); - MessageMetadata messageMetadata2 = new MessageMetadata(); - messageMetadata2.setSequenceId(1L); - messageMetadata2.setProducerName("producer1"); - messageMetadata2.setPublishTime(System.currentTimeMillis()); - ByteBuffer payload2 = ByteBuffer.wrap("payload2".getBytes(StandardCharsets.UTF_8)); - final MessageImpl message2 = MessageImpl.create(messageMetadata2, payload2, Schema.BYTES, null); - // after oom, our add can self-healing, won't throw exception - batchMessageContainer.add(message2, null); + try (MockedConstruction mocked = Mockito.mockConstruction(ByteBufAllocatorImpl.class, + (mockAllocator, context) -> { + doThrow(new OutOfMemoryError("test")).when(mockAllocator).buffer(anyInt(), anyInt()); + })) { + + final ProducerImpl producer = Mockito.mock(ProducerImpl.class); + final ProducerConfigurationData producerConfigurationData = new ProducerConfigurationData(); + producerConfigurationData.setCompressionType(CompressionType.NONE); + when(producer.getConfiguration()).thenReturn(producerConfigurationData); + final BatchMessageContainerImpl batchMessageContainer = new BatchMessageContainerImpl(); + batchMessageContainer.setProducer(producer); + MessageMetadata messageMetadata1 = new MessageMetadata(); + messageMetadata1.setSequenceId(1L); + messageMetadata1.setProducerName("producer1"); + messageMetadata1.setPublishTime(System.currentTimeMillis()); + ByteBuffer payload1 = ByteBuffer.wrap("payload1".getBytes(StandardCharsets.UTF_8)); + final MessageImpl message1 = MessageImpl.create(messageMetadata1, payload1, Schema.BYTES, null); + batchMessageContainer.add(message1, null); + MessageMetadata messageMetadata2 = new MessageMetadata(); + messageMetadata2.setSequenceId(1L); + messageMetadata2.setProducerName("producer1"); + messageMetadata2.setPublishTime(System.currentTimeMillis()); + ByteBuffer payload2 = ByteBuffer.wrap("payload2".getBytes(StandardCharsets.UTF_8)); + final MessageImpl message2 = MessageImpl.create(messageMetadata2, payload2, Schema.BYTES, null); + // after oom, our add can self-healing, won't throw exception + batchMessageContainer.add(message2, null); + } } } diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ControlledClusterFailoverTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ControlledClusterFailoverTest.java index 1fd4a017ec8d1..d7a631c571c08 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ControlledClusterFailoverTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/ControlledClusterFailoverTest.java @@ -27,10 +27,11 @@ import org.asynchttpclient.Request; import org.awaitility.Awaitility; import org.mockito.Mockito; -import org.powermock.api.mockito.PowerMockito; import org.testng.Assert; import org.testng.annotations.Test; +import static org.mockito.Mockito.mock; + @Test(groups = "broker-impl") public class ControlledClusterFailoverTest { @Test @@ -86,7 +87,7 @@ public void testControlledClusterFailoverSwitch() throws IOException { .build(); ControlledClusterFailover controlledClusterFailover = Mockito.spy((ControlledClusterFailover) provider); - PulsarClientImpl pulsarClient = PowerMockito.mock(PulsarClientImpl.class); + PulsarClientImpl pulsarClient = mock(PulsarClientImpl.class); controlledClusterFailover.initialize(pulsarClient); diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java index 49ab9995c0d4a..3d299865ac462 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java @@ -40,12 +40,13 @@ import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.schema.KeyValueEncodingType; import org.testng.Assert; + +import static org.mockito.Mockito.when; import static org.testng.AssertJUnit.fail; import org.testng.annotations.Test; import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; -import static org.powermock.api.mockito.PowerMockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertFalse; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/RoundRobinPartitionMessageRouterImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/RoundRobinPartitionMessageRouterImplTest.java index 6453b861a88e6..dee9144acafb1 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/RoundRobinPartitionMessageRouterImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/RoundRobinPartitionMessageRouterImplTest.java @@ -19,7 +19,7 @@ package org.apache.pulsar.client.impl; import static org.mockito.Mockito.mock; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; @@ -29,6 +29,7 @@ import org.apache.pulsar.client.api.HashingScheme; import org.apache.pulsar.client.api.Message; +import org.mockito.Mockito; import org.testng.annotations.Test; /** @@ -51,7 +52,7 @@ public void testChoosePartitionWithoutKey() { @Test public void testChoosePartitionWithoutKeyWithBatching() { Message msg = mock(Message.class); - when(msg.getKey()).thenReturn(null); + Mockito.when(msg.getKey()).thenReturn(null); // Fake clock, simulate 1 millisecond passes for each invocation Clock clock = new Clock() { diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningAvroSchemaTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningAvroSchemaTest.java index 21c8a0b021cb8..8161621e695a4 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningAvroSchemaTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningAvroSchemaTest.java @@ -28,7 +28,7 @@ import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningKeyValueSchemaTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningKeyValueSchemaTest.java index c40fcf09881ca..726eb558e43ab 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningKeyValueSchemaTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/schema/SupportVersioningKeyValueSchemaTest.java @@ -29,7 +29,7 @@ import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.when; public class SupportVersioningKeyValueSchemaTest { diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocator.java b/pulsar-common/src/main/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocator.java index 31f7bed45d5d2..30c46ac611482 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocator.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocator.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.common.allocator; +import com.google.common.annotations.VisibleForTesting; import io.netty.buffer.ByteBufAllocator; import io.netty.buffer.PooledByteBufAllocator; import java.util.List; @@ -50,19 +51,22 @@ public static void registerOOMListener(Consumer listener) { LISTENERS.add(listener); } - private static final boolean EXIT_ON_OOM; - static { - boolean isPooled = "true".equalsIgnoreCase(System.getProperty(PULSAR_ALLOCATOR_POOLED, "true")); - EXIT_ON_OOM = "true".equalsIgnoreCase(System.getProperty(PULSAR_ALLOCATOR_EXIT_ON_OOM, "false")); - OutOfMemoryPolicy outOfMemoryPolicy = OutOfMemoryPolicy.valueOf( + DEFAULT = createByteBufAllocator(); + } + + @VisibleForTesting + static ByteBufAllocator createByteBufAllocator() { + final boolean isPooled = "true".equalsIgnoreCase(System.getProperty(PULSAR_ALLOCATOR_POOLED, "true")); + final boolean isExitOnOutOfMemory = "true".equalsIgnoreCase( + System.getProperty(PULSAR_ALLOCATOR_EXIT_ON_OOM, "false")); + final OutOfMemoryPolicy outOfMemoryPolicy = OutOfMemoryPolicy.valueOf( System.getProperty(PULSAR_ALLOCATOR_OUT_OF_MEMORY_POLICY, "FallbackToHeap")); - LeakDetectionPolicy leakDetectionPolicy = LeakDetectionPolicy + final LeakDetectionPolicy leakDetectionPolicy = LeakDetectionPolicy .valueOf(System.getProperty(PULSAR_ALLOCATOR_LEAK_DETECTION, "Disabled")); - if (log.isDebugEnabled()) { - log.debug("Is Pooled: {} -- Exit on OOM: {}", isPooled, EXIT_ON_OOM); + log.debug("Is Pooled: {} -- Exit on OOM: {}", isPooled, isExitOnOutOfMemory); } ByteBufAllocatorBuilder builder = ByteBufAllocatorBuilder.create() @@ -78,7 +82,7 @@ public static void registerOOMListener(Consumer listener) { } }); - if (EXIT_ON_OOM) { + if (isExitOnOutOfMemory) { log.info("Exiting JVM process for OOM error: {}", oomException.getMessage(), oomException); Runtime.getRuntime().halt(1); } @@ -90,7 +94,7 @@ public static void registerOOMListener(Consumer listener) { builder.poolingPolicy(PoolingPolicy.UnpooledHeap); } builder.outOfMemoryPolicy(outOfMemoryPolicy); + return builder.build(); - DEFAULT = builder.build(); } } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorDefaultTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorDefaultTest.java index 503a17dca7eeb..92f3250aad423 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorDefaultTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorDefaultTest.java @@ -23,36 +23,44 @@ import org.apache.bookkeeper.common.allocator.LeakDetectionPolicy; import org.apache.bookkeeper.common.allocator.OutOfMemoryPolicy; import org.apache.bookkeeper.common.allocator.PoolingPolicy; -import org.apache.bookkeeper.common.allocator.impl.ByteBufAllocatorBuilderImpl; import org.apache.bookkeeper.common.allocator.impl.ByteBufAllocatorImpl; +import org.mockito.MockedConstruction; import org.mockito.Mockito; -import org.powermock.api.mockito.PowerMockito; -import org.powermock.core.classloader.annotations.PowerMockIgnore; -import org.powermock.core.classloader.annotations.PrepareForTest; -import org.testng.IObjectFactory; -import org.testng.annotations.ObjectFactory; import org.testng.annotations.Test; -@PrepareForTest({ByteBufAllocatorImpl.class, ByteBufAllocatorBuilderImpl.class}) -@PowerMockIgnore({"javax.management.*", "javax.ws.*", "org.apache.logging.log4j.*"}) +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; + @Slf4j public class PulsarByteBufAllocatorDefaultTest { - @ObjectFactory - public IObjectFactory getObjectFactory() { - return new org.powermock.modules.testng.PowerMockObjectFactory(); - } - @Test public void testDefaultConfig() throws Exception { - final ByteBufAllocatorImpl mockAllocator = PowerMockito.mock(ByteBufAllocatorImpl.class); - PowerMockito.whenNew(ByteBufAllocatorImpl.class).withAnyArguments().thenReturn(mockAllocator); - final ByteBufAllocatorImpl byteBufAllocator = (ByteBufAllocatorImpl) PulsarByteBufAllocator.DEFAULT; - // use the variable, in case the compiler optimization - log.trace("{}", byteBufAllocator); - PowerMockito.verifyNew(ByteBufAllocatorImpl.class).withArguments(Mockito.any(ByteBufAllocator.class), Mockito.any(), - Mockito.eq(PoolingPolicy.PooledDirect), Mockito.any(), Mockito.eq(OutOfMemoryPolicy.FallbackToHeap), - Mockito.any(), Mockito.eq(LeakDetectionPolicy.Advanced)); + AtomicBoolean called = new AtomicBoolean(); + try (MockedConstruction mocked = Mockito.mockConstruction(ByteBufAllocatorImpl.class, + (mock, context) -> { + called.set(true); + final List arguments = context.arguments(); + assertTrue(arguments.get(0) instanceof ByteBufAllocator); + assertEquals(arguments.get(2), PoolingPolicy.PooledDirect); + assertEquals(arguments.get(4), OutOfMemoryPolicy.FallbackToHeap); + assertEquals(arguments.get(6), LeakDetectionPolicy.Advanced); + + })) { + final ByteBufAllocatorImpl byteBufAllocator = (ByteBufAllocatorImpl) PulsarByteBufAllocator.DEFAULT; + // use the variable, in case the compiler optimization + log.trace("{}", byteBufAllocator); + + if (!called.get()) { + // maybe PulsarByteBufAllocator static initialization has already been called by a previous test + // let's rerun the same method + PulsarByteBufAllocator.createByteBufAllocator(); + } + assertTrue(called.get()); + } } } \ No newline at end of file diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorOomThrowExceptionTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorOomThrowExceptionTest.java index a23e6899bac94..d8c6c37a3d804 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorOomThrowExceptionTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/allocator/PulsarByteBufAllocatorOomThrowExceptionTest.java @@ -23,41 +23,48 @@ import org.apache.bookkeeper.common.allocator.LeakDetectionPolicy; import org.apache.bookkeeper.common.allocator.OutOfMemoryPolicy; import org.apache.bookkeeper.common.allocator.PoolingPolicy; -import org.apache.bookkeeper.common.allocator.impl.ByteBufAllocatorBuilderImpl; import org.apache.bookkeeper.common.allocator.impl.ByteBufAllocatorImpl; +import org.mockito.MockedConstruction; import org.mockito.Mockito; -import org.powermock.api.mockito.PowerMockito; -import org.powermock.core.classloader.annotations.PowerMockIgnore; -import org.powermock.core.classloader.annotations.PrepareForTest; -import org.testng.IObjectFactory; -import org.testng.annotations.ObjectFactory; import org.testng.annotations.Test; -@PrepareForTest({ByteBufAllocatorImpl.class, ByteBufAllocatorBuilderImpl.class}) -@PowerMockIgnore({"javax.management.*", "javax.ws.*", "org.apache.logging.log4j.*"}) +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; + @Slf4j public class PulsarByteBufAllocatorOomThrowExceptionTest { - @ObjectFactory - public IObjectFactory getObjectFactory() { - return new org.powermock.modules.testng.PowerMockObjectFactory(); - } - @Test public void testDefaultConfig() throws Exception { - try { - System.setProperty("pulsar.allocator.out_of_memory_policy", "ThrowException"); - final ByteBufAllocatorImpl mockAllocator = PowerMockito.mock(ByteBufAllocatorImpl.class); - PowerMockito.whenNew(ByteBufAllocatorImpl.class).withAnyArguments().thenReturn(mockAllocator); + AtomicBoolean called = new AtomicBoolean(); + System.setProperty("pulsar.allocator.out_of_memory_policy", "ThrowException"); + try (MockedConstruction mocked = Mockito.mockConstruction(ByteBufAllocatorImpl.class, + (mock, context) -> { + called.set(true); + final List arguments = context.arguments(); + assertTrue(arguments.get(0) instanceof ByteBufAllocator); + assertEquals(arguments.get(2), PoolingPolicy.PooledDirect); + assertEquals(arguments.get(4), OutOfMemoryPolicy.ThrowException); + assertEquals(arguments.get(6), LeakDetectionPolicy.Advanced); + + })) { final ByteBufAllocatorImpl byteBufAllocator = (ByteBufAllocatorImpl) PulsarByteBufAllocator.DEFAULT; // use the variable, in case the compiler optimization log.trace("{}", byteBufAllocator); - PowerMockito.verifyNew(ByteBufAllocatorImpl.class).withArguments(Mockito.any(ByteBufAllocator.class), Mockito.any(), - Mockito.eq(PoolingPolicy.PooledDirect), Mockito.any(), Mockito.eq(OutOfMemoryPolicy.ThrowException), - Mockito.any(), Mockito.eq(LeakDetectionPolicy.Advanced)); + if (!called.get()) { + // maybe PulsarByteBufAllocator static initialization has already been called by a previous test + // let's rerun the same method + PulsarByteBufAllocator.createByteBufAllocator(); + } + assertTrue(called.get()); } finally { System.clearProperty("pulsar.allocator.out_of_memory_policy"); } + + } } diff --git a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/FunctionResultRouterTest.java b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/FunctionResultRouterTest.java index d5bbf93555549..3bea1bf45bd42 100644 --- a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/FunctionResultRouterTest.java +++ b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/FunctionResultRouterTest.java @@ -19,7 +19,7 @@ package org.apache.pulsar.functions.instance; import static org.mockito.Mockito.mock; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import java.time.Clock; diff --git a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceRunnableTest.java b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceRunnableTest.java index b96e8134bf9d7..e37075b6ab768 100644 --- a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceRunnableTest.java +++ b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/instance/JavaInstanceRunnableTest.java @@ -39,7 +39,7 @@ import java.util.Map; import static org.mockito.Mockito.mock; -import static org.powermock.api.mockito.PowerMockito.when; +import static org.mockito.Mockito.when; public class JavaInstanceRunnableTest { @@ -133,7 +133,7 @@ public void testStatsManagerNull() throws Exception { @Test public void testSinkConfigParsingPreservesOriginalType() throws Exception { SinkSpecOrBuilder sinkSpec = mock(SinkSpecOrBuilder.class); - Mockito.when(sinkSpec.getConfigs()).thenReturn("{\"ttl\": 9223372036854775807}"); + when(sinkSpec.getConfigs()).thenReturn("{\"ttl\": 9223372036854775807}"); Map parsedConfig = new ObjectMapper().readValue(sinkSpec.getConfigs(), new TypeReference>() {}); Assert.assertEquals(parsedConfig.get("ttl").getClass(), Long.class); @@ -143,7 +143,7 @@ public void testSinkConfigParsingPreservesOriginalType() throws Exception { @Test public void testSourceConfigParsingPreservesOriginalType() throws Exception { SourceSpecOrBuilder sourceSpec = mock(SourceSpecOrBuilder.class); - Mockito.when(sourceSpec.getConfigs()).thenReturn("{\"ttl\": 9223372036854775807}"); + when(sourceSpec.getConfigs()).thenReturn("{\"ttl\": 9223372036854775807}"); Map parsedConfig = new ObjectMapper().readValue(sourceSpec.getConfigs(), new TypeReference>() {}); Assert.assertEquals(parsedConfig.get("ttl").getClass(), Long.class); diff --git a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactoryTest.java b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactoryTest.java index dc0119cab4625..40887e225a81f 100644 --- a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactoryTest.java +++ b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactoryTest.java @@ -52,8 +52,8 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; -import static org.powermock.api.mockito.PowerMockito.doNothing; -import static org.powermock.api.mockito.PowerMockito.spy; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.spy; import static org.testng.Assert.assertEquals; import static org.testng.Assert.fail; diff --git a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java index b295cf8a72a80..ef45d7fd1bba3 100644 --- a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java +++ b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java @@ -63,8 +63,8 @@ import static org.apache.pulsar.functions.runtime.RuntimeUtils.FUNCTIONS_INSTANCE_CLASSPATH; import static org.apache.pulsar.functions.utils.FunctionCommon.roundDecimal; -import static org.powermock.api.mockito.PowerMockito.doNothing; -import static org.powermock.api.mockito.PowerMockito.spy; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.spy; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertThrows; diff --git a/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SinkConfigUtilsTest.java b/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SinkConfigUtilsTest.java index 57a5052eda700..c124c3167acc5 100644 --- a/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SinkConfigUtilsTest.java +++ b/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SinkConfigUtilsTest.java @@ -31,7 +31,6 @@ import org.apache.pulsar.config.validation.ConfigValidationAnnotations; import org.apache.pulsar.functions.api.utils.IdentityFunction; import org.apache.pulsar.functions.proto.Function; -import org.powermock.modules.testng.PowerMockTestCase; import org.testng.annotations.Test; import java.io.IOException; @@ -52,7 +51,7 @@ /** * Unit test of {@link SinkConfigUtilsTest}. */ -public class SinkConfigUtilsTest extends PowerMockTestCase { +public class SinkConfigUtilsTest { private ConnectorDefinition defn; diff --git a/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SourceConfigUtilsTest.java b/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SourceConfigUtilsTest.java index a63cf14bc3d57..828cdd2587edb 100644 --- a/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SourceConfigUtilsTest.java +++ b/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/SourceConfigUtilsTest.java @@ -32,7 +32,6 @@ import org.apache.pulsar.functions.proto.Function; import org.apache.pulsar.io.core.BatchSourceTriggerer; import org.apache.pulsar.io.core.SourceContext; -import org.powermock.modules.testng.PowerMockTestCase; import org.testng.annotations.Test; import java.io.IOException; @@ -49,7 +48,7 @@ /** * Unit test of {@link SourceConfigUtilsTest}. */ -public class SourceConfigUtilsTest extends PowerMockTestCase { +public class SourceConfigUtilsTest { private ConnectorDefinition defn; diff --git a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/FunctionRuntimeManagerTest.java b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/FunctionRuntimeManagerTest.java index 26cebc3bf8dc5..0ab0b2728b829 100644 --- a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/FunctionRuntimeManagerTest.java +++ b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/FunctionRuntimeManagerTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.functions.worker; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.CALLS_REAL_METHODS; import static org.mockito.Mockito.any; import static org.mockito.Mockito.anyBoolean; import static org.mockito.Mockito.anyString; @@ -31,6 +32,7 @@ import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import static org.mockito.Mockito.withSettings; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; @@ -72,7 +74,6 @@ import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; -import org.powermock.api.mockito.PowerMockito; import org.testng.annotations.Test; import java.util.HashMap; @@ -938,90 +939,93 @@ public void testFunctionRuntimeFactoryConfigsBackwardsCompatibility() throws Exc WorkerConfig workerConfig = new WorkerConfig(); workerConfig.setKubernetesContainerFactory(kubernetesContainerFactory); - KubernetesRuntimeFactory mockedKubernetesRuntimeFactory = spy(KubernetesRuntimeFactory.class); - doNothing().when(mockedKubernetesRuntimeFactory).initialize( - any(WorkerConfig.class), - any(AuthenticationConfig.class), - any(SecretsProviderConfigurator.class), - any(), - any(), - any() - ); - doNothing().when(mockedKubernetesRuntimeFactory).setupClient(); - doReturn(true).when(mockedKubernetesRuntimeFactory).externallyManaged(); - PowerMockito.whenNew(KubernetesRuntimeFactory.class) - .withNoArguments().thenReturn(mockedKubernetesRuntimeFactory); + try (MockedConstruction mocked = Mockito.mockConstruction(KubernetesRuntimeFactory.class, + withSettings().defaultAnswer(CALLS_REAL_METHODS), + (mockedKubernetesRuntimeFactory, context) -> { + doNothing().when(mockedKubernetesRuntimeFactory).initialize( + any(WorkerConfig.class), + any(AuthenticationConfig.class), + any(SecretsProviderConfigurator.class), + any(), + any(), + any() + ); + doNothing().when(mockedKubernetesRuntimeFactory).setupClient(); + doReturn(true).when(mockedKubernetesRuntimeFactory).externallyManaged(); - FunctionRuntimeManager functionRuntimeManager = new FunctionRuntimeManager( - workerConfig, - mock(PulsarWorkerService.class), - mock(Namespace.class), - mock(MembershipManager.class), - mock(ConnectorsManager.class), - mock(FunctionsManager.class), - mock(FunctionMetaDataManager.class), - mock(WorkerStatsManager.class), - mock(ErrorNotifier.class)); - - KubernetesRuntimeFactory kubernetesRuntimeFactory = (KubernetesRuntimeFactory) functionRuntimeManager.getRuntimeFactory(); - assertEquals(kubernetesRuntimeFactory.getK8Uri(), "k8Uri"); - assertEquals(kubernetesRuntimeFactory.getJobNamespace(), "jobNamespace"); - assertEquals(kubernetesRuntimeFactory.getPulsarDockerImageName(), "pulsarDockerImageName"); - assertEquals(kubernetesRuntimeFactory.getImagePullPolicy(), "imagePullPolicy"); - assertEquals(kubernetesRuntimeFactory.getPulsarRootDir(), "pulsarRootDir"); - - // Test process runtime - - WorkerConfig.ProcessContainerFactory processContainerFactory - = new WorkerConfig.ProcessContainerFactory(); - processContainerFactory.setExtraFunctionDependenciesDir("extraDependenciesDir"); - processContainerFactory.setLogDirectory("logDirectory"); - processContainerFactory.setPythonInstanceLocation("pythonInstanceLocation"); - processContainerFactory.setJavaInstanceJarLocation("javaInstanceJarLocation"); - workerConfig = new WorkerConfig(); - workerConfig.setProcessContainerFactory(processContainerFactory); - - functionRuntimeManager = new FunctionRuntimeManager( - workerConfig, - mock(PulsarWorkerService.class), - mock(Namespace.class), - mock(MembershipManager.class), - mock(ConnectorsManager.class), - mock(FunctionsManager.class), - mock(FunctionMetaDataManager.class), - mock(WorkerStatsManager.class), - mock(ErrorNotifier.class)); - - assertEquals(functionRuntimeManager.getRuntimeFactory().getClass(), ProcessRuntimeFactory.class); - ProcessRuntimeFactory processRuntimeFactory = (ProcessRuntimeFactory) functionRuntimeManager.getRuntimeFactory(); - assertEquals(processRuntimeFactory.getExtraDependenciesDir(), "extraDependenciesDir"); - assertEquals(processRuntimeFactory.getLogDirectory(), "logDirectory/functions"); - assertEquals(processRuntimeFactory.getPythonInstanceFile(), "pythonInstanceLocation"); - assertEquals(processRuntimeFactory.getJavaInstanceJarFile(), "javaInstanceJarLocation"); - - // Test thread runtime - - WorkerConfig.ThreadContainerFactory threadContainerFactory - = new WorkerConfig.ThreadContainerFactory(); - threadContainerFactory.setThreadGroupName("threadGroupName"); - workerConfig = new WorkerConfig(); - workerConfig.setThreadContainerFactory(threadContainerFactory); - workerConfig.setPulsarServiceUrl(PULSAR_SERVICE_URL); + })) { - functionRuntimeManager = new FunctionRuntimeManager( - workerConfig, - mock(PulsarWorkerService.class), - mock(Namespace.class), - mock(MembershipManager.class), - mock(ConnectorsManager.class), - mock(FunctionsManager.class), - mock(FunctionMetaDataManager.class), - mock(WorkerStatsManager.class), - mock(ErrorNotifier.class)); - - assertEquals(functionRuntimeManager.getRuntimeFactory().getClass(), ThreadRuntimeFactory.class); - ThreadRuntimeFactory threadRuntimeFactory = (ThreadRuntimeFactory) functionRuntimeManager.getRuntimeFactory(); - assertEquals(threadRuntimeFactory.getThreadGroup().getName(), "threadGroupName"); + FunctionRuntimeManager functionRuntimeManager = new FunctionRuntimeManager( + workerConfig, + mock(PulsarWorkerService.class), + mock(Namespace.class), + mock(MembershipManager.class), + mock(ConnectorsManager.class), + mock(FunctionsManager.class), + mock(FunctionMetaDataManager.class), + mock(WorkerStatsManager.class), + mock(ErrorNotifier.class)); + + KubernetesRuntimeFactory kubernetesRuntimeFactory = (KubernetesRuntimeFactory) functionRuntimeManager.getRuntimeFactory(); + assertEquals(kubernetesRuntimeFactory.getK8Uri(), "k8Uri"); + assertEquals(kubernetesRuntimeFactory.getJobNamespace(), "jobNamespace"); + assertEquals(kubernetesRuntimeFactory.getPulsarDockerImageName(), "pulsarDockerImageName"); + assertEquals(kubernetesRuntimeFactory.getImagePullPolicy(), "imagePullPolicy"); + assertEquals(kubernetesRuntimeFactory.getPulsarRootDir(), "pulsarRootDir"); + + // Test process runtime + + WorkerConfig.ProcessContainerFactory processContainerFactory + = new WorkerConfig.ProcessContainerFactory(); + processContainerFactory.setExtraFunctionDependenciesDir("extraDependenciesDir"); + processContainerFactory.setLogDirectory("logDirectory"); + processContainerFactory.setPythonInstanceLocation("pythonInstanceLocation"); + processContainerFactory.setJavaInstanceJarLocation("javaInstanceJarLocation"); + workerConfig = new WorkerConfig(); + workerConfig.setProcessContainerFactory(processContainerFactory); + + functionRuntimeManager = new FunctionRuntimeManager( + workerConfig, + mock(PulsarWorkerService.class), + mock(Namespace.class), + mock(MembershipManager.class), + mock(ConnectorsManager.class), + mock(FunctionsManager.class), + mock(FunctionMetaDataManager.class), + mock(WorkerStatsManager.class), + mock(ErrorNotifier.class)); + + assertEquals(functionRuntimeManager.getRuntimeFactory().getClass(), ProcessRuntimeFactory.class); + ProcessRuntimeFactory processRuntimeFactory = (ProcessRuntimeFactory) functionRuntimeManager.getRuntimeFactory(); + assertEquals(processRuntimeFactory.getExtraDependenciesDir(), "extraDependenciesDir"); + assertEquals(processRuntimeFactory.getLogDirectory(), "logDirectory/functions"); + assertEquals(processRuntimeFactory.getPythonInstanceFile(), "pythonInstanceLocation"); + assertEquals(processRuntimeFactory.getJavaInstanceJarFile(), "javaInstanceJarLocation"); + + // Test thread runtime + + WorkerConfig.ThreadContainerFactory threadContainerFactory + = new WorkerConfig.ThreadContainerFactory(); + threadContainerFactory.setThreadGroupName("threadGroupName"); + workerConfig = new WorkerConfig(); + workerConfig.setThreadContainerFactory(threadContainerFactory); + workerConfig.setPulsarServiceUrl(PULSAR_SERVICE_URL); + + functionRuntimeManager = new FunctionRuntimeManager( + workerConfig, + mock(PulsarWorkerService.class), + mock(Namespace.class), + mock(MembershipManager.class), + mock(ConnectorsManager.class), + mock(FunctionsManager.class), + mock(FunctionMetaDataManager.class), + mock(WorkerStatsManager.class), + mock(ErrorNotifier.class)); + + assertEquals(functionRuntimeManager.getRuntimeFactory().getClass(), ThreadRuntimeFactory.class); + ThreadRuntimeFactory threadRuntimeFactory = (ThreadRuntimeFactory) functionRuntimeManager.getRuntimeFactory(); + assertEquals(threadRuntimeFactory.getThreadGroup().getName(), "threadGroupName"); + } } @Test diff --git a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java index 759d1504258ed..8dac514d41b2e 100644 --- a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java +++ b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java @@ -32,7 +32,6 @@ import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.functions.api.Context; import org.apache.pulsar.functions.instance.InstanceConfig; -import org.apache.pulsar.functions.instance.InstanceUtils; import org.apache.pulsar.functions.instance.JavaInstanceRunnable; import org.apache.pulsar.functions.proto.Function; import org.apache.pulsar.functions.proto.InstanceCommunication; @@ -48,13 +47,8 @@ import org.apache.pulsar.functions.worker.WorkerUtils; import org.apache.pulsar.functions.worker.WorkerConfig; import org.glassfish.jersey.media.multipart.FormDataContentDisposition; -import org.powermock.api.mockito.PowerMockito; -import org.powermock.core.classloader.annotations.PowerMockIgnore; -import org.powermock.core.classloader.annotations.PrepareForTest; -import org.testng.IObjectFactory; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; -import org.testng.annotations.ObjectFactory; import org.testng.annotations.Test; import java.io.InputStream; @@ -69,11 +63,11 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.any; import static org.mockito.Mockito.anyInt; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.when; -import static org.powermock.api.mockito.PowerMockito.doReturn; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; diff --git a/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSinkTest.java b/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSinkTest.java index 7f90a87d4c866..5fe9675ae5c8e 100644 --- a/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSinkTest.java +++ b/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSinkTest.java @@ -30,10 +30,8 @@ import org.mockito.Mock; import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; -import org.testng.IObjectFactory; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; -import org.testng.annotations.ObjectFactory; import org.testng.annotations.Test; import java.util.Arrays; @@ -47,6 +45,8 @@ import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.times; + + public class MongoSinkTest { @Mock @@ -73,11 +73,6 @@ public class MongoSinkTest { @Mock private Publisher mockPublisher; - @ObjectFactory - public IObjectFactory getObjectFactory() { - return new org.powermock.modules.testng.PowerMockObjectFactory(); - } - @BeforeMethod public void setUp() { diff --git a/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSourceTest.java b/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSourceTest.java index 1c317f6dafa70..06df54e164912 100644 --- a/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSourceTest.java +++ b/pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSourceTest.java @@ -33,10 +33,8 @@ import org.bson.Document; import org.mockito.Mock; import org.reactivestreams.Subscriber; -import org.testng.IObjectFactory; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; -import org.testng.annotations.ObjectFactory; import org.testng.annotations.Test; import java.util.Map; @@ -76,12 +74,6 @@ public class MongoSourceTest { private Map map; - - @ObjectFactory - public IObjectFactory getObjectFactory() { - return new org.powermock.modules.testng.PowerMockObjectFactory(); - } - @BeforeMethod public void setUp() { diff --git a/testmocks/pom.xml b/testmocks/pom.xml index c7e2100702283..f931052e1bff5 100644 --- a/testmocks/pom.xml +++ b/testmocks/pom.xml @@ -61,8 +61,9 @@ org.powermock - powermock-module-testng + powermock-reflect +