solution: track balance of a ERC20 token
This commit is contained in:
@@ -20,10 +20,7 @@ import com.fasterxml.jackson.core.Version
|
||||
import com.fasterxml.jackson.databind.DeserializationFeature
|
||||
import com.fasterxml.jackson.databind.ObjectMapper
|
||||
import com.fasterxml.jackson.databind.module.SimpleModule
|
||||
import io.emeraldpay.dshackle.config.CacheConfig
|
||||
import io.emeraldpay.dshackle.config.MainConfig
|
||||
import io.emeraldpay.dshackle.config.MainConfigReader
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import io.emeraldpay.dshackle.config.*
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.beans.factory.annotation.Qualifier
|
||||
@@ -103,4 +100,9 @@ open class Config(
|
||||
return mainConfig.cache ?: CacheConfig()
|
||||
}
|
||||
|
||||
@Bean
|
||||
open fun tokensConfig(@Autowired mainConfig: MainConfig): TokensConfig {
|
||||
return mainConfig.tokens ?: TokensConfig(emptyList())
|
||||
}
|
||||
|
||||
}
|
||||
@@ -64,9 +64,13 @@ class BlockchainRpc(
|
||||
override fun subscribeBalance(requestMono: Mono<BlockchainOuterClass.BalanceRequest>): Flux<BlockchainOuterClass.AddressBalance> {
|
||||
return requestMono.flatMapMany { request ->
|
||||
val chain = Chain.byId(request.asset.chainValue)
|
||||
val asset = request.asset.code.toLowerCase()
|
||||
try {
|
||||
trackAddress.find { it.isSupported(chain) }?.subscribe(request)
|
||||
?: Flux.error(SilentException.UnsupportedBlockchain(chain))
|
||||
trackAddress.find { it.isSupported(chain, asset) }?.subscribe(request)
|
||||
?: Flux.error<BlockchainOuterClass.AddressBalance>(SilentException.UnsupportedBlockchain(chain))
|
||||
.doOnSubscribe {
|
||||
log.error("Balance for $chain:$asset is not supported")
|
||||
}
|
||||
} catch (t: Throwable) {
|
||||
log.error("Internal error during Balance Subscription", t)
|
||||
Flux.error<BlockchainOuterClass.AddressBalance>(IllegalStateException("Internal Error"))
|
||||
@@ -77,9 +81,13 @@ class BlockchainRpc(
|
||||
override fun getBalance(requestMono: Mono<BlockchainOuterClass.BalanceRequest>): Flux<BlockchainOuterClass.AddressBalance> {
|
||||
return requestMono.flatMapMany { request ->
|
||||
val chain = Chain.byId(request.asset.chainValue)
|
||||
val asset = request.asset.code.toLowerCase()
|
||||
try {
|
||||
trackAddress.find { it.isSupported(chain) }?.getBalance(request)
|
||||
?: Flux.error(SilentException.UnsupportedBlockchain(chain))
|
||||
trackAddress.find { it.isSupported(chain, asset) }?.getBalance(request)
|
||||
?: Flux.error<BlockchainOuterClass.AddressBalance>(SilentException.UnsupportedBlockchain(chain))
|
||||
.doOnSubscribe {
|
||||
log.error("Balance for $chain:$asset is not supported")
|
||||
}
|
||||
} catch (t: Throwable) {
|
||||
log.error("Internal error during Balance Request", t)
|
||||
Flux.error<BlockchainOuterClass.AddressBalance>(IllegalStateException("Internal Error"))
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
package io.emeraldpay.dshackle.rpc
|
||||
|
||||
import io.emeraldpay.api.proto.Common
|
||||
import io.infinitape.etherjar.domain.Address
|
||||
import org.slf4j.LoggerFactory
|
||||
import reactor.core.publisher.Flux
|
||||
|
||||
class EthereumAddresses {
|
||||
|
||||
companion object {
|
||||
private val log = LoggerFactory.getLogger(EthereumAddresses::class.java)
|
||||
}
|
||||
|
||||
fun extract(addresses: Common.AnyAddress): Flux<Address> {
|
||||
return when (addresses.addrTypeCase) {
|
||||
Common.AnyAddress.AddrTypeCase.ADDRESS_SINGLE ->
|
||||
Flux.just(Address.from(addresses.addressSingle.address))
|
||||
Common.AnyAddress.AddrTypeCase.ADDRESS_MULTI ->
|
||||
Flux.fromIterable(addresses.addressMulti.addressesList)
|
||||
.map { Address.from(it.address) }
|
||||
else -> {
|
||||
log.error("Unsupported address type: ${addresses.addrTypeCase}")
|
||||
Flux.empty()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -25,7 +25,7 @@ import reactor.core.publisher.Mono
|
||||
*/
|
||||
interface TrackAddress {
|
||||
|
||||
fun isSupported(chain: Chain): Boolean
|
||||
fun isSupported(chain: Chain, asset: String): Boolean
|
||||
fun getBalance(request: BlockchainOuterClass.BalanceRequest): Flux<BlockchainOuterClass.AddressBalance>
|
||||
fun subscribe(request: BlockchainOuterClass.BalanceRequest): Flux<BlockchainOuterClass.AddressBalance>
|
||||
|
||||
|
||||
@@ -41,8 +41,9 @@ class TrackBitcoinAddress(
|
||||
private val log = LoggerFactory.getLogger(TrackBitcoinAddress::class.java)
|
||||
}
|
||||
|
||||
override fun isSupported(chain: Chain): Boolean {
|
||||
return BlockchainType.fromBlockchain(chain) == BlockchainType.BITCOIN && multistreamHolder.isAvailable(chain)
|
||||
override fun isSupported(chain: Chain, asset: String): Boolean {
|
||||
return (asset == "bitcoin" || asset == "btc" || asset == "satoshi")
|
||||
&& BlockchainType.fromBlockchain(chain) == BlockchainType.BITCOIN && multistreamHolder.isAvailable(chain)
|
||||
}
|
||||
|
||||
fun allAddresses(request: BlockchainOuterClass.BalanceRequest): List<String>? {
|
||||
|
||||
132
src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt
Normal file
132
src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackERC20Address.kt
Normal file
@@ -0,0 +1,132 @@
|
||||
package io.emeraldpay.dshackle.rpc
|
||||
|
||||
import io.emeraldpay.api.proto.BlockchainOuterClass
|
||||
import io.emeraldpay.api.proto.Common
|
||||
import io.emeraldpay.dshackle.BlockchainType
|
||||
import io.emeraldpay.dshackle.SilentException
|
||||
import io.emeraldpay.dshackle.config.TokensConfig
|
||||
import io.emeraldpay.dshackle.upstream.MultistreamHolder
|
||||
import io.emeraldpay.dshackle.upstream.Selector
|
||||
import io.emeraldpay.dshackle.upstream.ethereum.EthereumMultistream
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
|
||||
import io.emeraldpay.grpc.Chain
|
||||
import io.infinitape.etherjar.domain.Address
|
||||
import io.infinitape.etherjar.erc20.ERC20Token
|
||||
import io.infinitape.etherjar.hex.Hex32
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.stereotype.Service
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.core.publisher.Mono
|
||||
import java.math.BigInteger
|
||||
import javax.annotation.PostConstruct
|
||||
|
||||
@Service
|
||||
class TrackERC20Address(
|
||||
@Autowired private val multistreamHolder: MultistreamHolder,
|
||||
@Autowired private val tokensConfig: TokensConfig
|
||||
) : TrackAddress {
|
||||
|
||||
companion object {
|
||||
private val log = LoggerFactory.getLogger(TrackERC20Address::class.java)
|
||||
}
|
||||
|
||||
private val ethereumAddresses = EthereumAddresses()
|
||||
private val tokens: MutableMap<TokenId, TokenDefinition> = HashMap()
|
||||
|
||||
@PostConstruct
|
||||
fun init() {
|
||||
tokensConfig.tokens.forEach { token ->
|
||||
val chain = token.blockchain!!
|
||||
val asset = token.name!!.toLowerCase()
|
||||
val id = TokenId(chain, asset)
|
||||
val definition = TokenDefinition(
|
||||
chain, asset,
|
||||
ERC20Token(Address.from(token.address))
|
||||
)
|
||||
tokens[id] = definition
|
||||
log.info("Enable ERC20 balance for $chain:$asset")
|
||||
}
|
||||
}
|
||||
|
||||
override fun isSupported(chain: Chain, asset: String): Boolean {
|
||||
return tokens.containsKey(TokenId(chain, asset.toLowerCase())) &&
|
||||
BlockchainType.fromBlockchain(chain) == BlockchainType.ETHEREUM && multistreamHolder.isAvailable(chain)
|
||||
}
|
||||
|
||||
override fun getBalance(request: BlockchainOuterClass.BalanceRequest): Flux<BlockchainOuterClass.AddressBalance> {
|
||||
val chain = Chain.byId(request.asset.chainValue)
|
||||
val asset = request.asset.code.toLowerCase()
|
||||
val tokenDefinition = tokens[TokenId(chain, asset)] ?: return Flux.empty()
|
||||
return ethereumAddresses.extract(request.address)
|
||||
.map { TrackedAddress(chain, it, tokenDefinition.token, tokenDefinition.name) }
|
||||
.flatMap { addr -> getBalance(addr).map(addr::withBalance) }
|
||||
.map { buildResponse(it) }
|
||||
}
|
||||
|
||||
override fun subscribe(request: BlockchainOuterClass.BalanceRequest): Flux<BlockchainOuterClass.AddressBalance> {
|
||||
val chain = Chain.byId(request.asset.chainValue)
|
||||
val asset = request.asset.code.toLowerCase()
|
||||
val tokenDefinition = tokens[TokenId(chain, asset)] ?: return Flux.empty()
|
||||
val head = multistreamHolder.getUpstream(chain)?.getHead()?.getFlux() ?: Flux.empty()
|
||||
|
||||
return ethereumAddresses.extract(request.address)
|
||||
.map { TrackedAddress(chain, it, tokenDefinition.token, tokenDefinition.name) }
|
||||
.flatMap { addr ->
|
||||
val current = getBalance(addr)
|
||||
val updates = head.flatMap { getBalance(addr) }
|
||||
Flux.concat(current, updates)
|
||||
.distinctUntilChanged()
|
||||
.map { addr.withBalance(it) }
|
||||
}
|
||||
.map { buildResponse(it) }
|
||||
}
|
||||
|
||||
fun getBalance(addr: TrackedAddress): Mono<BigInteger> {
|
||||
return getUpstream(addr.chain)
|
||||
.getDirectApi(Selector.empty)
|
||||
.flatMap { api ->
|
||||
api.read(prepareEthCall(addr.token, addr.address))
|
||||
.flatMap(JsonRpcResponse::requireStringResult)
|
||||
.map {
|
||||
Hex32.from(it).asQuantity().value
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun prepareEthCall(token: ERC20Token, target: Address): JsonRpcRequest {
|
||||
val call = token
|
||||
.readBalanceOf(target)
|
||||
.toJson()
|
||||
return JsonRpcRequest("eth_call", listOf(call, "latest"))
|
||||
}
|
||||
|
||||
fun getUpstream(chain: Chain): EthereumMultistream {
|
||||
return multistreamHolder.getUpstream(chain)?.cast(EthereumMultistream::class.java)
|
||||
?: throw SilentException.UnsupportedBlockchain(chain)
|
||||
}
|
||||
|
||||
private fun buildResponse(address: TrackedAddress): BlockchainOuterClass.AddressBalance {
|
||||
return BlockchainOuterClass.AddressBalance.newBuilder()
|
||||
.setBalance(address.balance!!.toString(10))
|
||||
.setAsset(Common.Asset.newBuilder()
|
||||
.setChainValue(address.chain.id)
|
||||
.setCode(address.tokenName.toUpperCase()))
|
||||
.setAddress(Common.SingleAddress.newBuilder().setAddress(address.address.toHex()))
|
||||
.build()
|
||||
}
|
||||
|
||||
class TrackedAddress(val chain: Chain,
|
||||
val address: Address,
|
||||
val token: ERC20Token,
|
||||
val tokenName: String,
|
||||
val balance: BigInteger? = null
|
||||
) {
|
||||
fun withBalance(balance: BigInteger) = TrackedAddress(chain, address, token, tokenName, balance)
|
||||
}
|
||||
|
||||
data class TokenId(val chain: Chain, val name: String)
|
||||
data class TokenDefinition(val chain: Chain, val name: String, val token: ERC20Token)
|
||||
|
||||
}
|
||||
@@ -38,9 +38,11 @@ class TrackEthereumAddress(
|
||||
) : TrackAddress {
|
||||
|
||||
private val log = LoggerFactory.getLogger(TrackEthereumAddress::class.java)
|
||||
private val ethereumAddresses = EthereumAddresses()
|
||||
|
||||
override fun isSupported(chain: Chain): Boolean {
|
||||
return BlockchainType.fromBlockchain(chain) == BlockchainType.ETHEREUM && multistreamHolder.isAvailable(chain)
|
||||
override fun isSupported(chain: Chain, asset: String): Boolean {
|
||||
return asset == "ether" &&
|
||||
BlockchainType.fromBlockchain(chain) == BlockchainType.ETHEREUM && multistreamHolder.isAvailable(chain)
|
||||
}
|
||||
|
||||
override fun getBalance(request: BlockchainOuterClass.BalanceRequest): Flux<BlockchainOuterClass.AddressBalance> {
|
||||
@@ -99,16 +101,8 @@ class TrackEthereumAddress(
|
||||
if (request.asset.code?.toLowerCase() != "ether") {
|
||||
return Flux.error(SilentException("Unsupported asset ${request.asset.code}"))
|
||||
}
|
||||
return when (request.address.addrTypeCase) {
|
||||
Common.AnyAddress.AddrTypeCase.ADDRESS_SINGLE ->
|
||||
Flux.just(createAddress(request.address.addressSingle, chain))
|
||||
Common.AnyAddress.AddrTypeCase.ADDRESS_MULTI ->
|
||||
Flux.fromIterable(request.address.addressMulti.addressesList)
|
||||
.map { createAddress(it, chain) }
|
||||
else -> {
|
||||
log.error("Unsupported address type: ${request.address.addrTypeCase}")
|
||||
Flux.empty()
|
||||
}
|
||||
return ethereumAddresses.extract(request.address).map {
|
||||
TrackedAddress(chain, it)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -48,5 +48,7 @@ class JsonRpcRequest(
|
||||
return result
|
||||
}
|
||||
|
||||
|
||||
override fun toString(): String {
|
||||
return String(this.toJson())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user