Skip to content

[run-tests] Test polishing #209

New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Merged
merged 1 commit into from
Nov 17, 2022
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -40,59 +40,58 @@
*/
class PulsarListenerTests implements PulsarTestContainerSupport {

static CountDownLatch latch1 = new CountDownLatch(1);
static CountDownLatch latch2 = new CountDownLatch(10);
private static final CountDownLatch LATCH_1 = new CountDownLatch(1);

private static final CountDownLatch LATCH_2 = new CountDownLatch(10);

@Test
void testBasicListener() throws Exception {
void basicPulsarListener() throws Exception {
SpringApplication app = new SpringApplication(BasicListenerConfig.class);
app.setWebApplicationType(WebApplicationType.NONE);

try (ConfigurableApplicationContext context = app
.run("--spring.pulsar.client.serviceUrl=" + PulsarTestContainerSupport.getPulsarBrokerUrl())) {
@SuppressWarnings("unchecked")
final PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
pulsarTemplate.send("hello-pulsar-exclusive", "John Doe");
final boolean await = latch1.await(20, TimeUnit.SECONDS);
assertThat(await).isTrue();
PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
pulsarTemplate.send("plt-topic1", "John Doe");
assertThat(LATCH_1.await(20, TimeUnit.SECONDS)).isTrue();
}
}

@Test
void testBatchListener() throws Exception {
void batchPulsarListener() throws Exception {
SpringApplication app = new SpringApplication(BatchListenerConfig.class);
app.setWebApplicationType(WebApplicationType.NONE);

try (ConfigurableApplicationContext context = app
.run("--spring.pulsar.client.serviceUrl=" + PulsarTestContainerSupport.getPulsarBrokerUrl())) {
@SuppressWarnings("unchecked")
final PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
for (int i = 0; i < 10; i++) {
pulsarTemplate.send("hello-pulsar-exclusive", "John Doe");
pulsarTemplate.send("plt-topic2", "John Doe");
}
final boolean await = latch2.await(10, TimeUnit.SECONDS);
assertThat(await).isTrue();
assertThat(LATCH_2.await(10, TimeUnit.SECONDS)).isTrue();
}
}

@Configuration
@Configuration(proxyBeanMethods = false)
@Import(PulsarAutoConfiguration.class)
static class BasicListenerConfig {

@PulsarListener(subscriptionName = "test-exclusive-sub-1", topics = "hello-pulsar-exclusive")
@PulsarListener(subscriptionName = "plt-subscription1", topics = "plt-topic1")
public void listen(String foo) {
latch1.countDown();
LATCH_1.countDown();
}

}

@Configuration
@Configuration(proxyBeanMethods = false)
@Import(PulsarAutoConfiguration.class)
static class BatchListenerConfig {

@PulsarListener(subscriptionName = "test-exclusive-sub-2", topics = "hello-pulsar-exclusive", batch = true)
@PulsarListener(subscriptionName = "plt-subscription2", topics = "plt-topic2", batch = true)
public void listen(List<String> foo) {
foo.forEach(t -> latch2.countDown());
foo.forEach(t -> LATCH_2.countDown());
}

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,74 +43,72 @@
* Tests for {@link ReactivePulsarListener}.
*
* @author Christophe Bornet
* @author Chris Bono
*/
class ReactivePulsarListenerTests implements PulsarTestContainerSupport {

static CountDownLatch latch1 = new CountDownLatch(1);
static CountDownLatch latch2 = new CountDownLatch(10);
private static final CountDownLatch LATCH1 = new CountDownLatch(1);

private static final CountDownLatch LATCH2 = new CountDownLatch(10);

@Test
void testBasicListener() throws Exception {
void basicListener() throws Exception {
SpringApplication app = new SpringApplication(BasicListenerConfig.class);
app.setWebApplicationType(WebApplicationType.NONE);
app.setAllowCircularReferences(true);

try (ConfigurableApplicationContext context = app
.run("--spring.pulsar.client.serviceUrl=" + PulsarTestContainerSupport.getPulsarBrokerUrl())) {
@SuppressWarnings("unchecked")
final PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
pulsarTemplate.send("hello-pulsar-exclusive", "John Doe");
final boolean await = latch1.await(20, TimeUnit.SECONDS);
assertThat(await).isTrue();
PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
pulsarTemplate.send("rplt-topic1", "John Doe");
assertThat(LATCH1.await(20, TimeUnit.SECONDS)).isTrue();
}
}

@Test
void testFluxListener() throws Exception {
void fluxListener() throws Exception {
SpringApplication app = new SpringApplication(FluxListenerConfig.class);
app.setWebApplicationType(WebApplicationType.NONE);
app.setAllowCircularReferences(true);

try (ConfigurableApplicationContext context = app
.run("--spring.pulsar.client.serviceUrl=" + PulsarTestContainerSupport.getPulsarBrokerUrl())) {
@SuppressWarnings("unchecked")
final PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
for (int i = 0; i < 10; i++) {
pulsarTemplate.send("hello-pulsar-exclusive", "John Doe");
pulsarTemplate.send("rplt-topic2", "John Doe");
}
final boolean await = latch2.await(10, TimeUnit.SECONDS);
assertThat(await).isTrue();
assertThat(LATCH2.await(10, TimeUnit.SECONDS)).isTrue();
}
}

@Configuration
@Import({ PulsarAutoConfiguration.class, PulsarReactiveAutoConfiguration.class })
@Configuration(proxyBeanMethods = false)
@Import({ PulsarAutoConfiguration.class, PulsarReactiveAutoConfiguration.class, ConsumerCustomizerConfig.class })
static class BasicListenerConfig {

@ReactivePulsarListener(subscriptionName = "test-exclusive-sub-1", topics = "hello-pulsar-exclusive",
@ReactivePulsarListener(subscriptionName = "rplt-subscription1", topics = "rplt-topic1",
consumerCustomizer = "consumerCustomizer")
public Mono<Void> listen(String foo) {
latch1.countDown();
LATCH1.countDown();
return Mono.empty();
}

@Bean
ReactiveMessageConsumerBuilderCustomizer<String> consumerCustomizer() {
return b -> b.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest);
}

}

@Configuration
@Import({ PulsarAutoConfiguration.class, PulsarReactiveAutoConfiguration.class })
@Configuration(proxyBeanMethods = false)
@Import({ PulsarAutoConfiguration.class, PulsarReactiveAutoConfiguration.class, ConsumerCustomizerConfig.class })
static class FluxListenerConfig {

@ReactivePulsarListener(subscriptionName = "test-exclusive-sub-2", topics = "hello-pulsar-exclusive",
stream = true, consumerCustomizer = "consumerCustomizer")
@ReactivePulsarListener(subscriptionName = "rplt-subscription2", topics = "rplt-topic2", stream = true,
consumerCustomizer = "consumerCustomizer")
public Flux<MessageResult<Void>> listen(Flux<Message<String>> messages) {
return messages.doOnNext(t -> latch2.countDown()).map(m -> MessageResult.acknowledge(m.getMessageId()));
return messages.doOnNext(t -> LATCH2.countDown()).map(m -> MessageResult.acknowledge(m.getMessageId()));
}

}

@Configuration(proxyBeanMethods = false)
static class ConsumerCustomizerConfig {

@Bean
ReactiveMessageConsumerBuilderCustomizer<String> consumerCustomizer() {
return b -> b.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest);
Expand Down