-
Notifications
You must be signed in to change notification settings - Fork 100
/
Turbine.kt
227 lines (191 loc) · 7.02 KB
/
Turbine.kt
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
/*
* Copyright (C) 2022 Square, Inc.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package app.cash.turbine
import kotlin.time.Duration
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.Job
import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
import kotlinx.coroutines.channels.ChannelResult
internal const val debug = false
/**
* A standalone [Turbine] suitable for usage in fakes or other external test components.
*/
public interface Turbine<T> : ReceiveTurbine<T> {
/**
* Returns the underlying [Channel]. The [Channel] will have a buffer size of [UNLIMITED].
*/
public override fun asChannel(): Channel<T>
/**
* Closes the underlying [Channel]. After all items have been consumed, this [Turbine] will yield
* [Event.Complete] if [cause] is null, and [Event.Error] otherwise.
*/
public fun close(cause: Throwable? = null)
/**
* Add an item to the underlying [Channel] without blocking.
*
* This method is equivalent to:
*
* ```
* if (!asChannel().trySend(item).isSuccess) error()
* ```
*/
public fun add(item: T)
/**
* Assert that the next event received was non-null and return it.
* This function will not suspend. On JVM and Android, it will attempt to throw if invoked in a suspending context.
*
* @throws AssertionError if the next event was completion or an error.
*/
public fun takeEvent(): Event<T>
/**
* Assert that the next event received was an item and return it.
* This function will not suspend. On JVM and Android, it will attempt to throw if invoked in a suspending context.
*
* @throws AssertionError if the next event was completion or an error, or no event.
*/
public fun takeItem(): T
/**
* Assert that the next event received is [Event.Complete].
* This function will not suspend. On JVM and Android, it will attempt to throw if invoked in a suspending context.
*
* @throws AssertionError if the next event was completion or an error.
*/
public fun takeComplete()
/**
* Assert that the next event received is [Event.Error], and return the error.
* This function will not suspend. On JVM and Android, it will attempt to throw if invoked in a suspending context.
*
* @throws AssertionError if the next event was completion or an error.
*/
public fun takeError(): Throwable
}
public operator fun <T> Turbine<T>.plusAssign(value: T) { add(value) }
/**
* Construct a standalone [Turbine].
*
* @param timeout If non-null, overrides the current Turbine timeout for this [Turbine]. See also:
* [withTurbineTimeout].
* @param name If non-null, name is added to any exceptions thrown to help identify which [Turbine] failed.
*/
@Suppress("FunctionName") // Interface constructor pattern.
public fun <T> Turbine(
timeout: Duration? = null,
name: String? = null,
): Turbine<T> = ChannelTurbine(timeout = timeout, name = name)
internal class ChannelTurbine<T>(
channel: Channel<T> = Channel(UNLIMITED),
private val job: Job? = null,
private val timeout: Duration?,
private val name: String?,
) : Turbine<T> {
private suspend fun <T> withTurbineTimeout(block: suspend () -> T): T {
return if (timeout != null) {
withTurbineTimeout(timeout) { block() }
} else {
block()
}
}
private val channel = object : Channel<T> by channel {
override fun tryReceive(): ChannelResult<T> {
val result = channel.tryReceive()
val event = result.toEvent()
if (event is Event.Error || event is Event.Complete) ignoreRemainingEvents = true
return result
}
override suspend fun receive(): T = try {
channel.receive()
} catch (e: Throwable) {
ignoreRemainingEvents = true
throw e
}
}
override fun asChannel(): Channel<T> = channel
override fun add(item: T) {
if (!channel.trySend(item).isSuccess) throw IllegalStateException("Added when closed")
}
@OptIn(ExperimentalCoroutinesApi::class)
override suspend fun cancel() {
if (!channel.isClosedForSend) ignoreTerminalEvents = true
channel.cancel()
job?.cancelAndJoin()
}
@OptIn(ExperimentalCoroutinesApi::class)
override fun close(cause: Throwable?) {
if (!channel.isClosedForSend) ignoreTerminalEvents = true
channel.close(cause)
job?.cancel()
}
override fun takeEvent(): Event<T> = channel.takeEvent(name = name)
override fun takeItem(): T = channel.takeItem(name = name)
override fun takeComplete() = channel.takeComplete(name = name)
override fun takeError(): Throwable = channel.takeError(name = name)
private var ignoreTerminalEvents = false
private var ignoreRemainingEvents = false
override suspend fun cancelAndIgnoreRemainingEvents() {
cancel()
ignoreRemainingEvents = true
}
override suspend fun cancelAndConsumeRemainingEvents(): List<Event<T>> {
val events = buildList {
while (true) {
val event = channel.takeEventUnsafe() ?: break
add(event)
if (event is Event.Error || event is Event.Complete) break
}
}
ignoreRemainingEvents = true
cancel()
return events
}
override fun expectNoEvents() {
channel.expectNoEvents(name = name)
}
override fun expectMostRecentItem(): T = channel.expectMostRecentItem(name = name)
override suspend fun awaitEvent(): Event<T> = withTurbineTimeout { channel.awaitEvent(name = name) }
override suspend fun awaitItem(): T = withTurbineTimeout { channel.awaitItem(name = name) }
override suspend fun skipItems(count: Int) = withTurbineTimeout { channel.skipItems(count, name) }
override suspend fun awaitComplete() = withTurbineTimeout { channel.awaitComplete(name = name) }
override suspend fun awaitError(): Throwable = withTurbineTimeout { channel.awaitError(name = name) }
override fun ensureAllEventsConsumed() {
if (ignoreRemainingEvents) return
val unconsumed = mutableListOf<Event<T>>()
var cause: Throwable? = null
while (true) {
val event = channel.takeEventUnsafe() ?: break
if (!(ignoreTerminalEvents && event.isTerminal)) unconsumed += event
if (event is Event.Error) {
cause = event.throwable
break
} else if (event is Event.Complete) {
break
}
}
if (unconsumed.isNotEmpty()) {
throw TurbineAssertionError(
buildString {
append("Unconsumed events found".qualifiedBy(name))
append(":")
for (event in unconsumed) {
append("\n - $event")
}
},
cause
)
}
}
}