Fix RSocketRequester API for requests without payload
This commit makes it possible to send requests without requiring to call data(Mono.empty()). It introduces a dedicated MetadataSpec interface and merge ResponseSpec into RequestSpec for more flexibility. Closes gh-23649
This commit is contained in:
@@ -38,7 +38,6 @@ import reactor.test.StepVerifier;
|
||||
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.messaging.rsocket.RSocketRequester.RequestSpec;
|
||||
import org.springframework.messaging.rsocket.RSocketRequester.ResponseSpec;
|
||||
|
||||
import static java.util.concurrent.TimeUnit.MILLISECONDS;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
@@ -87,7 +86,7 @@ public class DefaultRSocketRequesterTests {
|
||||
testSendMono(spec -> spec.data(Mono.delay(MILLIS_10).then(), Void.class), "");
|
||||
}
|
||||
|
||||
private void testSendMono(Function<RequestSpec, ResponseSpec> mapper, String expectedValue) {
|
||||
private void testSendMono(Function<RequestSpec, RequestSpec> mapper, String expectedValue) {
|
||||
mapper.apply(this.requester.route("toA")).send().block(Duration.ofSeconds(5));
|
||||
|
||||
assertThat(this.rsocket.getSavedMethodName()).isEqualTo("fireAndForget");
|
||||
@@ -111,7 +110,7 @@ public class DefaultRSocketRequesterTests {
|
||||
testSendFlux(spec -> spec.data(stringFlux.cast(Object.class), Object.class), values);
|
||||
}
|
||||
|
||||
private void testSendFlux(Function<RequestSpec, ResponseSpec> mapper, String... expectedValues) {
|
||||
private void testSendFlux(Function<RequestSpec, RequestSpec> mapper, String... expectedValues) {
|
||||
this.rsocket.reset();
|
||||
mapper.apply(this.requester.route("toA")).retrieveFlux(String.class).blockLast(Duration.ofSeconds(5));
|
||||
|
||||
|
||||
@@ -74,83 +74,78 @@ class RSocketRequesterExtensionsTests {
|
||||
@Test
|
||||
fun `dataWithType with Publisher`() {
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
val data = mockk<Publisher<String>>()
|
||||
every { requestSpec.data(any<Publisher<String>>(), match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns responseSpec
|
||||
assertThat(requestSpec.dataWithType(data)).isEqualTo(responseSpec)
|
||||
every { requestSpec.data(any<Publisher<String>>(), match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns requestSpec
|
||||
assertThat(requestSpec.dataWithType(data)).isEqualTo(requestSpec)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `dataWithType with Flow`() {
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
val data = mockk<Flow<String>>()
|
||||
every { requestSpec.data(any<Publisher<String>>(), match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns responseSpec
|
||||
assertThat(requestSpec.dataWithType(data)).isEqualTo(responseSpec)
|
||||
every { requestSpec.data(any<Publisher<String>>(), match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns requestSpec
|
||||
assertThat(requestSpec.dataWithType(data)).isEqualTo(requestSpec)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `dataWithType with CompletableFuture`() {
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
val data = mockk<CompletableFuture<String>>()
|
||||
every { requestSpec.data(any<Publisher<String>>(), match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns responseSpec
|
||||
assertThat(requestSpec.dataWithType<String>(data)).isEqualTo(responseSpec)
|
||||
every { requestSpec.data(any<Publisher<String>>(), match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns requestSpec
|
||||
assertThat(requestSpec.dataWithType<String>(data)).isEqualTo(requestSpec)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun dataFlowWithoutType() {
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
every { requestSpec.data(any()) } returns responseSpec
|
||||
assertThat(requestSpec.data(mockk())).isEqualTo(responseSpec)
|
||||
every { requestSpec.data(any()) } returns requestSpec
|
||||
assertThat(requestSpec.data(mockk())).isEqualTo(requestSpec)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun sendAndAwait() {
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
every { responseSpec.send() } returns Mono.empty()
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
every { requestSpec.send() } returns Mono.empty()
|
||||
runBlocking {
|
||||
responseSpec.sendAndAwait()
|
||||
requestSpec.sendAndAwait()
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun retrieveAndAwait() {
|
||||
val response = "foo"
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
every { responseSpec.retrieveMono(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Mono.just("foo")
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
every { requestSpec.retrieveMono(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Mono.just("foo")
|
||||
runBlocking {
|
||||
assertThat(responseSpec.retrieveAndAwait<String>()).isEqualTo(response)
|
||||
assertThat(requestSpec.retrieveAndAwait<String>()).isEqualTo(response)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@ExperimentalCoroutinesApi
|
||||
fun retrieveFlow() {
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
every { responseSpec.retrieveFlux(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Flux.just("foo", "bar")
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
every { requestSpec.retrieveFlux(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Flux.just("foo", "bar")
|
||||
runBlocking {
|
||||
assertThat(responseSpec.retrieveFlow<String>().toList()).contains("foo", "bar")
|
||||
assertThat(requestSpec.retrieveFlow<String>().toList()).contains("foo", "bar")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun retrieveMono() {
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
every { responseSpec.retrieveMono(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Mono.just("foo")
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
every { requestSpec.retrieveMono(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Mono.just("foo")
|
||||
runBlocking {
|
||||
assertThat(responseSpec.retrieveMono<String>().block()).isEqualTo("foo")
|
||||
assertThat(requestSpec.retrieveMono<String>().block()).isEqualTo("foo")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun retrieveFlux() {
|
||||
val responseSpec = mockk<RSocketRequester.ResponseSpec>()
|
||||
every { responseSpec.retrieveFlux(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Flux.just("foo", "bar")
|
||||
val requestSpec = mockk<RSocketRequester.RequestSpec>()
|
||||
every { requestSpec.retrieveFlux(match<ParameterizedTypeReference<*>>(stringTypeRefMatcher)) } returns Flux.just("foo", "bar")
|
||||
runBlocking {
|
||||
assertThat(responseSpec.retrieveFlux<String>().collectList().block()).contains("foo", "bar")
|
||||
assertThat(requestSpec.retrieveFlux<String>().collectList().block()).contains("foo", "bar")
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user