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 | |
| 5 | package kotlinx.coroutines.experimental.reactive |
| 6 | |
| 7 | import kotlinx.coroutines.experimental.CommonPool |
| 8 | import kotlinx.coroutines.experimental.TestBase |
| 9 | import kotlinx.coroutines.experimental.launch |
| 10 | import kotlinx.coroutines.experimental.runBlocking |
| 11 | import org.hamcrest.core.IsEqual |
| 12 | import org.junit.Assert.assertThat |
| 13 | import org.junit.Assert.assertTrue |
| 14 | import org.junit.Test |
| 15 | |
| 16 | /** |
| 17 | * Test emitting multiple values with [publish]. |
| 18 | */ |
| 19 | class PublisherMultiTest : TestBase() { |
| 20 | @Test |
| 21 | fun testConcurrentStress() = runBlocking<Unit> { |
| 22 | val n = 10_000 * stressTestMultiplier |
| 23 | val observable = publish<Int>(CommonPool) { |
| 24 | // concurrent emitters (many coroutines) |
| 25 | val jobs = List(n) { |
| 26 | // launch |
| 27 | launch(CommonPool) { |
| 28 | send(it) |
| 29 | } |
| 30 | } |
| 31 | jobs.forEach { it.join() } |
| 32 | } |
| 33 | val resultSet = mutableSetOf<Int>() |
Roman Elizarov | 86349be | 2017-03-17 16:47:37 +0300 | [diff] [blame] | 34 | observable.consumeEach { |
| 35 | assertTrue(resultSet.add(it)) |
Roman Elizarov | 331750b | 2017-02-15 17:59:17 +0300 | [diff] [blame] | 36 | } |
| 37 | assertThat(resultSet.size, IsEqual(n)) |
| 38 | } |
| 39 | } |