Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 1 | /* |
Roman Elizarov | 1f74a2d | 2018-06-29 19:19:45 +0300 | [diff] [blame] | 2 | * Copyright 2016-2018 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license. |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 3 | */ |
| 4 | |
Roman Elizarov | 0950dfa | 2018-07-13 10:33:25 +0300 | [diff] [blame] | 5 | package kotlinx.coroutines.rx2 |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 6 | |
Vsevolod Tolstopyatov | bbaf99d | 2018-09-11 15:55:56 +0300 | [diff] [blame] | 7 | import io.reactivex.* |
Roman Elizarov | 0950dfa | 2018-07-13 10:33:25 +0300 | [diff] [blame] | 8 | import kotlinx.coroutines.* |
| 9 | import kotlinx.coroutines.channels.* |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 10 | import kotlinx.coroutines.flow.* |
SokolovaMaria | 1dcfd97 | 2019-08-09 17:35:14 +0300 | [diff] [blame] | 11 | import kotlinx.coroutines.reactive.* |
Roman Elizarov | 0950dfa | 2018-07-13 10:33:25 +0300 | [diff] [blame] | 12 | import kotlin.coroutines.* |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 13 | |
| 14 | /** |
| 15 | * Converts this job to the hot reactive completable that signals |
| 16 | * with [onCompleted][CompletableSubscriber.onCompleted] when the corresponding job completes. |
| 17 | * |
| 18 | * Every subscriber gets the signal at the same time. |
| 19 | * Unsubscribing from the resulting completable **does not** affect the original job in any way. |
| 20 | * |
Roman Elizarov | 27b8f45 | 2018-09-20 21:23:41 +0300 | [diff] [blame] | 21 | * **Note: This is an experimental api.** Conversion of coroutines primitives to reactive entities may change |
| 22 | * in the future to account for the concept of structured concurrency. |
| 23 | * |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 24 | * @param context -- the coroutine context from which the resulting completable is going to be signalled |
| 25 | */ |
Roman Elizarov | 27b8f45 | 2018-09-20 21:23:41 +0300 | [diff] [blame] | 26 | @ExperimentalCoroutinesApi |
Vsevolod Tolstopyatov | d100a3f | 2019-07-16 16:14:35 +0300 | [diff] [blame] | 27 | public fun Job.asCompletable(context: CoroutineContext): Completable = rxCompletable(context) { |
Roman Elizarov | 3c3aed7 | 2017-03-09 12:31:59 +0300 | [diff] [blame] | 28 | this@asCompletable.join() |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 29 | } |
| 30 | |
| 31 | /** |
Konrad Kamiński | d6bb148 | 2017-04-07 09:26:40 +0200 | [diff] [blame] | 32 | * Converts this deferred value to the hot reactive maybe that signals |
| 33 | * [onComplete][MaybeEmitter.onComplete], [onSuccess][MaybeEmitter.onSuccess] or [onError][MaybeEmitter.onError]. |
| 34 | * |
| 35 | * Every subscriber gets the same completion value. |
| 36 | * Unsubscribing from the resulting maybe **does not** affect the original deferred value in any way. |
| 37 | * |
Roman Elizarov | 27b8f45 | 2018-09-20 21:23:41 +0300 | [diff] [blame] | 38 | * **Note: This is an experimental api.** Conversion of coroutines primitives to reactive entities may change |
| 39 | * in the future to account for the concept of structured concurrency. |
| 40 | * |
Konrad Kamiński | d6bb148 | 2017-04-07 09:26:40 +0200 | [diff] [blame] | 41 | * @param context -- the coroutine context from which the resulting maybe is going to be signalled |
| 42 | */ |
Roman Elizarov | 27b8f45 | 2018-09-20 21:23:41 +0300 | [diff] [blame] | 43 | @ExperimentalCoroutinesApi |
Vsevolod Tolstopyatov | d100a3f | 2019-07-16 16:14:35 +0300 | [diff] [blame] | 44 | public fun <T> Deferred<T?>.asMaybe(context: CoroutineContext): Maybe<T> = rxMaybe(context) { |
Konrad Kamiński | d6bb148 | 2017-04-07 09:26:40 +0200 | [diff] [blame] | 45 | this@asMaybe.await() |
| 46 | } |
| 47 | |
| 48 | /** |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 49 | * Converts this deferred value to the hot reactive single that signals either |
| 50 | * [onSuccess][SingleSubscriber.onSuccess] or [onError][SingleSubscriber.onError]. |
| 51 | * |
| 52 | * Every subscriber gets the same completion value. |
| 53 | * Unsubscribing from the resulting single **does not** affect the original deferred value in any way. |
| 54 | * |
Roman Elizarov | 27b8f45 | 2018-09-20 21:23:41 +0300 | [diff] [blame] | 55 | * **Note: This is an experimental api.** Conversion of coroutines primitives to reactive entities may change |
| 56 | * in the future to account for the concept of structured concurrency. |
| 57 | * |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 58 | * @param context -- the coroutine context from which the resulting single is going to be signalled |
| 59 | */ |
Roman Elizarov | 27b8f45 | 2018-09-20 21:23:41 +0300 | [diff] [blame] | 60 | @ExperimentalCoroutinesApi |
Vsevolod Tolstopyatov | d100a3f | 2019-07-16 16:14:35 +0300 | [diff] [blame] | 61 | public fun <T : Any> Deferred<T>.asSingle(context: CoroutineContext): Single<T> = rxSingle(context) { |
Roman Elizarov | 3c3aed7 | 2017-03-09 12:31:59 +0300 | [diff] [blame] | 62 | this@asSingle.await() |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 63 | } |
| 64 | |
| 65 | /** |
| 66 | * Converts a stream of elements received from the channel to the hot reactive observable. |
| 67 | * |
| 68 | * Every subscriber receives values from this channel in **fan-out** fashion. If the are multiple subscribers, |
| 69 | * they'll receive values in round-robin way. |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 70 | */ |
Vsevolod Tolstopyatov | 5378b80 | 2019-09-19 17:39:08 +0300 | [diff] [blame] | 71 | @Deprecated( |
| 72 | message = "Deprecated in the favour of Flow", |
| 73 | level = DeprecationLevel.WARNING, replaceWith = ReplaceWith("this.consumeAsFlow().asObservable()") |
| 74 | ) |
Vsevolod Tolstopyatov | d100a3f | 2019-07-16 16:14:35 +0300 | [diff] [blame] | 75 | public fun <T : Any> ReceiveChannel<T>.asObservable(context: CoroutineContext): Observable<T> = rxObservable(context) { |
Roman Elizarov | 3c3aed7 | 2017-03-09 12:31:59 +0300 | [diff] [blame] | 76 | for (t in this@asObservable) |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 77 | send(t) |
| 78 | } |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 79 | |
| 80 | /** |
| 81 | * Converts the given flow to a cold observable. |
SokolovaMaria | 1dcfd97 | 2019-08-09 17:35:14 +0300 | [diff] [blame] | 82 | * The original flow is cancelled when the observable subscriber is disposed. |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 83 | */ |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 84 | @JvmName("from") |
Vsevolod Tolstopyatov | 15c7d0f | 2019-06-06 11:47:19 +0300 | [diff] [blame] | 85 | @ExperimentalCoroutinesApi |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 86 | public fun <T: Any> Flow<T>.asObservable() : Observable<T> = Observable.create { emitter -> |
| 87 | /* |
| 88 | * ATOMIC is used here to provide stable behaviour of subscribe+dispose pair even if |
| 89 | * asObservable is already invoked from unconfined |
| 90 | */ |
| 91 | val job = GlobalScope.launch(Dispatchers.Unconfined, start = CoroutineStart.ATOMIC) { |
| 92 | try { |
| 93 | collect { value -> emitter.onNext(value) } |
| 94 | emitter.onComplete() |
| 95 | } catch (e: Throwable) { |
| 96 | // 'create' provides safe emitter, so we can unconditionally call on* here if exception occurs in `onComplete` |
Vsevolod Tolstopyatov | a930b0c | 2019-12-05 19:33:44 +0300 | [diff] [blame^] | 97 | if (e !is CancellationException) { |
| 98 | if (!emitter.tryOnError(e)) { |
| 99 | handleUndeliverableException(e, coroutineContext) |
| 100 | } |
| 101 | } else { |
| 102 | emitter.onComplete() |
| 103 | } |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 104 | } |
| 105 | } |
| 106 | emitter.setCancellable(RxCancellable(job)) |
| 107 | } |
| 108 | |
| 109 | /** |
SokolovaMaria | 1dcfd97 | 2019-08-09 17:35:14 +0300 | [diff] [blame] | 110 | * Converts the given flow to a cold flowable. |
| 111 | * The original flow is cancelled when the flowable subscriber is disposed. |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 112 | */ |
Vsevolod Tolstopyatov | 170690f | 2019-04-09 12:33:57 +0300 | [diff] [blame] | 113 | @JvmName("from") |
Vsevolod Tolstopyatov | 15c7d0f | 2019-06-06 11:47:19 +0300 | [diff] [blame] | 114 | @ExperimentalCoroutinesApi |
Vsevolod Tolstopyatov | 61c64cc | 2019-04-12 16:05:58 +0300 | [diff] [blame] | 115 | public fun <T: Any> Flow<T>.asFlowable(): Flowable<T> = Flowable.fromPublisher(asPublisher()) |