Reflect grouping in directory structure
This commit is contained in:
@@ -0,0 +1,42 @@
|
||||
package com.example.rsocket;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
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 static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
@ApplicationTest
|
||||
class RSocketApplicationAotTests {
|
||||
|
||||
@Test
|
||||
void messageIsReceivedAndAnswered(AssertableOutput output) {
|
||||
Awaitility.await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
|
||||
assertThat(output).hasSingleLineContaining("Server: message(): Message{origin='client', message='Hello!'}")
|
||||
.hasSingleLineContaining("Client: message(): Message{origin='server', message='Hello!'}");
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void reactiveMessageIsReceivedAndAnswered(AssertableOutput output) {
|
||||
Awaitility.await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
|
||||
assertThat(output)
|
||||
.hasSingleLineContaining("Server: reactiveMessage: Message{origin='client', message='Hello!'}")
|
||||
.hasSingleLineContaining("Client: reactiveMessage(): Message{origin='server', message='Hello!'}");
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void messageRecordIsReceivedAndAnswered(AssertableOutput output) {
|
||||
Awaitility.await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
|
||||
assertThat(output)
|
||||
.hasSingleLineContaining("Server: messageRecord(): MessageRecord[origin=client, message=Hello!]")
|
||||
.hasSingleLineContaining("Client: messageRecord(): MessageRecord[origin=server, message=Hello!]");
|
||||
});
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
package com.example.rsocket;
|
||||
|
||||
import org.springframework.boot.SpringApplication;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
|
||||
@SpringBootApplication
|
||||
public class RSocketApplication {
|
||||
|
||||
public static void main(String[] args) {
|
||||
SpringApplication.run(RSocketApplication.class, args);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package com.example.rsocket;
|
||||
|
||||
import com.example.rsocket.dto.Message;
|
||||
import com.example.rsocket.dto.MessageRecord;
|
||||
|
||||
import org.springframework.boot.CommandLineRunner;
|
||||
import org.springframework.boot.rsocket.context.RSocketServerInitializedEvent;
|
||||
import org.springframework.context.event.EventListener;
|
||||
import org.springframework.messaging.rsocket.RSocketRequester;
|
||||
import org.springframework.messaging.rsocket.RSocketRequester.Builder;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
@Component
|
||||
class RSocketClient implements CommandLineRunner {
|
||||
|
||||
private int port;
|
||||
|
||||
private final RSocketRequester.Builder builder;
|
||||
|
||||
RSocketClient(Builder builder) {
|
||||
this.builder = builder;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void run(String... args) {
|
||||
RSocketRequester requester = this.builder.tcp("localhost", this.port);
|
||||
|
||||
message(requester);
|
||||
reactiveMessage(requester);
|
||||
messageRecord(requester);
|
||||
}
|
||||
|
||||
@EventListener
|
||||
public void onRSocketServerStart(RSocketServerInitializedEvent event) {
|
||||
this.port = event.getServer().address().getPort();
|
||||
}
|
||||
|
||||
private void message(RSocketRequester requester) {
|
||||
Message message = requester.route("message").data(new Message("client", "Hello!")).retrieveMono(Message.class)
|
||||
.block();
|
||||
System.out.printf("Client: message(): %s%n", message);
|
||||
}
|
||||
|
||||
private void reactiveMessage(RSocketRequester requester) {
|
||||
Message message = requester.route("reactive-message").data(new Message("client", "Hello!"))
|
||||
.retrieveMono(Message.class).block();
|
||||
System.out.printf("Client: reactiveMessage(): %s%n", message);
|
||||
}
|
||||
|
||||
private void messageRecord(RSocketRequester requester) {
|
||||
MessageRecord messageRecord = requester.route("message-record").data(new MessageRecord("client", "Hello!"))
|
||||
.retrieveMono(MessageRecord.class).block();
|
||||
System.out.printf("Client: messageRecord(): %s%n", messageRecord);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
package com.example.rsocket;
|
||||
|
||||
import com.example.rsocket.dto.Message;
|
||||
import com.example.rsocket.dto.MessageRecord;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.messaging.handler.annotation.MessageMapping;
|
||||
import org.springframework.stereotype.Controller;
|
||||
|
||||
@Controller
|
||||
public class RSocketController {
|
||||
|
||||
@MessageMapping("message")
|
||||
public Message message(Message request) {
|
||||
System.out.printf("Server: message(): %s%n", request);
|
||||
return new Message("server", request.getMessage());
|
||||
}
|
||||
|
||||
@MessageMapping("reactive-message")
|
||||
public Mono<Message> reactiveMessage(Message request) {
|
||||
System.out.printf("Server: reactiveMessage: %s%n", request);
|
||||
return Mono.just(new Message("server", request.getMessage()));
|
||||
}
|
||||
|
||||
@MessageMapping("message-record")
|
||||
public MessageRecord messageRecord(MessageRecord request) {
|
||||
System.out.printf("Server: messageRecord(): %s%n", request);
|
||||
return new MessageRecord("server", request.message());
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
package com.example.rsocket.dto;
|
||||
|
||||
public class Message {
|
||||
|
||||
private String origin;
|
||||
|
||||
private String message;
|
||||
|
||||
public Message() {
|
||||
}
|
||||
|
||||
public Message(String origin, String message) {
|
||||
this.origin = origin;
|
||||
this.message = message;
|
||||
}
|
||||
|
||||
public String getOrigin() {
|
||||
return origin;
|
||||
}
|
||||
|
||||
public void setOrigin(String origin) {
|
||||
this.origin = origin;
|
||||
}
|
||||
|
||||
public String getMessage() {
|
||||
return message;
|
||||
}
|
||||
|
||||
public void setMessage(String message) {
|
||||
this.message = message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "Message{" + "origin='" + origin + '\'' + ", message='" + message + '\'' + '}';
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,4 @@
|
||||
package com.example.rsocket.dto;
|
||||
|
||||
public record MessageRecord(String origin, String message) {
|
||||
}
|
||||
@@ -0,0 +1 @@
|
||||
spring.rsocket.server.port=0
|
||||
Reference in New Issue
Block a user