2020-10-04 10:57:39 +00:00
|
|
|
package dev.inmo.tgbotapi.extensions.utils
|
2020-08-10 05:53:54 +00:00
|
|
|
|
2021-01-09 12:25:11 +00:00
|
|
|
import dev.inmo.micro_utils.coroutines.safely
|
2020-08-10 05:53:54 +00:00
|
|
|
import kotlinx.coroutines.CoroutineScope
|
|
|
|
import kotlinx.coroutines.channels.BroadcastChannel
|
|
|
|
import kotlinx.coroutines.flow.*
|
|
|
|
|
2020-08-13 06:53:04 +00:00
|
|
|
/**
|
|
|
|
* Analog of [merge] function for [Flow]s. The difference is in the usage of [BroadcastChannel] in this case
|
|
|
|
*/
|
2020-08-10 05:53:54 +00:00
|
|
|
fun <T> aggregateFlows(
|
|
|
|
withScope: CoroutineScope,
|
|
|
|
vararg flows: Flow<T>,
|
2020-10-27 09:26:32 +00:00
|
|
|
internalBufferSize: Int = 64
|
2020-08-10 05:53:54 +00:00
|
|
|
): Flow<T> {
|
2020-10-27 09:11:57 +00:00
|
|
|
val sharedFlow = MutableSharedFlow<T>(extraBufferCapacity = internalBufferSize)
|
2020-08-10 05:53:54 +00:00
|
|
|
flows.forEach {
|
|
|
|
it.onEach {
|
2020-10-27 09:11:57 +00:00
|
|
|
safely { sharedFlow.emit(it) }
|
2020-08-10 05:53:54 +00:00
|
|
|
}.launchIn(withScope)
|
|
|
|
}
|
2020-10-27 09:51:48 +00:00
|
|
|
return sharedFlow
|
2020-08-10 05:53:54 +00:00
|
|
|
}
|
2020-08-13 08:55:21 +00:00
|
|
|
|
2022-05-21 19:27:52 +00:00
|
|
|
fun <T> Flow<Iterable<T>>.flatten(): Flow<T> = flow {
|
2020-08-13 08:55:21 +00:00
|
|
|
collect {
|
|
|
|
it.forEach {
|
|
|
|
emit(it)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2020-08-13 09:28:10 +00:00
|
|
|
|
2022-08-24 09:03:54 +00:00
|
|
|
fun <T, R> Flow<T>.flatMap(mapper: suspend (T) -> Iterable<R>): Flow<R> = flow {
|
2020-08-13 09:28:10 +00:00
|
|
|
collect {
|
|
|
|
mapper(it).forEach {
|
|
|
|
emit(it)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|