Default MimeType selection for RSocketRequester
Remove the dataMimeType argument on connect methods. Applications can still configure it through the ClientRSocketFactory but it shouldn't be necessary as we now choose a default MimeType from the supported encoders and decoders. Add an option to provide the RSocketStrategies instance (vs customizing it) which is expected in Spring config where an RSocketStrategies instance may be shared between client and server setups.
This commit is contained in:
@@ -28,14 +28,8 @@ import org.reactivestreams.Publisher;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import org.springframework.util.MimeTypeUtils;
|
||||
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyInt;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.verifyZeroInteractions;
|
||||
import static org.mockito.Mockito.when;
|
||||
import static org.mockito.ArgumentMatchers.*;
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
/**
|
||||
* Unit tests for {@link DefaultRSocketRequesterBuilder}.
|
||||
@@ -46,6 +40,7 @@ public class DefaultRSocketRequesterBuilderTests {
|
||||
|
||||
private ClientTransport transport;
|
||||
|
||||
|
||||
@Before
|
||||
public void setup() {
|
||||
this.transport = mock(ClientTransport.class);
|
||||
@@ -57,10 +52,10 @@ public class DefaultRSocketRequesterBuilderTests {
|
||||
public void shouldApplyCustomizationsAtSubscription() {
|
||||
Consumer<RSocketFactory.ClientRSocketFactory> factoryConfigurer = mock(Consumer.class);
|
||||
Consumer<RSocketStrategies.Builder> strategiesConfigurer = mock(Consumer.class);
|
||||
Mono<RSocketRequester> requester = RSocketRequester.builder()
|
||||
RSocketRequester.builder()
|
||||
.rsocketFactory(factoryConfigurer)
|
||||
.rsocketStrategies(strategiesConfigurer)
|
||||
.connect(this.transport, MimeTypeUtils.APPLICATION_JSON);
|
||||
.connect(this.transport);
|
||||
verifyZeroInteractions(this.transport, factoryConfigurer, strategiesConfigurer);
|
||||
}
|
||||
|
||||
@@ -69,10 +64,10 @@ public class DefaultRSocketRequesterBuilderTests {
|
||||
public void shouldApplyCustomizations() {
|
||||
Consumer<RSocketFactory.ClientRSocketFactory> factoryConfigurer = mock(Consumer.class);
|
||||
Consumer<RSocketStrategies.Builder> strategiesConfigurer = mock(Consumer.class);
|
||||
RSocketRequester requester = RSocketRequester.builder()
|
||||
RSocketRequester.builder()
|
||||
.rsocketFactory(factoryConfigurer)
|
||||
.rsocketStrategies(strategiesConfigurer)
|
||||
.connect(this.transport, MimeTypeUtils.APPLICATION_JSON)
|
||||
.connect(this.transport)
|
||||
.block();
|
||||
verify(this.transport).connect(anyInt());
|
||||
verify(factoryConfigurer).accept(any(RSocketFactory.ClientRSocketFactory.class));
|
||||
|
||||
@@ -32,7 +32,6 @@ import io.rsocket.RSocket;
|
||||
import io.rsocket.RSocketFactory;
|
||||
import io.rsocket.frame.decoder.PayloadDecoder;
|
||||
import io.rsocket.plugins.RSocketInterceptor;
|
||||
import io.rsocket.transport.netty.client.TcpClientTransport;
|
||||
import io.rsocket.transport.netty.server.CloseableChannel;
|
||||
import io.rsocket.transport.netty.server.TcpServerTransport;
|
||||
import org.junit.After;
|
||||
@@ -60,7 +59,6 @@ import org.springframework.messaging.handler.annotation.MessageExceptionHandler;
|
||||
import org.springframework.messaging.handler.annotation.MessageMapping;
|
||||
import org.springframework.messaging.handler.annotation.Payload;
|
||||
import org.springframework.stereotype.Controller;
|
||||
import org.springframework.util.MimeTypeUtils;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
|
||||
import static org.junit.Assert.*;
|
||||
@@ -78,8 +76,6 @@ public class RSocketBufferLeakTests {
|
||||
|
||||
private static CloseableChannel server;
|
||||
|
||||
private static RSocket client;
|
||||
|
||||
private static RSocketRequester requester;
|
||||
|
||||
|
||||
@@ -96,21 +92,19 @@ public class RSocketBufferLeakTests {
|
||||
.start()
|
||||
.block();
|
||||
|
||||
client = RSocketFactory.connect()
|
||||
.frameDecoder(PayloadDecoder.ZERO_COPY)
|
||||
.addClientPlugin(payloadInterceptor) // intercept outgoing requests
|
||||
.dataMimeType(MimeTypeUtils.TEXT_PLAIN_VALUE)
|
||||
.transport(TcpClientTransport.create("localhost", 7000))
|
||||
.start()
|
||||
requester = RSocketRequester.builder()
|
||||
.rsocketFactory(factory -> {
|
||||
factory.frameDecoder(PayloadDecoder.ZERO_COPY);
|
||||
factory.addClientPlugin(payloadInterceptor); // intercept outgoing requests
|
||||
})
|
||||
.rsocketStrategies(context.getBean(RSocketStrategies.class))
|
||||
.connectTcp("localhost", 7000)
|
||||
.block();
|
||||
|
||||
requester = RSocketRequester.create(
|
||||
client, MimeTypeUtils.TEXT_PLAIN, context.getBean(RSocketStrategies.class));
|
||||
}
|
||||
|
||||
@AfterClass
|
||||
public static void tearDownOnce() {
|
||||
client.dispose();
|
||||
requester.rsocket().dispose();
|
||||
server.dispose();
|
||||
}
|
||||
|
||||
|
||||
@@ -40,7 +40,6 @@ 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 org.springframework.util.MimeTypeUtils;
|
||||
|
||||
import static org.junit.Assert.*;
|
||||
|
||||
@@ -75,11 +74,8 @@ public class RSocketClientToServerIntegrationTests {
|
||||
|
||||
requester = RSocketRequester.builder()
|
||||
.rsocketFactory(factory -> factory.frameDecoder(PayloadDecoder.ZERO_COPY))
|
||||
.rsocketStrategies(strategies -> strategies
|
||||
.decoder(StringDecoder.allMimeTypes())
|
||||
.encoder(CharSequenceEncoder.allMimeTypes())
|
||||
.dataBufferFactory(new NettyDataBufferFactory(PooledByteBufAllocator.DEFAULT)))
|
||||
.connectTcp("localhost", 7000, MimeTypeUtils.TEXT_PLAIN)
|
||||
.rsocketStrategies(context.getBean(RSocketStrategies.class))
|
||||
.connectTcp("localhost", 7000)
|
||||
.block();
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user