Updated on 2026-08-14

This commit is contained in:
Tangem 2025-03-12 20:44:27 +04:00
parent 80d661ab55
commit 918ab4613f
6 changed files with 274 additions and 0 deletions

View file

@ -10,4 +10,9 @@ dependencies {
api(deps.arrow.fx)
implementation(deps.kotlin.serialization)
testImplementation(deps.test.coroutine)
testImplementation(deps.test.junit)
testImplementation(deps.test.mockk)
testImplementation(deps.test.truth)
}

View file

@ -0,0 +1,59 @@
package com.tangem.domain.core.flow
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.update
private typealias FlowsStore<T> = MutableStateFlow<Map<String, Flow<T>>>
/**
* [Flow] supplier with caching mechanism
*
* @param Producer producer of flow [Flow]
* @param Params type of data that required to get [Flow]
* @param Data data type of [Flow]
* @property flowsStore store of flows
*/
abstract class FlowCachingSupplier<Producer : FlowProducer<Data>, Params : Any, Data : Any>(
private val flowsStore: FlowsStore<Data> = MutableStateFlow(value = emptyMap()),
) : FlowSupplier<Params, Data> {
/** Factory of [FlowProducer] */
abstract val factory: FlowProducer.Factory<Params, Producer>
/** Key creator */
abstract val keyCreator: (Params) -> String
/**
* Supply [Flow] by [params].
*/
override operator fun invoke(params: Params): Flow<Data> {
val key = keyCreator(params)
val saved = flowsStore.value[key]
if (saved != null) return saved
val flowProducer = factory.create(params = params)
return runCatching(flowProducer::produceWithFallback)
.onSuccess { flow ->
flowsStore.update {
it.toMutableMap().apply {
put(key = key, value = flow)
}
}
}
.getOrThrow()
.catch { cause ->
flowsStore.update {
it.toMutableMap().apply {
remove(key)
}
}
throw cause
}
}
}

View file

@ -0,0 +1,15 @@
package com.tangem.domain.core.flow
import arrow.core.Either
/**
* Flow fetcher
*
* @param Params data that required to fetch flow
*
[REDACTED_AUTHOR]
*/
interface FlowFetcher<Params : Any> {
suspend operator fun invoke(params: Params): Either<Throwable, Unit>
}

View file

@ -0,0 +1,39 @@
package com.tangem.domain.core.flow
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.retryWhen
/**
* [Flow] producer
*
* @param Data data type of [Flow]
*
[REDACTED_AUTHOR]
*/
interface FlowProducer<Data : Any> {
/** Fallback value if [Flow] throws exception */
val fallback: Data
/** Produce [Flow] */
fun produce(): Flow<Data>
/** Produce [Flow] with retry mechanism */
fun produceWithFallback(): Flow<Data> {
return produce().retryWhen { _, _ ->
emit(value = fallback)
delay(timeMillis = 2000)
true
}
}
/** Factory for creating [Producer]. It helps to provide [Params] by constructor */
interface Factory<Params : Any, Producer : FlowProducer<*>> {
/** Create [Producer] using [params] */
fun create(params: Params): Producer
}
}

View file

@ -0,0 +1,17 @@
package com.tangem.domain.core.flow
import kotlinx.coroutines.flow.Flow
/**
* [Flow] supplier
*
* @param Params type of data that required to get [Flow]
* @param Data data type of [Flow]
*
[REDACTED_AUTHOR]
*/
interface FlowSupplier<Params : Any, Data : Any> {
/** Supply [Flow] by [params] */
operator fun invoke(params: Params): Flow<Data>
}

View file

