Updated on 2026-08-14

This commit is contained in:
Tangem 2026-06-30 15:21:16 +02:00
parent f642bd1e28
commit de3cdf0c47
8 changed files with 380 additions and 48 deletions

View file

@ -8,7 +8,6 @@ import com.tangem.data.common.api.safeApiCall
import com.tangem.data.staking.store.P2PEthPoolBalancesStore
import com.tangem.data.staking.store.StakeKitBalancesStore
import com.tangem.data.staking.utils.YieldBalanceRequestBodyFactory
import com.tangem.datasource.api.common.response.ApiResponse
import com.tangem.datasource.api.ethpool.P2PEthPoolApi
import com.tangem.datasource.api.ethpool.models.response.P2PEthPoolAccountResponse
import com.tangem.datasource.api.stakekit.StakeKitApi
@ -25,7 +24,6 @@ import com.tangem.domain.models.wallet.UserWallet
import com.tangem.domain.models.wallet.UserWalletId
import com.tangem.domain.models.wallet.isMultiCurrency
import com.tangem.domain.staking.model.StakingIntegrationID
import com.tangem.domain.staking.model.ethpool.P2PEthPoolStakingConfig
import com.tangem.domain.staking.model.ethpool.P2PEthPoolVault
import com.tangem.domain.staking.multi.MultiStakingBalanceFetcher
import com.tangem.utils.coroutines.CoroutineDispatcherProvider
@ -65,6 +63,11 @@ internal class DefaultMultiStakingBalanceFetcher @Inject constructor(
private val dispatchers: CoroutineDispatcherProvider,
) : MultiStakingBalanceFetcher {
private val p2pAccountsFetcher = P2PEthPoolAccountsFetcher(
p2pEthPoolApi = p2pEthPoolApi,
dispatchers = dispatchers,
)
override suspend fun invoke(params: MultiStakingBalanceFetcher.Params): Either<Throwable, Unit> {
TangemLogger.i("Start fetching staking balances for params:\n$params")
@ -195,50 +198,7 @@ internal class DefaultMultiStakingBalanceFetcher @Inject constructor(
vaults: List<P2PEthPoolVault>,
addresses: Set<String>,
): Set<P2PEthPoolAccountResponse> {
val responses = mutableSetOf<P2PEthPoolAccountResponse>()
for (vault in vaults) {
for (address in addresses) {
runSuspendCatching {
val response = p2pEthPoolApi.getAccountInfo(
network = P2PEthPoolStakingConfig.activeNetwork.value,
delegatorAddress = address,
vaultAddress = vault.vaultAddress,
)
when (response) {
is ApiResponse.Success -> {
val data = response.data
if (data.error != null) {
TangemLogger.w(
"P2PEthPool API returned error for vault ${vault.vaultAddress}, " +
"address $address: ${data.error ?: "error"}",
)
} else {
val result = requireNotNull(data.result) {
"Result is null in successful response"
}
responses.add(result)
}
}
is ApiResponse.Error -> {
TangemLogger.w(
"Failed to fetch P2PEthPool balance for vault ${vault.vaultAddress}, " +
"address $address",
response.cause,
)
}
}
}.onFailure { error ->
TangemLogger.w(
"Failed to fetch P2PEthPool balance for vault ${vault.vaultAddress}, address $address",
error,
)
}
}
}
return responses
return p2pAccountsFetcher.fetchBatch(vaults = vaults, addresses = addresses)
}
private inline fun checkIsSupportedByWalletOrElse(userWalletId: UserWalletId, ifNotSupported: (Throwable) -> Unit) {

View file

@ -0,0 +1,102 @@
package com.tangem.data.staking.multi
import com.tangem.datasource.api.common.response.ApiResponse
import com.tangem.datasource.api.ethpool.P2PEthPoolApi
import com.tangem.datasource.api.ethpool.models.request.P2PEthPoolAccountsListRequest
import com.tangem.datasource.api.ethpool.models.response.P2PEthPoolAccountResponse
import com.tangem.datasource.api.ethpool.models.response.P2PEthPoolAccountsListResponse
import com.tangem.datasource.api.ethpool.models.response.P2PEthPoolResponse
import com.tangem.domain.staking.model.ethpool.P2PEthPoolStakingConfig
import com.tangem.domain.staking.model.ethpool.P2PEthPoolVault
import com.tangem.utils.coroutines.CoroutineDispatcherProvider
import com.tangem.utils.coroutines.runSuspendCatching
import com.tangem.utils.logging.TangemLogger
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.coroutineScope
/**
* Fetches P2P ETH Pool account responses via the batch strategy:
* one POST per vault sending all delegator addresses at once.
*
* @property p2pEthPoolApi P2PEthPool API
* @property dispatchers coroutine dispatcher provider
*
[REDACTED_AUTHOR]
*/
internal class P2PEthPoolAccountsFetcher(
private val p2pEthPoolApi: P2PEthPoolApi,
private val dispatchers: CoroutineDispatcherProvider,
) {
suspend fun fetchBatch(vaults: List<P2PEthPoolVault>, addresses: Set<String>): Set<P2PEthPoolAccountResponse> =
coroutineScope {
val request = P2PEthPoolAccountsListRequest(delegatorAddresses = addresses.toList())
vaults
.map { vault ->
async(dispatchers.io) {
runSuspendCatching {
val response = p2pEthPoolApi.getAccountsList(
network = P2PEthPoolStakingConfig.activeNetwork.value,
vaultAddress = vault.vaultAddress,
body = request,
)
mapBatchVaultResponse(vault = vault, response = response)
}.getOrElse { error ->
TangemLogger.w(
"Failed to fetch P2PEthPool batch balances for vault ${vault.vaultAddress}",
error,
)
emptyList()
}
}
}
.awaitAll()
.flatten()
.toSet()
}
private fun mapBatchVaultResponse(
vault: P2PEthPoolVault,
response: ApiResponse<P2PEthPoolResponse<P2PEthPoolAccountsListResponse>>,
): List<P2PEthPoolAccountResponse> {
return when (response) {
is ApiResponse.Success -> {
val data = response.data
if (data.error != null) {
TangemLogger.w(
"P2PEthPool batch API returned error for vault " +
"${vault.vaultAddress}: ${data.error}",
)
emptyList()
} else {
val result = requireNotNull(data.result) {
"Result is null in successful response"
}
result.list.mapNotNull { item ->
if (item.error != null) {
TangemLogger.w(
"P2PEthPool batch item error for vault " +
"${vault.vaultAddress}, address " +
"${item.delegatorAddress}: ${item.error}",
)
null
} else {
item.account
}
}
}
}
is ApiResponse.Error -> {
TangemLogger.w(
"Failed to fetch P2PEthPool batch balances for vault " +
"${vault.vaultAddress}",
response.cause,
)
emptyList()
}
}
}
}

View file

@ -10,12 +10,15 @@ import com.tangem.data.staking.utils.YieldBalanceRequestBodyFactory
import com.tangem.datasource.api.common.response.ApiResponse
import com.tangem.datasource.api.common.response.ApiResponseError
import com.tangem.datasource.api.ethpool.P2PEthPoolApi
import com.tangem.datasource.api.ethpool.models.request.P2PEthPoolAccountsListRequest
import com.tangem.datasource.api.ethpool.models.response.*
import com.tangem.datasource.api.stakekit.StakeKitApi
import com.tangem.datasource.api.stakekit.models.response.model.YieldBalanceWrapperDTO
import com.tangem.datasource.local.token.P2PEthPoolVaultsStore
import com.tangem.datasource.local.token.StakingYieldsStore
import com.tangem.domain.common.wallets.UserWalletsListRepository
import com.tangem.domain.models.staking.StakingID
import com.tangem.domain.staking.model.ethpool.P2PEthPoolVault
import com.tangem.domain.staking.multi.MultiStakingBalanceFetcher
import com.tangem.test.core.assertEitherLeft
import com.tangem.test.core.assertEitherRight
@ -26,6 +29,7 @@ import kotlinx.coroutines.test.runTest
import org.junit.jupiter.api.BeforeEach
import org.junit.jupiter.api.Test
import org.junit.jupiter.api.TestInstance
import java.math.BigDecimal
/**
[REDACTED_AUTHOR]
@ -54,7 +58,14 @@ internal class DefaultMultiStakingBalanceFetcherTest {
@BeforeEach
fun resetMocks() {
clearMocks(userWalletsListRepository, stakingYieldsStore, stakeKitBalancesStore, stakeKitApi)
clearMocks(
userWalletsListRepository,
stakingYieldsStore,
stakeKitBalancesStore,
stakeKitApi,
p2pEthPoolApi,
p2pEthPoolVaultsStore,
)
}
@Test
@ -359,6 +370,88 @@ internal class DefaultMultiStakingBalanceFetcherTest {
assertEitherLeft(actual, expected)
}
@Test
fun `fetch P2P balances via batch endpoint`() = runTest {
// Arrange
val params = MultiStakingBalanceFetcher.Params(userWalletId, setOf(p2pId1, p2pId2))
every { userWalletsListRepository.userWallets } returns MutableStateFlow(listOf(userWallet))
coEvery { p2pEthPoolVaultsStore.getSync() } returns listOf(vault(VAULT_A), vault(VAULT_B))
coEvery { p2pEthPoolApi.getAccountsList(any(), VAULT_A, any()) } returns
accountsListSuccess(accountResponse(ADDR_1, VAULT_A))
coEvery { p2pEthPoolApi.getAccountsList(any(), VAULT_B, any()) } returns
accountsListSuccess(accountResponse(ADDR_2, VAULT_B))
// Actual
val actual = fetcher.invoke(params)
// Assert
coVerify(exactly = 1) {
p2pEthPoolApi.getAccountsList(
network = any(),
vaultAddress = VAULT_A,
body = match { it.delegatorAddresses.containsAll(listOf(ADDR_1, ADDR_2)) },
)
}
coVerify(exactly = 1) {
p2pEthPoolApi.getAccountsList(network = any(), vaultAddress = VAULT_B, body = any())
}
coVerify(inverse = true) { p2pEthPoolApi.getAccountInfo(any(), any(), any()) }
coVerify { p2PEthPoolBalancesStore.storeActual(userWalletId = userWalletId, values = any()) }
assertEitherRight(actual)
}
@Test
fun `fetch P2P batch maps per-item error to missing stakingId`() = runTest {
// Arrange
val params = MultiStakingBalanceFetcher.Params(userWalletId, setOf(p2pId1, p2pId2))
every { userWalletsListRepository.userWallets } returns MutableStateFlow(listOf(userWallet))
coEvery { p2pEthPoolVaultsStore.getSync() } returns listOf(vault(VAULT_A))
coEvery { p2pEthPoolApi.getAccountsList(any(), VAULT_A, any()) } returns
ApiResponse.Success(
P2PEthPoolResponse(
error = null,
result = P2PEthPoolAccountsListResponse(
list = listOf(
P2PEthPoolAccountListItem(
delegatorAddress = ADDR_1,
account = accountResponse(ADDR_1, VAULT_A),
error = null,
),
P2PEthPoolAccountListItem(
delegatorAddress = ADDR_2,
account = null,
error = P2PEthPoolErrorDetailsDTO(
code = 127108,
message = "invalid",
name = null,
errors = null,
),
),
),
),
),
)
// Actual
val actual = fetcher.invoke(params)
// Assert
coVerify { p2PEthPoolBalancesStore.storeActual(userWalletId = userWalletId, values = any()) }
coVerify {
p2PEthPoolBalancesStore.storeError(
userWalletId = userWalletId,
stakingIds = match { it == setOf(p2pId2) },
)
}
assertEitherRight(actual)
}
private companion object {
val userWallet = MockUserWalletFactory.create()
val userWalletId = userWallet.walletId
@ -370,5 +463,55 @@ internal class DefaultMultiStakingBalanceFetcherTest {
)
val tonAndSolanaIds = setOf(tonId, solanaId)
const val ADDR_1 = "0x1111111111111111111111111111111111111111"
const val ADDR_2 = "0x2222222222222222222222222222222222222222"
const val VAULT_A = "0xVaultAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"
const val VAULT_B = "0xVaultBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB"
val p2pId1 = StakingID(integrationId = "p2p-ethereum-pooled", address = ADDR_1)
val p2pId2 = StakingID(integrationId = "p2p-ethereum-pooled", address = ADDR_2)
fun vault(address: String) = P2PEthPoolVault(
vaultAddress = address,
displayName = "Vault",
apy = BigDecimal("4.5"),
baseApy = BigDecimal("4.0"),
capacity = BigDecimal("1000"),
totalAssets = BigDecimal("100"),
feePercent = BigDecimal("10"),
isPrivate = false,
isGenesis = false,
isSmoothingPool = true,
isErc20 = false,
tokenName = null,
tokenSymbol = null,
createdAt = 0L,
)
fun accountResponse(address: String, vaultAddress: String) = P2PEthPoolAccountResponse(
delegatorAddress = address,
vaultAddress = vaultAddress,
stake = P2PEthPoolStakeDTO(assets = BigDecimal("1.5"), totalEarnedAssets = BigDecimal("0.1")),
availableToUnstake = BigDecimal.ZERO,
availableToWithdraw = BigDecimal.ZERO,
exitQueue = P2PEthPoolExitQueueDTO(total = BigDecimal.ZERO, requests = emptyList()),
)
fun accountsListSuccess(vararg accounts: P2PEthPoolAccountResponse) =
ApiResponse.Success(
P2PEthPoolResponse(
error = null,
result = P2PEthPoolAccountsListResponse(
list = accounts.map {
P2PEthPoolAccountListItem(
delegatorAddress = it.delegatorAddress,
account = it,
error = null,
)
},
),
),
)
}
}