diff --git a/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java b/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java index c13510c4..c09cb2f0 100644 --- a/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java +++ b/spring-pulsar-sample-apps/sample-reactive/src/main/java/org.springframework.pulsar.example/ReactiveSpringPulsarBootApp.java @@ -38,7 +38,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.pulsar.annotation.PulsarListener; import org.springframework.pulsar.core.reactive.DefaultReactivePulsarConsumerFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory; -import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate; +import org.springframework.pulsar.core.reactive.ReactivePulsarTemplate; import reactor.core.publisher.Flux; @@ -59,7 +59,7 @@ public class ReactiveSpringPulsarBootApp { private final Logger logger = LoggerFactory.getLogger(this.getClass()); @Autowired - private ReactivePulsarSenderTemplate reactivePulsarTemplate; + private ReactivePulsarTemplate reactivePulsarTemplate; // TODO remove this once the auto-config is available @Bean @@ -103,7 +103,7 @@ public class ReactiveSpringPulsarBootApp { private final Logger logger = LoggerFactory.getLogger(this.getClass()); @Bean - ApplicationRunner sendSimple(ReactivePulsarSenderTemplate reactivePulsarTemplate) { + ApplicationRunner sendSimple(ReactivePulsarTemplate reactivePulsarTemplate) { return args -> reactivePulsarTemplate .send("sample-reactive-topic2", Flux.range(0, 10).map((i) -> "msg-from-sendSimple-" + i)) .subscribe(); diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java index 6c088341..a5516cd5 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfiguration.java @@ -37,7 +37,7 @@ import org.springframework.pulsar.core.reactive.DefaultReactivePulsarSenderFacto import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarReaderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarSenderFactory; -import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate; +import org.springframework.pulsar.core.reactive.ReactivePulsarTemplate; import com.github.benmanes.caffeine.cache.Caffeine; @@ -48,7 +48,7 @@ import com.github.benmanes.caffeine.cache.Caffeine; * @author Christophe Bornet */ @AutoConfiguration(after = PulsarAutoConfiguration.class) -@ConditionalOnClass({ ReactivePulsarSenderTemplate.class, ReactivePulsarClient.class }) +@ConditionalOnClass({ ReactivePulsarTemplate.class, ReactivePulsarClient.class }) @EnableConfigurationProperties(PulsarReactiveProperties.class) public class PulsarReactiveAutoConfiguration { @@ -111,9 +111,9 @@ public class PulsarReactiveAutoConfiguration { @Bean @ConditionalOnMissingBean - public ReactivePulsarSenderTemplate pulsarReactiveSenderTemplate( + public ReactivePulsarTemplate pulsarReactiveTemplate( ReactivePulsarSenderFactory reactivePulsarSenderFactory) { - return new ReactivePulsarSenderTemplate<>(reactivePulsarSenderFactory); + return new ReactivePulsarTemplate<>(reactivePulsarSenderFactory); } } diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java index 8fed9309..ed3436fa 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactiveAutoConfigurationTests.java @@ -51,7 +51,7 @@ import org.springframework.pulsar.core.reactive.DefaultReactivePulsarSenderFacto import org.springframework.pulsar.core.reactive.ReactivePulsarConsumerFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarReaderFactory; import org.springframework.pulsar.core.reactive.ReactivePulsarSenderFactory; -import org.springframework.pulsar.core.reactive.ReactivePulsarSenderTemplate; +import org.springframework.pulsar.core.reactive.ReactivePulsarTemplate; /** * Autoconfiguration tests for {@link PulsarReactiveAutoConfiguration}. @@ -72,23 +72,23 @@ class PulsarReactiveAutoConfigurationTests { } @Test - void autoConfigurationSkippedWhenReactivePulsarSenderTemplateNotOnClasspath() { - this.contextRunner.withClassLoader(new FilteredClassLoader(ReactivePulsarSenderTemplate.class)).run( + void autoConfigurationSkippedWhenReactivePulsarTemplateNotOnClasspath() { + this.contextRunner.withClassLoader(new FilteredClassLoader(ReactivePulsarTemplate.class)).run( (context) -> assertThat(context).hasNotFailed().doesNotHaveBean(PulsarReactiveAutoConfiguration.class)); } @Test void defaultBeansAreAutoConfigured() { this.contextRunner.run((context) -> assertThat(context).hasNotFailed() - .hasSingleBean(ReactivePulsarSenderTemplate.class).hasSingleBean(ReactivePulsarClient.class) + .hasSingleBean(ReactivePulsarTemplate.class).hasSingleBean(ReactivePulsarClient.class) .hasSingleBean(ProducerCacheProvider.class).hasSingleBean(ReactiveMessageSenderCache.class) - .hasSingleBean(ReactivePulsarSenderFactory.class).getBean(ReactivePulsarSenderTemplate.class)); + .hasSingleBean(ReactivePulsarSenderFactory.class).getBean(ReactivePulsarTemplate.class)); } @ParameterizedTest @ValueSource(classes = { ReactivePulsarClient.class, ProducerCacheProvider.class, ReactiveMessageSenderCache.class, ReactivePulsarSenderFactory.class, ReactivePulsarConsumerFactory.class, ReactivePulsarReaderFactory.class, - ReactivePulsarSenderTemplate.class }) + ReactivePulsarTemplate.class }) void customBeanIsRespected(Class beanClass) { T bean = mock(beanClass); this.contextRunner.withBean(beanClass.getName(), beanClass, () -> bean) @@ -100,7 +100,7 @@ class PulsarReactiveAutoConfigurationTests { ReactivePulsarSenderFactory senderFactory = mock(ReactivePulsarSenderFactory.class); this.contextRunner .withBean("customReactivePulsarSenderFactory", ReactivePulsarSenderFactory.class, () -> senderFactory) - .run((context -> assertThat(context).hasNotFailed().getBean(ReactivePulsarSenderTemplate.class) + .run((context -> assertThat(context).hasNotFailed().getBean(ReactivePulsarTemplate.class) .extracting("reactiveMessageSenderFactory") .asInstanceOf(InstanceOfAssertFactories.type(ReactivePulsarSenderFactory.class)) .isSameAs(senderFactory))); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarOperations.java similarity index 96% rename from spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java rename to spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarOperations.java index 4b2a5188..145a60d9 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderOperations.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarOperations.java @@ -28,7 +28,7 @@ import reactor.core.publisher.Mono; * @param the message payload type * @author Christophe Bornet */ -public interface ReactivePulsarSenderOperations { +public interface ReactivePulsarOperations { /** * Sends a message to the default topic in a reactive manner. @@ -74,7 +74,7 @@ public interface ReactivePulsarSenderOperations { /** * Builder that can be used to configure and send a message. Provides more options - * than the send methods provided by {@link ReactivePulsarSenderOperations}. + * than the send methods provided by {@link ReactivePulsarOperations}. * * @param the message payload type */ diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java similarity index 94% rename from spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java rename to spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java index 307f0918..b1dba2b2 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarSenderTemplate.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplate.java @@ -37,7 +37,7 @@ import reactor.core.publisher.Mono; * @param the message payload type * @author Christophe Bornet */ -public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOperations { +public class ReactivePulsarTemplate implements ReactivePulsarOperations { private final LogAccessor logger = new LogAccessor(this.getClass()); @@ -50,7 +50,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper * @param reactiveMessageSenderFactory the factory used to create the backing Pulsar * reactive senders */ - public ReactivePulsarSenderTemplate(ReactivePulsarSenderFactory reactiveMessageSenderFactory) { + public ReactivePulsarTemplate(ReactivePulsarSenderFactory reactiveMessageSenderFactory) { this.reactiveMessageSenderFactory = reactiveMessageSenderFactory; } @@ -141,7 +141,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper public static class SendMessageBuilderImpl implements SendMessageBuilder { - private final ReactivePulsarSenderTemplate template; + private final ReactivePulsarTemplate template; private final T message; @@ -151,7 +151,7 @@ public class ReactivePulsarSenderTemplate implements ReactivePulsarSenderOper private ReactiveMessageSenderBuilderCustomizer senderCustomizer; - SendMessageBuilderImpl(ReactivePulsarSenderTemplate template, T message) { + SendMessageBuilderImpl(ReactivePulsarTemplate template, T message) { this.template = template; this.message = message; } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java index 2b9e9c51..2b553877 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/reactive/ReactivePulsarTemplateTests.java @@ -45,7 +45,7 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; /** - * Tests for {@link ReactivePulsarSenderTemplate}. + * Tests for {@link ReactivePulsarTemplate}. * * @author Christophe Bornet */ @@ -62,7 +62,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { senderSpec.setTopicName(topic); ReactivePulsarSenderFactory producerFactory = new DefaultReactivePulsarSenderFactory<>(client, senderSpec, null); - ReactivePulsarSenderTemplate pulsarTemplate = new ReactivePulsarSenderTemplate<>(producerFactory); + ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(producerFactory); pulsarTemplate.setSchema(Schema.JSON(Foo.class)); List foos = new ArrayList<>(); @@ -104,7 +104,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { } ReactivePulsarSenderFactory senderFactory = new DefaultReactivePulsarSenderFactory<>(client, senderSpec, null); - ReactivePulsarSenderTemplate pulsarTemplate = new ReactivePulsarSenderTemplate<>(senderFactory); + ReactivePulsarTemplate pulsarTemplate = new ReactivePulsarTemplate<>(senderFactory); Mono sendResponse; if (testArgs.useTemplateSchema) { pulsarTemplate.setSchema(Schema.STRING); @@ -114,7 +114,7 @@ class ReactivePulsarTemplateTests implements PulsarTestContainerSupport { : pulsarTemplate.send(msgPayload); } else { - ReactivePulsarSenderTemplate.SendMessageBuilderImpl messageBuilder = pulsarTemplate + ReactivePulsarTemplate.SendMessageBuilderImpl messageBuilder = pulsarTemplate .newMessage(msgPayload); if (testArgs.useSpecificTopic) { messageBuilder = messageBuilder.withTopic(topic);