Updated on 2026-08-14
This commit is contained in:
parent
cb26294b01
commit
27995dca39
1 changed files with 50 additions and 67 deletions
|
|
@ -25,12 +25,19 @@ import com.tangem.domain.models.network.Network
|
||||||
import com.tangem.domain.models.wallet.UserWalletId
|
import com.tangem.domain.models.wallet.UserWalletId
|
||||||
import com.tangem.utils.coroutines.CoroutineDispatcherProvider
|
import com.tangem.utils.coroutines.CoroutineDispatcherProvider
|
||||||
import com.tangem.utils.logging.TangemLogger
|
import com.tangem.utils.logging.TangemLogger
|
||||||
import kotlinx.coroutines.async
|
import kotlinx.coroutines.CancellationException
|
||||||
import kotlinx.coroutines.awaitAll
|
import kotlinx.coroutines.TimeoutCancellationException
|
||||||
import kotlinx.coroutines.coroutineScope
|
import kotlinx.coroutines.coroutineScope
|
||||||
|
import kotlinx.coroutines.ensureActive
|
||||||
import kotlinx.coroutines.flow.Flow
|
import kotlinx.coroutines.flow.Flow
|
||||||
import kotlinx.coroutines.flow.MutableStateFlow
|
import kotlinx.coroutines.flow.MutableStateFlow
|
||||||
|
import kotlinx.coroutines.launch
|
||||||
|
import kotlinx.coroutines.sync.Mutex
|
||||||
|
import kotlinx.coroutines.sync.Semaphore
|
||||||
|
import kotlinx.coroutines.sync.withLock
|
||||||
|
import kotlinx.coroutines.sync.withPermit
|
||||||
import kotlinx.coroutines.withContext
|
import kotlinx.coroutines.withContext
|
||||||
|
import kotlinx.coroutines.withTimeout
|
||||||
import java.math.BigDecimal
|
import java.math.BigDecimal
|
||||||
import java.util.concurrent.ConcurrentHashMap
|
import java.util.concurrent.ConcurrentHashMap
|
||||||
|
|
||||||
|
|
@ -116,22 +123,29 @@ internal class DefaultAssetsDiscoveryRepository(
|
||||||
val assetsDiscoveryStore = assetsDiscoveryStoreFactory.provide(userWalletId)
|
val assetsDiscoveryStore = assetsDiscoveryStoreFactory.provide(userWalletId)
|
||||||
assetsDiscoveryStore.clear()
|
assetsDiscoveryStore.clear()
|
||||||
|
|
||||||
val batches = networks.chunked(MAX_CONCURRENT_REQUESTS)
|
val totalNetworks = networks.size
|
||||||
var completedNetworks = 0
|
val progressFlow = getProgressFlow(userWalletId)
|
||||||
getProgressFlow(userWalletId).value = AssetsDiscoveryProgress.InProgress(
|
progressFlow.value = AssetsDiscoveryProgress.InProgress(completedNetworks = 0, totalNetworks = totalNetworks)
|
||||||
completedNetworks = 0,
|
|
||||||
totalNetworks = networks.size,
|
|
||||||
)
|
|
||||||
|
|
||||||
for (batch in batches) {
|
val semaphore = Semaphore(permits = MAX_CONCURRENT_REQUESTS)
|
||||||
val batchResults = processBatch(userWalletId, batch)
|
val progressMutex = Mutex()
|
||||||
completedNetworks = handleBatchResults(
|
var completedNetworks = 0
|
||||||
userWalletId = userWalletId,
|
|
||||||
results = batchResults,
|
coroutineScope {
|
||||||
assetsDiscoveryStore = assetsDiscoveryStore,
|
networks.forEach { network ->
|
||||||
completedNetworks = completedNetworks,
|
launch(dispatchers.io) {
|
||||||
totalNetworks = networks.size,
|
val result = semaphore.withPermit { processNetwork(userWalletId, network) }
|
||||||
)
|
handleNetworkResult(result, assetsDiscoveryStore)
|
||||||
|
progressMutex.withLock {
|
||||||
|
completedNetworks++
|
||||||
|
val inProgress = AssetsDiscoveryProgress.InProgress(
|
||||||
|
completedNetworks = completedNetworks,
|
||||||
|
totalNetworks = totalNetworks,
|
||||||
|
)
|
||||||
|
progressFlow.value = inProgress
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -140,38 +154,6 @@ internal class DefaultAssetsDiscoveryRepository(
|
||||||
getProgressFlow(userWalletId).value = AssetsDiscoveryProgress.Completed
|
getProgressFlow(userWalletId).value = AssetsDiscoveryProgress.Completed
|
||||||
}
|
}
|
||||||
|
|
||||||
private suspend fun processBatch(userWalletId: UserWalletId, batch: List<Network>): List<NetworkResult> {
|
|
||||||
return coroutineScope {
|
|
||||||
batch.map { network ->
|
|
||||||
async(dispatchers.io) {
|
|
||||||
processNetwork(userWalletId, network)
|
|
||||||
}
|
|
||||||
}.awaitAll()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private suspend fun handleBatchResults(
|
|
||||||
userWalletId: UserWalletId,
|
|
||||||
results: List<NetworkResult>,
|
|
||||||
assetsDiscoveryStore: AssetsDiscoveryStore,
|
|
||||||
completedNetworks: Int,
|
|
||||||
totalNetworks: Int,
|
|
||||||
): Int {
|
|
||||||
var completed = completedNetworks
|
|
||||||
val progressFlow = getProgressFlow(userWalletId)
|
|
||||||
|
|
||||||
for (result in results) {
|
|
||||||
completed++
|
|
||||||
handleNetworkResult(result, assetsDiscoveryStore)
|
|
||||||
progressFlow.value = AssetsDiscoveryProgress.InProgress(
|
|
||||||
completedNetworks = completed,
|
|
||||||
totalNetworks = totalNetworks,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
return completed
|
|
||||||
}
|
|
||||||
|
|
||||||
private suspend fun handleNetworkResult(result: NetworkResult, assetsDiscoveryStore: AssetsDiscoveryStore) {
|
private suspend fun handleNetworkResult(result: NetworkResult, assetsDiscoveryStore: AssetsDiscoveryStore) {
|
||||||
when (result) {
|
when (result) {
|
||||||
is NetworkResult.Success -> {
|
is NetworkResult.Success -> {
|
||||||
|
|
@ -191,25 +173,25 @@ internal class DefaultAssetsDiscoveryRepository(
|
||||||
|
|
||||||
private suspend fun processNetwork(userWalletId: UserWalletId, network: Network): NetworkResult {
|
private suspend fun processNetwork(userWalletId: UserWalletId, network: Network): NetworkResult {
|
||||||
return try {
|
return try {
|
||||||
val discoveredAssets = discoverAndFilterAssets(userWalletId, network)
|
withTimeout(NETWORK_DISCOVERY_TIMEOUT_MILLIS) {
|
||||||
|
val discoveredAssets = discoverAndFilterAssets(userWalletId, network)
|
||||||
|
|
||||||
if (discoveredAssets.isEmpty()) {
|
val responseTokens = if (discoveredAssets.isEmpty()) {
|
||||||
return NetworkResult.Success(
|
emptyList()
|
||||||
networkId = network.rawId,
|
} else {
|
||||||
responseTokens = emptyList(),
|
enrichTokensWithCatalog(discoveredAssets, network)
|
||||||
)
|
.filter { it.contractAddress != null }
|
||||||
|
.map { it.toResponseToken() }
|
||||||
|
}
|
||||||
|
|
||||||
|
ensureActive()
|
||||||
|
|
||||||
|
NetworkResult.Success(networkId = network.rawId, responseTokens = responseTokens)
|
||||||
}
|
}
|
||||||
|
} catch (e: TimeoutCancellationException) {
|
||||||
val enrichedTokens = enrichTokensWithCatalog(discoveredAssets, network)
|
NetworkResult.Error(networkId = network.rawId, cause = e)
|
||||||
|
} catch (e: CancellationException) {
|
||||||
val responseTokens = enrichedTokens
|
throw e
|
||||||
.filter { it.contractAddress != null }
|
|
||||||
.map { it.toResponseToken() }
|
|
||||||
|
|
||||||
NetworkResult.Success(
|
|
||||||
networkId = network.rawId,
|
|
||||||
responseTokens = responseTokens,
|
|
||||||
)
|
|
||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
NetworkResult.Error(networkId = network.rawId, cause = e)
|
NetworkResult.Error(networkId = network.rawId, cause = e)
|
||||||
}
|
}
|
||||||
|
|
@ -352,6 +334,7 @@ internal class DefaultAssetsDiscoveryRepository(
|
||||||
}
|
}
|
||||||
|
|
||||||
companion object {
|
companion object {
|
||||||
private const val MAX_CONCURRENT_REQUESTS = 3
|
private const val MAX_CONCURRENT_REQUESTS = 20
|
||||||
|
private const val NETWORK_DISCOVERY_TIMEOUT_MILLIS = 30_000L
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Loading…
Add table
Add a link
Reference in a new issue