diff --git a/data/networks/src/main/java/com/tangem/data/networks/store/DefaultNetworksStatusesStore.kt b/data/networks/src/main/java/com/tangem/data/networks/store/DefaultNetworksStatusesStore.kt index 1131c43dfc..7961e06c77 100644 --- a/data/networks/src/main/java/com/tangem/data/networks/store/DefaultNetworksStatusesStore.kt +++ b/data/networks/src/main/java/com/tangem/data/networks/store/DefaultNetworksStatusesStore.kt @@ -13,13 +13,8 @@ import com.tangem.domain.models.network.NetworkStatus import com.tangem.domain.models.wallet.UserWalletId import com.tangem.utils.coroutines.CoroutineDispatcherProvider import com.tangem.utils.extensions.addOrReplace -import kotlinx.coroutines.CoroutineScope -import kotlinx.coroutines.SupervisorJob -import kotlinx.coroutines.coroutineScope -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.firstOrNull -import kotlinx.coroutines.flow.mapNotNull -import kotlinx.coroutines.launch +import kotlinx.coroutines.* +import kotlinx.coroutines.flow.* import timber.log.Timber import java.io.File @@ -67,6 +62,8 @@ internal class DefaultNetworksStatusesStore( override fun get(userWalletId: UserWalletId): Flow> { return runtimeStore.get().mapNotNull { it[userWalletId.stringValue] } + .adaptiveThrottle() + .conflate() } override suspend fun getSyncOrNull(userWalletId: UserWalletId, network: Network): SimpleNetworkStatus? { @@ -180,4 +177,61 @@ internal class DefaultNetworksStatusesStore( } } } +} + +@Suppress("MagicNumber") +internal fun Flow>.adaptiveThrottle(): Flow> = channelFlow { + var accumulator: Set? = null + var lastEmitTime = 0L + + // params that control maximum emissions that can be throttled + var densityLevel = 0 + val maxDensity = 10 + + // params that control maximum delay and growth of delay between emissions + var lastDelay = 0L + val maxDelay = 1500L + val growthFactor = 250L + + fun resetThrottling() { + lastDelay = 0L + densityLevel = 0 + } + + this@adaptiveThrottle.collectLatest { newSet -> + val previousSet: Collection? = accumulator + accumulator = newSet + + when { + // first value, just emit + previousSet == null -> resetThrottling() + // changed size, just emit + previousSet.size != newSet.size -> resetThrottling() + + // apply adaptive throttling + else -> { + val networksCount = newSet.size + // more networks - more throttling + val cooldownThreshold = when { + networksCount in 10..25 -> 300L + networksCount > 25 -> 500L + // 0..9 networks + else -> 100L + } + + val now = System.currentTimeMillis() + val timeSinceLastEmit = now - lastEmitTime + if (timeSinceLastEmit < cooldownThreshold && densityLevel < maxDensity) { + lastDelay = (lastDelay + growthFactor).coerceAtMost(maximumValue = maxDelay) + densityLevel += 1 + delay(lastDelay) + } else { + resetThrottling() + } + } + } + + lastEmitTime = System.currentTimeMillis() + channel.send(newSet) + } } \ No newline at end of file diff --git a/data/networks/src/test/java/com/tangem/data/networks/store/StoreAdaptiveThrottleTest.kt b/data/networks/src/test/java/com/tangem/data/networks/store/StoreAdaptiveThrottleTest.kt new file mode 100644 index 0000000000..719df24b15 --- /dev/null +++ b/data/networks/src/test/java/com/tangem/data/networks/store/StoreAdaptiveThrottleTest.kt @@ -0,0 +1,109 @@ +package com.tangem.data.networks.store + +import app.cash.turbine.test +import com.google.common.truth.Truth.assertThat +import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.flowOf +import kotlinx.coroutines.launch +import kotlinx.coroutines.test.advanceTimeBy +import kotlinx.coroutines.test.runTest +import org.junit.jupiter.api.Test + +internal class StoreAdaptiveThrottleTest { + + @Test + fun `first value is emitted immediately`() = runTest { + val flow = flowOf(setOf(1, 2, 3)).adaptiveThrottle() + + flow.test { + val item = awaitItem() + assertThat(item).isEqualTo(setOf(1, 2, 3)) + awaitComplete() + } + } + + @Test + fun `size change bypasses throttling`() = runTest { + val upstream = MutableSharedFlow>() + + upstream.adaptiveThrottle().test { + + upstream.emit(setOf(1, 2)) + assertThat(awaitItem()).isEqualTo(setOf(1, 2)) + + upstream.emit(setOf(1, 2, 3)) + assertThat(awaitItem()).isEqualTo(setOf(1, 2, 3)) + + upstream.emit(setOf(1)) + assertThat(awaitItem()).isEqualTo(setOf(1)) + } + } + + @Test + fun `same size events trigger throttling delay`() = runTest { + val upstream = MutableSharedFlow>() + + upstream.adaptiveThrottle().test { + + upstream.emit(setOf(1, 2)) + awaitItem() + + upstream.emit(setOf(3, 4)) + + // delay should happen + expectNoEvents() + advanceTimeBy(250) + + val item = awaitItem() + assertThat(item).isEqualTo(setOf(3, 4)) + } + } + + @Test + fun `rapid events result in only latest emission due to collectLatest`() = runTest { + val upstream = MutableSharedFlow>() + + upstream.adaptiveThrottle().test { + + upstream.emit(setOf(1, 2)) + awaitItem() + + launch { + upstream.emit(setOf(3, 4)) + upstream.emit(setOf(5, 6)) + upstream.emit(setOf(7, 8)) + } + + // delay should happen + expectNoEvents() + advanceTimeBy(250) + + val item = awaitItem() + assertThat(item).isEqualTo(setOf(7, 8)) + } + } + + @Test + fun `throttling resets when cooldown window passed`() = runTest { + val upstream = MutableSharedFlow>() + + upstream.adaptiveThrottle().test { + + upstream.emit(setOf(1, 2)) + awaitItem() + + upstream.emit(setOf(3, 4)) + expectNoEvents() + advanceTimeBy(250) + awaitItem() + + // wait long enough to reset throttling + advanceTimeBy(2000) + + upstream.emit(setOf(5, 6)) + + val item = awaitItem() + assertThat(item).isEqualTo(setOf(5, 6)) + } + } +} \ No newline at end of file diff --git a/domain/account/status/src/main/java/com/tangem/domain/account/status/producer/DefaultSingleAccountStatusListProducer.kt b/domain/account/status/src/main/java/com/tangem/domain/account/status/producer/DefaultSingleAccountStatusListProducer.kt index be5c3f927d..5250d458de 100644 --- a/domain/account/status/src/main/java/com/tangem/domain/account/status/producer/DefaultSingleAccountStatusListProducer.kt +++ b/domain/account/status/src/main/java/com/tangem/domain/account/status/producer/DefaultSingleAccountStatusListProducer.kt @@ -164,7 +164,7 @@ internal class DefaultSingleAccountStatusListProducer @AssistedInject constructo flattenCurrency: MutableSharedFlow>, ): Flow> { val walletId = userWallet.walletId - val networkStatusFlow: SharedFlow> = networkStatusFlow(walletId, flattenCurrency) + val networkStatusFlow: SharedFlow> = networkStatusFlow(walletId) .shareIn(this, started = SharingStarted.Eagerly, replay = 1) val stakingBalanceFlow: SharedFlow>> = stakingFlow(userWallet) .shareIn(this, started = SharingStarted.Eagerly, replay = 1) @@ -213,27 +213,10 @@ internal class DefaultSingleAccountStatusListProducer @AssistedInject constructo } } - private fun networkStatusFlow( - walletId: UserWalletId, - flattenCurrency: MutableSharedFlow>, - ): Flow> = channelFlow { - val currencyCount = flattenCurrency - .map { map -> map.size } - .stateIn(this, SharingStarted.Eagerly, 0) - + private fun networkStatusFlow(walletId: UserWalletId): Flow> = networkStatusSupplier(MultiNetworkStatusProducer.Params(walletId)) - // todo accounts high frequency, investigate better debounce - .debounce { - val count = currencyCount.value - @Suppress("MagicNumber") when { - count in 10..25 -> 50L - count > 25 -> 100L - else -> 0 - } - } .mapLatest { statuses -> statuses.associateBy { status -> status.network.id } } - .distinctUntilChanged().collect { result -> channel.send(result) } - } + .distinctUntilChanged() private fun stakingFlow(wallet: UserWallet): Flow>> = if (!wallet.isMultiCurrency) { diff --git a/test/core/build.gradle.kts b/test/core/build.gradle.kts index d34ffdd2aa..ad834bcedb 100644 --- a/test/core/build.gradle.kts +++ b/test/core/build.gradle.kts @@ -10,4 +10,5 @@ dependencies { api(deps.test.junit5) api(deps.test.mockk) api(deps.test.truth) + api(deps.test.turbine) } \ No newline at end of file