diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/SpringPulsarBootAppSanityTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/SpringPulsarBootAppSanityTests.java index 42ebd788..4687a82c 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/SpringPulsarBootAppSanityTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/SpringPulsarBootAppSanityTests.java @@ -20,7 +20,6 @@ import static org.assertj.core.api.Assertions.assertThat; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClientException; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.ObjectProvider; @@ -44,7 +43,6 @@ import org.springframework.web.bind.annotation.RestController; * @author Chris Bono */ @SpringBootTest(classes = SpringPulsarBootTestApp.class, webEnvironment = WebEnvironment.RANDOM_PORT) -@Disabled("temporarily") class SpringPulsarBootAppSanityTests implements PulsarTestContainerSupport { @DynamicPropertySource diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java index f91c94f9..fc2b9de7 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/CachingPulsarProducerFactory.java @@ -26,17 +26,17 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.function.Consumer; -import org.aopalliance.intercept.MethodInterceptor; -import org.aopalliance.intercept.MethodInvocation; import org.apache.commons.logging.LogFactory; +import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.ProducerStats; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.TypedMessageBuilder; +import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.common.protocol.schema.SchemaHash; -import org.springframework.aop.framework.AopProxyUtils; -import org.springframework.aop.framework.ProxyFactory; import org.springframework.beans.factory.DisposableBean; import org.springframework.core.log.LogAccessor; import org.springframework.lang.Nullable; @@ -106,7 +106,7 @@ public class CachingPulsarProducerFactory extends DefaultPulsarProducerFactor @Nullable Collection encryptionKeys, @Nullable List> customizers) { try { Producer producer = super.doCreateProducer(schema, topic, encryptionKeys, customizers); - return wrapProducerWithCloseCallback(producer, + return new ProducerWithCloseCallback<>(producer, (p) -> this.logger .trace(() -> String.format("Client closed producer %s but will skip actual closing", ProducerUtils.formatProducer(producer)))); @@ -116,27 +116,6 @@ public class CachingPulsarProducerFactory extends DefaultPulsarProducerFactor } } - @SuppressWarnings("unchecked") - private Producer wrapProducerWithCloseCallback(Producer producer, Consumer> closeCallback) { - ProxyFactory factory = new ProxyFactory(producer); - factory.addAdvice(new MethodInterceptor() { - @Nullable - @Override - public Object invoke(MethodInvocation invocation) throws Throwable { - if (invocation.getMethod().getName().equals("close")) { - closeCallback.accept((Producer) invocation.getThis()); - return null; - } - if (invocation.getMethod().getName().equals("closeAsync")) { - closeCallback.accept((Producer) invocation.getThis()); - return CompletableFuture.completedFuture(null); - } - return invocation.proceed(); - } - }); - return (Producer) factory.getProxy(); - } - @Override public void destroy() { this.producerCache.asMap().forEach((producerCacheKey, producer) -> { @@ -147,7 +126,10 @@ public class CachingPulsarProducerFactory extends DefaultPulsarProducerFactor @SuppressWarnings("unchecked") private void closeProducer(Producer producer) { - Producer actualProducer = (Producer) AopProxyUtils.getSingletonTarget(producer); + Producer actualProducer = null; + if (producer instanceof ProducerWithCloseCallback wrappedProducer) { + actualProducer = wrappedProducer.getActualProducer(); + } if (actualProducer == null) { this.logger.warn(() -> String.format("Unable to get actual producer for %s - will skip closing it", ProducerUtils.formatProducer(producer))); @@ -216,4 +198,108 @@ public class CachingPulsarProducerFactory extends DefaultPulsarProducerFactor } + /** + * A producer that does not actually close when the user calls + * {@link Producer#close()}. + * + * @param producer type. + */ + static class ProducerWithCloseCallback implements Producer { + + private final Producer producer; + + private final Consumer> closeCallback; + + ProducerWithCloseCallback(Producer producer, Consumer> closeCallback) { + this.producer = producer; + this.closeCallback = closeCallback; + } + + public Producer getActualProducer() { + return this.producer; + } + + @Override + public String getTopic() { + return this.producer.getTopic(); + } + + @Override + public String getProducerName() { + return this.producer.getProducerName(); + } + + @Override + public MessageId send(T message) throws PulsarClientException { + return this.producer.send(message); + } + + @Override + public CompletableFuture sendAsync(T message) { + return this.producer.sendAsync(message); + } + + @Override + public void flush() throws PulsarClientException { + this.producer.flush(); + } + + @Override + public CompletableFuture flushAsync() { + return this.producer.flushAsync(); + } + + @Override + public TypedMessageBuilder newMessage() { + return this.producer.newMessage(); + } + + @Override + public TypedMessageBuilder newMessage(Schema schema) { + return this.producer.newMessage(schema); + } + + @Override + public TypedMessageBuilder newMessage(Transaction txn) { + return this.producer.newMessage(txn); + } + + @Override + public long getLastSequenceId() { + return this.producer.getLastSequenceId(); + } + + @Override + public ProducerStats getStats() { + return this.producer.getStats(); + } + + @Override + public void close() throws PulsarClientException { + this.closeCallback.accept(this.producer); + } + + @Override + public CompletableFuture closeAsync() { + this.closeCallback.accept(this.producer); + return CompletableFuture.completedFuture(null); + } + + @Override + public boolean isConnected() { + return this.producer.isConnected(); + } + + @Override + public long getLastDisconnectedTimestamp() { + return this.producer.getLastDisconnectedTimestamp(); + } + + @Override + public int getNumOfPartitions() { + return this.producer.getNumOfPartitions(); + } + + } + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java index bb54c8f3..89dcc3ea 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/CachingPulsarProducerFactoryTests.java @@ -46,8 +46,8 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; -import org.springframework.aop.framework.AopProxyUtils; import org.springframework.pulsar.core.CachingPulsarProducerFactory.ProducerCacheKey; +import org.springframework.pulsar.core.CachingPulsarProducerFactory.ProducerWithCloseCallback; import org.springframework.test.util.ReflectionTestUtils; import org.springframework.util.ObjectUtils; @@ -86,19 +86,19 @@ class CachingPulsarProducerFactoryTests extends PulsarProducerFactoryTests { Cache, Producer> producerCache = getAssertedProducerCache(producerFactory, Collections.singletonList(cacheKey)); - Producer cachedProducerProxy = producerCache.asMap().get(cacheKey); - assertThat(cachedProducerProxy).isSameAs(producer1); + Producer cachedProducerWrapper = producerCache.asMap().get(cacheKey); + assertThat(cachedProducerWrapper).isSameAs(producer1); } @Test - void cachedProducerIsCloseSafeProxy() throws PulsarClientException { + void cachedProducerIsCloseSafeWrapper() throws PulsarClientException { PulsarProducerFactory producerFactory = producerFactory(pulsarClient, Collections.emptyMap()); - Producer proxyProducer = producerFactory.createProducer(schema, "topic1"); - Producer actualProducer = actualProducerFrom(proxyProducer); + Producer wrappedProducer = producerFactory.createProducer(schema, "topic1"); + Producer actualProducer = actualProducerFrom(wrappedProducer); assertThat(actualProducer.isConnected()).isTrue(); - proxyProducer.close(); + wrappedProducer.close(); assertThat(actualProducer.isConnected()).isTrue(); actualProducer.close(); assertThat(actualProducer.isConnected()).isFalse(); @@ -222,10 +222,9 @@ class CachingPulsarProducerFactoryTests extends PulsarProducerFactoryTests { } @SuppressWarnings("unchecked") - private Producer actualProducerFrom(Producer proxyProducer) { - Producer actualProducer = (Producer) AopProxyUtils.getSingletonTarget(proxyProducer); - assertThat(actualProducer).isNotNull(); - return actualProducer; + private Producer actualProducerFrom(Producer wrappedProducer) { + assertThat(wrappedProducer).isInstanceOf(ProducerWithCloseCallback.class); + return ((ProducerWithCloseCallback) wrappedProducer).getActualProducer(); } @SuppressWarnings("unchecked")