Smoke test for Spring Pulsar reactive native app
This commit is contained in:
24
integration/spring-pulsar-reactive/build.gradle
Normal file
24
integration/spring-pulsar-reactive/build.gradle
Normal file
@@ -0,0 +1,24 @@
|
||||
plugins {
|
||||
id 'java'
|
||||
id 'org.springframework.boot'
|
||||
id 'org.springframework.aot.smoke-test'
|
||||
id 'org.graalvm.buildtools.native'
|
||||
}
|
||||
|
||||
ext {
|
||||
set('springPulsarVersion', "0.1.0-SNAPSHOT")
|
||||
}
|
||||
|
||||
repositories {
|
||||
maven { url 'https://repository.apache.org/content/repositories/snapshots' }
|
||||
}
|
||||
|
||||
dependencies {
|
||||
implementation(platform(org.springframework.boot.gradle.plugin.SpringBootPlugin.BOM_COORDINATES))
|
||||
implementation("org.springframework.pulsar:spring-pulsar-reactive-spring-boot-starter:${springPulsarVersion}")
|
||||
|
||||
testImplementation("org.springframework.boot:spring-boot-starter-test")
|
||||
|
||||
appTestImplementation(project(":aot-smoke-test-support"))
|
||||
appTestImplementation("org.awaitility:awaitility:4.2.0")
|
||||
}
|
||||
8
integration/spring-pulsar-reactive/docker-compose.yml
Normal file
8
integration/spring-pulsar-reactive/docker-compose.yml
Normal file
@@ -0,0 +1,8 @@
|
||||
version: '3'
|
||||
services:
|
||||
pulsar:
|
||||
image: apachepulsar/pulsar:2.10.2
|
||||
ports:
|
||||
- '8080'
|
||||
- '6650'
|
||||
command: bin/pulsar standalone
|
||||
@@ -0,0 +1,21 @@
|
||||
package com.example.pulsar;
|
||||
|
||||
import org.awaitility.Awaitility;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.aot.smoketest.support.assertj.AssertableOutput;
|
||||
import org.springframework.aot.smoketest.support.junit.ApplicationTest;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
@ApplicationTest
|
||||
public class SpringPulsarReactiveApplicationAotTests {
|
||||
|
||||
@Test
|
||||
void reactivePulsarListenerMethodReceivesMessage(AssertableOutput output) {
|
||||
Awaitility.await().atMost(Duration.ofSeconds(30))
|
||||
.untilAsserted(() -> assertThat(output).hasSingleLineContaining("Message Received: sample-message-50"));
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,36 @@
|
||||
package com.example.pulsar;
|
||||
|
||||
import org.springframework.boot.ApplicationRunner;
|
||||
import org.springframework.boot.SpringApplication;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.pulsar.reactive.config.annotation.ReactivePulsarListener;
|
||||
import org.springframework.pulsar.reactive.core.ReactivePulsarTemplate;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
@SpringBootApplication
|
||||
public class SpringPulsarReactiveApplication {
|
||||
|
||||
public static void main(String[] args) {
|
||||
SpringApplication.run(SpringPulsarReactiveApplication.class, args);
|
||||
}
|
||||
|
||||
@Bean
|
||||
ApplicationRunner sendMessageToTopicOnAppStartup(ReactivePulsarTemplate<String> reactivePulsarTemplate) {
|
||||
String topic = "graalvm-demo-topic-reactive";
|
||||
return args -> {
|
||||
Flux.range(0, 100).map((i) -> "sample-message-" + i)
|
||||
.as(messages -> reactivePulsarTemplate.send(topic, messages)).subscribe();
|
||||
};
|
||||
}
|
||||
|
||||
@ReactivePulsarListener(subscriptionName = "graalvm-demo-subscription-reactive",
|
||||
topics = "graalvm-demo-topic-reactive")
|
||||
public Mono<Void> listenReactive(String message) {
|
||||
System.out.println("Message Received: " + message);
|
||||
return Mono.empty();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,2 @@
|
||||
spring.pulsar.client.service-url=pulsar://${PULSAR_HOST:localhost}:${PULSAR_PORT_6650:6650}
|
||||
spring.pulsar.administration.service-url=http://${PULSAR_HOST:localhost}:${PULSAR_PORT_8080:8080}
|
||||
Reference in New Issue
Block a user