Improve RSocketRequester.ResponseSpec Kotlin extensions
This commit adds retrieveMono and retrieveFlux reified variants, and turns dataFlow(flow: Flow) extension into a general purpose reified data(producer: Any) one. Closes gh-23164
This commit is contained in:
@@ -22,8 +22,9 @@ import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.reactive.awaitFirstOrNull
|
||||
import kotlinx.coroutines.reactive.awaitSingle
|
||||
import kotlinx.coroutines.reactive.flow.asFlow
|
||||
import kotlinx.coroutines.reactive.flow.asPublisher
|
||||
import org.springframework.core.ParameterizedTypeReference
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.core.publisher.Mono
|
||||
import java.net.URI
|
||||
|
||||
/**
|
||||
@@ -53,16 +54,18 @@ suspend fun RSocketRequester.Builder.connectTcpAndAwait(host: String, port: Int)
|
||||
suspend fun RSocketRequester.Builder.connectWebSocketAndAwait(uri: URI): RSocketRequester =
|
||||
connectWebSocket(uri).awaitSingle()
|
||||
|
||||
|
||||
/**
|
||||
* Kotlin [Flow] variant of [RSocketRequester.RequestSpec.data].
|
||||
* Extension for [RSocketRequester.RequestSpec.data] providing a `data<Foo>(producer)`
|
||||
* variant leveraging Kotlin reified type parameters. This extension is not subject to type
|
||||
* erasure and retains actual generic type arguments.
|
||||
*
|
||||
* @author Sebastien Deleuze
|
||||
* @since 5.2
|
||||
*/
|
||||
@Suppress("EXTENSION_SHADOWED_BY_MEMBER")
|
||||
@FlowPreview
|
||||
fun <T : Any> RSocketRequester.RequestSpec.dataFlow(data: Flow<T>): RSocketRequester.ResponseSpec =
|
||||
data(data.asPublisher(), object : ParameterizedTypeReference<T>() {})
|
||||
fun <T : Any> RSocketRequester.RequestSpec.data(producer: Any): RSocketRequester.ResponseSpec =
|
||||
data(producer, object : ParameterizedTypeReference<T>() {})
|
||||
|
||||
/**
|
||||
* Coroutines variant of [RSocketRequester.ResponseSpec.send].
|
||||
@@ -92,3 +95,26 @@ suspend fun <T : Any> RSocketRequester.ResponseSpec.retrieveAndAwait(): T =
|
||||
@FlowPreview
|
||||
fun <T : Any> RSocketRequester.ResponseSpec.retrieveFlow(batchSize: Int = 1): Flow<T> =
|
||||
retrieveFlux(object : ParameterizedTypeReference<T>() {}).asFlow(batchSize)
|
||||
|
||||
/**
|
||||
* Extension for [RSocketRequester.ResponseSpec.retrieveMono] providing a `retrieveMono<Foo>()`
|
||||
* variant leveraging Kotlin reified type parameters. This extension is not subject to type
|
||||
* erasure and retains actual generic type arguments.
|
||||
*
|
||||
* @author Sebastien Deleuze
|
||||
* @since 5.2
|
||||
*/
|
||||
fun <T : Any> RSocketRequester.ResponseSpec.retrieveMono(): Mono<T> =
|
||||
retrieveMono(object : ParameterizedTypeReference<T>() {})
|
||||
|
||||
|
||||
/**
|
||||
* Extension for [RSocketRequester.ResponseSpec.retrieveFlux] providing a `retrieveFlux<Foo>()`
|
||||
* variant leveraging Kotlin reified type parameters. This extension is not subject to type
|
||||
* erasure and retains actual generic type arguments.
|
||||
*
|
||||
* @author Sebastien Deleuze
|
||||
* @since 5.2
|
||||
*/
|
||||
fun <T : Any> RSocketRequester.ResponseSpec.retrieveFlux(): Flux<T> =
|
||||
retrieveFlux(object : ParameterizedTypeReference<T>() {})
|
||||
|
||||
Reference in New Issue
Block a user