Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 1 | /* |
Roman Elizarov | db0ef0c | 2019-07-03 15:02:44 +0300 | [diff] [blame] | 2 | * Copyright 2016-2019 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license. |
Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 3 | */ |
| 4 | |
| 5 | // This file was automatically generated from coroutines-guide-reactive.md by Knit tool. Do not edit. |
Roman Elizarov | 0950dfa | 2018-07-13 10:33:25 +0300 | [diff] [blame] | 6 | package kotlinx.coroutines.rx2.guide.operators03 |
Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 7 | |
Roman Elizarov | 0950dfa | 2018-07-13 10:33:25 +0300 | [diff] [blame] | 8 | import kotlinx.coroutines.channels.* |
| 9 | import kotlinx.coroutines.* |
| 10 | import kotlinx.coroutines.reactive.* |
| 11 | import kotlinx.coroutines.selects.* |
Roman Elizarov | 9fe5f46 | 2018-02-21 19:05:52 +0300 | [diff] [blame] | 12 | import org.reactivestreams.* |
Roman Elizarov | 0950dfa | 2018-07-13 10:33:25 +0300 | [diff] [blame] | 13 | import kotlin.coroutines.* |
Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 14 | |
Vsevolod Tolstopyatov | d100a3f | 2019-07-16 16:14:35 +0300 | [diff] [blame] | 15 | fun <T, U> Publisher<T>.takeUntil(context: CoroutineContext, other: Publisher<U>) = publish<T>(context) { |
Vsevolod Tolstopyatov | 313978c | 2018-06-01 15:30:34 +0300 | [diff] [blame] | 16 | this@takeUntil.openSubscription().consume { // explicitly open channel to Publisher<T> |
| 17 | val current = this |
| 18 | other.openSubscription().consume { // explicitly open channel to Publisher<U> |
| 19 | val other = this |
| 20 | whileSelect { |
Roman Elizarov | c961fb6 | 2019-04-24 13:01:49 +0300 | [diff] [blame] | 21 | other.onReceive { false } // bail out on any received element from `other` |
| 22 | current.onReceive { send(it); true } // resend element from this channel and continue |
Vsevolod Tolstopyatov | 313978c | 2018-06-01 15:30:34 +0300 | [diff] [blame] | 23 | } |
| 24 | } |
Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 25 | } |
| 26 | } |
| 27 | |
Vsevolod Tolstopyatov | bbaf99d | 2018-09-11 15:55:56 +0300 | [diff] [blame] | 28 | fun CoroutineScope.rangeWithInterval(time: Long, start: Int, count: Int) = publish<Int> { |
Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 29 | for (x in start until start + count) { |
| 30 | delay(time) // wait before sending each number |
| 31 | send(x) |
| 32 | } |
| 33 | } |
| 34 | |
Prendota | 65e6c8c | 2018-10-17 11:51:08 +0300 | [diff] [blame] | 35 | fun main() = runBlocking<Unit> { |
Vsevolod Tolstopyatov | bbaf99d | 2018-09-11 15:55:56 +0300 | [diff] [blame] | 36 | val slowNums = rangeWithInterval(200, 1, 10) // numbers with 200ms interval |
| 37 | val stop = rangeWithInterval(500, 1, 10) // the first one after 500ms |
Vsevolod Tolstopyatov | d100a3f | 2019-07-16 16:14:35 +0300 | [diff] [blame] | 38 | slowNums.takeUntil(Dispatchers.Unconfined, stop).collect { println(it) } // let's test it |
Roman Elizarov | 8a4a8e1 | 2017-03-09 19:52:58 +0300 | [diff] [blame] | 39 | } |