package kotlinx.coroutines.flow import kotlinx.coroutines.testing.* import kotlinx.coroutines.* import kotlinx.coroutines.channels.* import kotlin.test.* class ScanTest : TestBase() { @Test fun testScan() = runTest { val flow = flowOf(1, 2, 3, 4, 5) val result = flow.runningReduce { acc, v -> acc + v }.toList() assertEquals(listOf(1, 3, 6, 10, 15), result) } @Test fun testScanWithInitial() = runTest { val flow = flowOf(1, 2, 3) val result = flow.scan(emptyList()) { acc, value -> acc + value }.toList() assertEquals(listOf(emptyList(), listOf(1), listOf(1, 2), listOf(1, 2, 3)), result) } @Test fun testFoldWithInitial() = runTest { val flow = flowOf(1, 2, 3) val result = flow.runningFold(emptyList()) { acc, value -> acc + value }.toList() assertEquals(listOf(emptyList(), listOf(1), listOf(1, 2), listOf(1, 2, 3)), result) } @Test fun testNulls() = runTest { val flow = flowOf(null, 2, null, null, null, 5) val result = flow.runningReduce { acc, v -> if (v == null) acc else (if (acc == null) v else acc + v) }.toList() assertEquals(listOf(null, 2, 2, 2, 2, 7), result) } @Test fun testEmptyFlow() = runTest { val result = emptyFlow().runningReduce { _, _ -> 1 }.toList() assertTrue(result.isEmpty()) } @Test fun testErrorCancelsUpstream() = runTest { expect(1) val latch = Channel() val flow = flow { coroutineScope { launch { latch.send(Unit) hang { expect(3) } } emit(1) emit(2) } }.runningReduce { _, value -> expect(value) // 2 latch.receive() throw TestException() }.catch { /* ignore */ } assertEquals(1, flow.single()) finish(4) } private operator fun Collection.plus(element: T): List { val result = ArrayList(size + 1) result.addAll(this) result.add(element) return result } }