SokolovaMaria | 1dcfd97 | 2019-08-09 17:35:14 +0300 | [diff] [blame^] | 1 | @file:JvmName("FlowKt") |
| 2 | |
| 3 | package kotlinx.coroutines.reactor |
| 4 | |
| 5 | import kotlinx.coroutines.* |
| 6 | import kotlinx.coroutines.flow.Flow |
| 7 | import kotlinx.coroutines.flow.flowOn |
| 8 | import kotlinx.coroutines.reactive.FlowSubscription |
| 9 | import reactor.core.CoreSubscriber |
| 10 | import reactor.core.publisher.Flux |
| 11 | |
| 12 | /** |
| 13 | * Converts the given flow to a cold flux. |
| 14 | * The original flow is cancelled when the flux subscriber is disposed. |
| 15 | */ |
| 16 | @ExperimentalCoroutinesApi |
| 17 | public fun <T: Any> Flow<T>.asFlux(): Flux<T> = FlowAsFlux(this) |
| 18 | |
| 19 | private class FlowAsFlux<T : Any>(private val flow: Flow<T>) : Flux<T>() { |
| 20 | override fun subscribe(subscriber: CoreSubscriber<in T>?) { |
| 21 | if (subscriber == null) throw NullPointerException() |
| 22 | val hasContext = subscriber.currentContext().isEmpty |
| 23 | val source = if (hasContext) flow.flowOn(subscriber.currentContext().asCoroutineContext()) else flow |
| 24 | subscriber.onSubscribe(FlowSubscription(source, subscriber)) |
| 25 | } |
| 26 | } |