@ -0,0 +1,139 @@
package com.tangem.domain.core.flow
import com.google.common.truth.Truth
import io.mockk.every
import io.mockk.mockk
import io.mockk.verify
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.test.runTest
import org.junit.Test
/**
[REDACTED_AUTHOR]
*/
internal class FlowCachingSupplierTest {
private val factory = mockk<MockFlowProducer.Factory>()
private val errorFactory = mockk<MockErrorFlowProducer.Factory>()
@Test
fun `flow was created before`() = runTest {
val flowsStore = MutableStateFlow<Map<String, Flow<String>>>(value = emptyMap())
val supplier = MockFlowCachingSupplier(factory = factory, flowsStore = flowsStore)
val params = 1
every { factory.create(params = params) } returns MockFlowProducer(params = params)
val actual = supplier(params = params)
verify { factory.create(params = params) }
Truth.assertThat(actual.first()).isEqualTo("test_1")
Truth.assertThat(flowsStore.value.size).isEqualTo(1)
Truth.assertThat(flowsStore.value.containsKey("mock_1")).isTrue()
Truth.assertThat(flowsStore.value["mock_1"]!!.first()).isEqualTo("test_1")
}
@Test
fun `flow wasn't created before`() = runTest {
val storedFlow = flowOf("test_1")
val storedMap = mapOf("mock_1" to storedFlow)
val flowsStore = MutableStateFlow(value = storedMap)
val supplier = MockFlowCachingSupplier(factory = factory, flowsStore = flowsStore)
val params = 1
val actual = supplier(params = params)
verify(inverse = true) { factory.create(params = params) }
Truth.assertThat(actual).isEqualTo(storedFlow)
Truth.assertThat(flowsStore.value).isEqualTo(storedMap)
}
@Test
fun `flow wasn't created before and store isn't empty`() = runTest {
val flowsStore = MutableStateFlow(
value = mapOf("mock_1" to flowOf("test_1")),
)
val supplier = MockFlowCachingSupplier(factory = factory, flowsStore = flowsStore)
val params = 2
every { factory.create(params = params) } returns MockFlowProducer(params = params)
val actual = supplier(params = params)
verify { factory.create(params = params) }
Truth.assertThat(actual.first()).isEqualTo("test_2")
Truth.assertThat(flowsStore.value.size).isEqualTo(2)
Truth.assertThat(flowsStore.value.keys).isEqualTo(setOf("mock_1", "mock_2"))
Truth.assertThat(flowsStore.value["mock_1"]!!.first()).isEqualTo("test_1")
Truth.assertThat(flowsStore.value["mock_2"]!!.first()).isEqualTo("test_2")
}
@Test
fun `flow throws exception`() = runTest {
val flowsStore = MutableStateFlow<Map<String, Flow<String>>>(value = emptyMap())
val supplier = MockErrorFlowCachingSupplier(factory = errorFactory, flowsStore = flowsStore)
val params = 1
every { errorFactory.create(params = params) } returns MockErrorFlowProducer()
val actual = supplier(params = params)
verify { errorFactory.create(params = params) }
Truth.assertThat(actual.first()).isEqualTo("fallback")
Truth.assertThat(flowsStore.value.size).isEqualTo(1)
}
private class MockFlowCachingSupplier(
override val factory: FlowProducer.Factory<Int, MockFlowProducer>,
flowsStore: MutableStateFlow<Map<String, Flow<String>>>,
) : FlowCachingSupplier<MockFlowProducer, Int, String>(flowsStore = flowsStore) {
override val keyCreator: (Int) -> String = { "mock_$it" }
}
private class MockFlowProducer(private val params: Int) : FlowProducer<String> {
override val fallback: String
get() = "fallback"
override fun produce(): Flow<String> = flowOf("test_$params")
class Factory : FlowProducer.Factory<Int, MockFlowProducer> {
override fun create(params: Int): MockFlowProducer = MockFlowProducer(params = params)
}
}
private class MockErrorFlowCachingSupplier(
override val factory: FlowProducer.Factory<Int, MockErrorFlowProducer>,
flowsStore: MutableStateFlow<Map<String, Flow<String>>>,
) : FlowCachingSupplier<MockErrorFlowProducer, Int, String>(flowsStore = flowsStore) {
override val keyCreator: (Int) -> String = { "mock_$it" }
}
private class MockErrorFlowProducer : FlowProducer<String> {
override val fallback: String
get() = "fallback"
override fun produce(): Flow<String> = flow {
throw IllegalStateException()
}
class Factory : FlowProducer.Factory<Int, MockErrorFlowProducer> {
override fun create(params: Int): MockErrorFlowProducer = MockErrorFlowProducer()
}
}
}