Improve WebFlux suspending handler method support
This commit leverages Flux instead of Flow to support suspending handler methods returning Flow in order to avoid multiple invocations of the suspending function on every collect(). See gh-22820
This commit is contained in:
parent
b33d2f4634
commit
cd5dc84832
|
@ -23,9 +23,8 @@ import kotlinx.coroutines.FlowPreview
|
|||
import kotlinx.coroutines.GlobalScope
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.collect
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.reactive.awaitFirstOrNull
|
||||
import kotlinx.coroutines.reactive.flow.asPublisher
|
||||
|
||||
import kotlinx.coroutines.reactor.mono
|
||||
import reactor.core.publisher.Mono
|
||||
|
@ -33,8 +32,6 @@ import reactor.core.publisher.onErrorMap
|
|||
import java.lang.reflect.InvocationTargetException
|
||||
import java.lang.reflect.Method
|
||||
import kotlin.reflect.full.callSuspend
|
||||
import kotlin.reflect.full.isSubtypeOf
|
||||
import kotlin.reflect.full.starProjectedType
|
||||
import kotlin.reflect.jvm.kotlinFunction
|
||||
|
||||
/**
|
||||
|
@ -56,7 +53,8 @@ internal fun <T: Any> monoToDeferred(source: Mono<T>) =
|
|||
GlobalScope.async(Dispatchers.Unconfined) { source.awaitFirstOrNull() }
|
||||
|
||||
/**
|
||||
* Invoke an handler method converting suspending method to [Mono] or [Flow] if necessary.
|
||||
* Invoke an handler method converting suspending method to [Mono] or
|
||||
* [reactor.core.publisher.Flux] if necessary.
|
||||
*
|
||||
* @author Sebastien Deleuze
|
||||
* @since 5.2
|
||||
|
@ -66,18 +64,15 @@ internal fun <T: Any> monoToDeferred(source: Mono<T>) =
|
|||
internal fun invokeHandlerMethod(method: Method, bean: Any, vararg args: Any?): Any? {
|
||||
val function = method.kotlinFunction!!
|
||||
return if (function.isSuspend) {
|
||||
if (function.returnType.isSubtypeOf(Flow::class.starProjectedType)) {
|
||||
flow {
|
||||
(function.callSuspend(bean, *args.sliceArray(0..(args.size-2))) as Flow<*>).collect {
|
||||
emit(it)
|
||||
}
|
||||
}
|
||||
val mono = GlobalScope.mono(Dispatchers.Unconfined) {
|
||||
function.callSuspend(bean, *args.sliceArray(0..(args.size-2)))
|
||||
.let { if (it == Unit) null else it }
|
||||
}.onErrorMap(InvocationTargetException::class) { it.targetException }
|
||||
if (function.returnType.classifier == Flow::class) {
|
||||
mono.flatMapMany { (it as Flow<Any>).asPublisher() }
|
||||
}
|
||||
else {
|
||||
GlobalScope.mono(Dispatchers.Unconfined) {
|
||||
function.callSuspend(bean, *args.sliceArray(0..(args.size-2)))
|
||||
.let { if (it == Unit) null else it}
|
||||
}.onErrorMap(InvocationTargetException::class) { it.targetException }
|
||||
mono
|
||||
}
|
||||
}
|
||||
else {
|
||||
|
|
Loading…
Reference in New Issue