Remove MessageHandlerAcceptor sub-class
This commit removes the MessageHandlerAcceptor sub-class of RSocketMessageHandler, and rather than implementing directly the contracts for RSocket client and server acceptors, RSocketMessageHandler now exposes clientAcceptor() and serverAcceptor() methods that return the required adapter instances. This provides better separation between the RSocketMessageHandler and the RSocket adapter code, and also avoids implementing generic interfaces like the BiFunction required for the client acceptor.
This commit is contained in:
@@ -87,7 +87,7 @@ public class RSocketBufferLeakTests {
|
||||
server = RSocketFactory.receive()
|
||||
.frameDecoder(PayloadDecoder.ZERO_COPY)
|
||||
.addServerPlugin(payloadInterceptor) // intercept responding
|
||||
.acceptor(context.getBean(MessageHandlerAcceptor.class))
|
||||
.acceptor(context.getBean(RSocketMessageHandler.class).serverAcceptor())
|
||||
.transport(TcpServerTransport.create("localhost", 7000))
|
||||
.start()
|
||||
.block();
|
||||
@@ -214,10 +214,10 @@ public class RSocketBufferLeakTests {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public MessageHandlerAcceptor messageHandlerAcceptor() {
|
||||
MessageHandlerAcceptor acceptor = new MessageHandlerAcceptor();
|
||||
acceptor.setRSocketStrategies(rsocketStrategies());
|
||||
return acceptor;
|
||||
public RSocketMessageHandler messageHandler() {
|
||||
RSocketMessageHandler handler = new RSocketMessageHandler();
|
||||
handler.setRSocketStrategies(rsocketStrategies());
|
||||
return handler;
|
||||
}
|
||||
|
||||
@Bean
|
||||
|
||||
@@ -67,7 +67,7 @@ public class RSocketClientToServerIntegrationTests {
|
||||
server = RSocketFactory.receive()
|
||||
.addServerPlugin(interceptor)
|
||||
.frameDecoder(PayloadDecoder.ZERO_COPY)
|
||||
.acceptor(context.getBean(MessageHandlerAcceptor.class))
|
||||
.acceptor(context.getBean(RSocketMessageHandler.class).serverAcceptor())
|
||||
.transport(TcpServerTransport.create("localhost", 7000))
|
||||
.start()
|
||||
.block();
|
||||
@@ -257,10 +257,10 @@ public class RSocketClientToServerIntegrationTests {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public MessageHandlerAcceptor messageHandlerAcceptor() {
|
||||
MessageHandlerAcceptor acceptor = new MessageHandlerAcceptor();
|
||||
acceptor.setRSocketStrategies(rsocketStrategies());
|
||||
return acceptor;
|
||||
public RSocketMessageHandler messageHandler() {
|
||||
RSocketMessageHandler handler = new RSocketMessageHandler();
|
||||
handler.setRSocketStrategies(rsocketStrategies());
|
||||
return handler;
|
||||
}
|
||||
|
||||
@Bean
|
||||
|
||||
@@ -65,7 +65,7 @@ public class RSocketServerToClientIntegrationTests {
|
||||
|
||||
server = RSocketFactory.receive()
|
||||
.frameDecoder(PayloadDecoder.ZERO_COPY)
|
||||
.acceptor(context.getBean("serverAcceptor", MessageHandlerAcceptor.class))
|
||||
.acceptor(context.getBean("serverMessageHandler", RSocketMessageHandler.class).serverAcceptor())
|
||||
.transport(TcpServerTransport.create("localhost", 7000))
|
||||
.start()
|
||||
.block();
|
||||
@@ -110,7 +110,7 @@ public class RSocketServerToClientIntegrationTests {
|
||||
.dataMimeType("text/plain")
|
||||
.setupPayload(DefaultPayload.create("", destination))
|
||||
.frameDecoder(PayloadDecoder.ZERO_COPY)
|
||||
.acceptor(context.getBean("clientAcceptor", MessageHandlerAcceptor.class))
|
||||
.acceptor(context.getBean("clientMessageHandler", RSocketMessageHandler.class).clientAcceptor())
|
||||
.transport(TcpClientTransport.create("localhost", 7000))
|
||||
.start()
|
||||
.block();
|
||||
@@ -260,17 +260,16 @@ public class RSocketServerToClientIntegrationTests {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public MessageHandlerAcceptor clientAcceptor() {
|
||||
MessageHandlerAcceptor acceptor = new MessageHandlerAcceptor();
|
||||
acceptor.setHandlers(Collections.singletonList(clientHandler()));
|
||||
acceptor.setAutoDetectDisabled();
|
||||
acceptor.setRSocketStrategies(rsocketStrategies());
|
||||
return acceptor;
|
||||
public RSocketMessageHandler clientMessageHandler() {
|
||||
RSocketMessageHandler handler = new RSocketMessageHandler();
|
||||
handler.setHandlers(Collections.singletonList(clientHandler()));
|
||||
handler.setRSocketStrategies(rsocketStrategies());
|
||||
return handler;
|
||||
}
|
||||
|
||||
@Bean
|
||||
public MessageHandlerAcceptor serverAcceptor() {
|
||||
MessageHandlerAcceptor handler = new MessageHandlerAcceptor();
|
||||
public RSocketMessageHandler serverMessageHandler() {
|
||||
RSocketMessageHandler handler = new RSocketMessageHandler();
|
||||
handler.setRSocketStrategies(rsocketStrategies());
|
||||
return handler;
|
||||
}
|
||||
|
||||
@@ -16,8 +16,6 @@
|
||||
|
||||
package org.springframework.messaging.rsocket
|
||||
|
||||
import java.time.Duration
|
||||
|
||||
import io.netty.buffer.PooledByteBufAllocator
|
||||
import io.rsocket.RSocketFactory
|
||||
import io.rsocket.frame.decoder.PayloadDecoder
|
||||
@@ -31,9 +29,6 @@ import kotlinx.coroutines.flow.map
|
||||
import org.junit.AfterClass
|
||||
import org.junit.BeforeClass
|
||||
import org.junit.Test
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.test.StepVerifier
|
||||
|
||||
import org.springframework.context.annotation.AnnotationConfigApplicationContext
|
||||
import org.springframework.context.annotation.Bean
|
||||
import org.springframework.context.annotation.Configuration
|
||||
@@ -43,6 +38,9 @@ import org.springframework.core.io.buffer.NettyDataBufferFactory
|
||||
import org.springframework.messaging.handler.annotation.MessageExceptionHandler
|
||||
import org.springframework.messaging.handler.annotation.MessageMapping
|
||||
import org.springframework.stereotype.Controller
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.test.StepVerifier
|
||||
import java.time.Duration
|
||||
|
||||
/**
|
||||
* Coroutines server-side handling of RSocket requests.
|
||||
@@ -167,10 +165,10 @@ class RSocketClientToServerCoroutinesIntegrationTests {
|
||||
}
|
||||
|
||||
@Bean
|
||||
open fun messageHandlerAcceptor(): MessageHandlerAcceptor {
|
||||
val acceptor = MessageHandlerAcceptor()
|
||||
acceptor.rSocketStrategies = rsocketStrategies()
|
||||
return acceptor
|
||||
open fun messageHandler(): RSocketMessageHandler {
|
||||
val handler = RSocketMessageHandler()
|
||||
handler.rSocketStrategies = rsocketStrategies()
|
||||
return handler
|
||||
}
|
||||
|
||||
@Bean
|
||||
@@ -202,7 +200,7 @@ class RSocketClientToServerCoroutinesIntegrationTests {
|
||||
server = RSocketFactory.receive()
|
||||
.addServerPlugin(interceptor)
|
||||
.frameDecoder(PayloadDecoder.ZERO_COPY)
|
||||
.acceptor(context.getBean(MessageHandlerAcceptor::class.java))
|
||||
.acceptor(context.getBean(RSocketMessageHandler::class.java).serverAcceptor())
|
||||
.transport(TcpServerTransport.create("localhost", 7000))
|
||||
.start()
|
||||
.block()!!
|
||||
|
||||
Reference in New Issue
Block a user