From bd58ba977440947b2f959488a0a9ae72d56ea8ff Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Mon, 27 Apr 2020 00:05:02 -0400 Subject: [PATCH] problem: bitcoin mempool can be queries too often solution: introduce data access wrappers, in this case caching mempool state in memory --- .../dshackle/rpc/TrackBitcoinAddress.kt | 8 +-- .../emeraldpay/dshackle/rpc/TrackBitcoinTx.kt | 27 +++---- .../dshackle/startup/ConfiguredUpstreams.kt | 7 +- .../dshackle/upstream/CurrentUpstreams.kt | 8 +-- .../upstream/bitcoin/BitcoinChainUpstreams.kt | 4 +- .../dshackle/upstream/bitcoin/BitcoinData.kt | 33 +++++++++ .../upstream/bitcoin/BitcoinRpcHead.kt | 3 +- .../upstream/bitcoin/BitcoinUpstream.kt | 22 ++++-- .../bitcoin/BitcoinUpstreamValidator.kt | 5 +- .../upstream/bitcoin/CachingMempoolData.kt | 71 +++++++++++++++++++ .../{BitcoinApi.kt => DirectBitcoinApi.kt} | 7 +- .../rpc/TrackBitcoinAddressSpec.groovy | 4 +- .../dshackle/rpc/TrackBitcoinTxSpec.groovy | 63 ++++++++++------ .../bitcoin/BitcoinRpcHeadSpec.groovy | 2 +- ...pec.groovy => DirectBitcoinApiSpec.groovy} | 6 +- 15 files changed, 197 insertions(+), 73 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinData.kt create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt rename src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/{BitcoinApi.kt => DirectBitcoinApi.kt} (94%) rename src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/{BitcoinApiSpec.groovy => DirectBitcoinApiSpec.groovy} (96%) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt index 3ecde39b..0ce071b1 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt @@ -6,7 +6,7 @@ import io.emeraldpay.dshackle.BlockchainType import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstreams -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.DirectBitcoinApi import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory import org.springframework.beans.factory.annotation.Autowired @@ -48,7 +48,7 @@ class TrackBitcoinAddress( } } - fun requestBalances(chain: Chain, api: BitcoinApi, addresses: List): Flux { + fun requestBalances(chain: Chain, api: DirectBitcoinApi, addresses: List): Flux { return api.executeAndResult(0, "listunspent", emptyList(), List::class.java) .flatMapMany { unspents -> val result = getTotal(chain, addresses, unspents) @@ -58,7 +58,7 @@ class TrackBitcoinAddress( override fun getBalance(request: BlockchainOuterClass.BalanceRequest): Flux { val chain = Chain.byId(request.asset.chainValue) - val upstream = upstreams.getUpstream(chain)?.castApi(BitcoinApi::class.java) + val upstream = upstreams.getUpstream(chain)?.castApi(DirectBitcoinApi::class.java) ?: return Flux.error(SilentException.UnsupportedBlockchain(request.asset.chainValue)) val addresses = allAddresses(request) ?: return Flux.error(SilentException("Unsupported address")) if (addresses.isEmpty()) { @@ -107,7 +107,7 @@ class TrackBitcoinAddress( override fun subscribe(request: BlockchainOuterClass.BalanceRequest): Flux { val chain = Chain.byId(request.asset.chainValue) - val upstream = upstreams.getUpstream(chain)?.castApi(BitcoinApi::class.java) + val upstream = upstreams.getUpstream(chain)?.castApi(DirectBitcoinApi::class.java) ?: return Flux.error(SilentException.UnsupportedBlockchain(request.asset.chainValue)) val addresses = allAddresses(request) ?: return Flux.error(SilentException("Unsupported address")) if (addresses.isEmpty()) { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinTx.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinTx.kt index 4fb191df..1a0411d3 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinTx.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinTx.kt @@ -6,9 +6,9 @@ import io.emeraldpay.api.proto.Common import io.emeraldpay.dshackle.BlockchainType import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.upstream.Selector -import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstreams -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.DirectBitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream import io.emeraldpay.dshackle.upstream.bitcoin.ExtractBlock import io.emeraldpay.grpc.Chain import org.slf4j.LoggerFactory @@ -37,7 +37,7 @@ class TrackBitcoinTx( override fun subscribe(request: BlockchainOuterClass.TxStatusRequest): Flux { val chain = Chain.byId(request.chainValue) - val upstream = upstreams.getUpstream(chain)?.castApi(BitcoinApi::class.java) + val upstream = upstreams.getUpstream(chain)?.cast(BitcoinUpstream::class.java, DirectBitcoinApi::class.java) ?: return Flux.error(SilentException.UnsupportedBlockchain(chain)) val txid = request.txId val confirmations = max(min(1, request.confirmationLimit), 12) @@ -48,7 +48,7 @@ class TrackBitcoinTx( }.map(this::asProto) } - fun subscribe(chain: Chain, api: BitcoinApi, upstream: Upstream, txid: String): Flux { + fun subscribe(chain: Chain, api: DirectBitcoinApi, upstream: BitcoinUpstream, txid: String): Flux { return loadExisting(api, txid) .flatMapMany { status -> if (status.mined) { @@ -56,7 +56,7 @@ class TrackBitcoinTx( //without publishing an empty TxStatus first continueWithMined(api, upstream, status) } else { - loadMempool(api, txid) + loadMempool(upstream, txid) .flatMapMany { tx -> val next = if (tx.found) { untilMined(upstream, tx) @@ -70,7 +70,7 @@ class TrackBitcoinTx( } } - fun continueWithMined(api: BitcoinApi, upstream: Upstream, status: TxStatus): Flux { + fun continueWithMined(api: DirectBitcoinApi, upstream: BitcoinUpstream, status: TxStatus): Flux { return api.getBlock(status.blockHash!!) .map { block -> TxStatus(status.txid, true, ExtractBlock.getHeight(block), true, status.blockHash, ExtractBlock.getTime(block), ExtractBlock.getDifficulty(block)) @@ -79,10 +79,10 @@ class TrackBitcoinTx( } } - fun untilFound(chain: Chain, api: BitcoinApi, upstream: Upstream, txid: String): Flux { + fun untilFound(chain: Chain, api: DirectBitcoinApi, upstream: BitcoinUpstream, txid: String): Flux { return Flux.interval(Duration.ofSeconds(1)) .take(Duration.ofMinutes(10)) - .flatMap { loadMempool(api, txid) } + .flatMap { loadMempool(upstream, txid) } .skipUntil { it.found } .flatMap { subscribe(chain, api, upstream, txid) } .doOnError { t -> @@ -90,7 +90,7 @@ class TrackBitcoinTx( } } - fun untilMined(upstream: Upstream, tx: TxStatus): Mono { + fun untilMined(upstream: BitcoinUpstream, tx: TxStatus): Mono { return upstream.getHead().getFlux().flatMap { upstream.getApi(Selector.empty).flatMap { api -> loadExisting(api, tx.txid) @@ -98,13 +98,13 @@ class TrackBitcoinTx( }.single() } - fun withConfirmations(upstream: Upstream, tx: TxStatus): Flux { + fun withConfirmations(upstream: BitcoinUpstream, tx: TxStatus): Flux { return upstream.getHead().getFlux().map { tx.withHead(it.height) } } - fun loadExisting(api: BitcoinApi, txid: String): Mono { + fun loadExisting(api: DirectBitcoinApi, txid: String): Mono { val mined = api.getTx(txid) return mined.map { val block = it["blockhash"] as String? @@ -112,8 +112,9 @@ class TrackBitcoinTx( } } - fun loadMempool(api: BitcoinApi, txid: String): Mono { - val mempool = api.getMempool() + fun loadMempool(upstream: BitcoinUpstream, txid: String): Mono { + println("access: ${upstream.getData()}") + val mempool = upstream.getData().getMempool().get() return mempool.map { if (it.contains(txid)) { TxStatus(txid, found = true, mined = false) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index 4f2735b4..2e09a20f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -20,10 +20,9 @@ import io.emeraldpay.dshackle.BlockchainType import io.emeraldpay.dshackle.FileResolver import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.upstream.CurrentUpstreams -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.DirectBitcoinApi import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinRpcClient import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream -import io.emeraldpay.dshackle.upstream.bitcoin.DefaultBitcoinMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.ManagedCallMethods import io.emeraldpay.dshackle.upstream.ethereum.DirectEthereumApi @@ -134,11 +133,11 @@ open class ConfiguredUpstreams( options: UpstreamsConfig.Options) { val conn = config.connection!! - var rpcApi: BitcoinApi? = null + var rpcApi: DirectBitcoinApi? = null val methods = buildMethods(config, chain) conn.rpc?.let { endpoint -> val rpcClient = BitcoinRpcClient(endpoint.url.toString(), endpoint.basicAuth!!) - rpcApi = BitcoinApi(rpcClient, objectMapper, methods) + rpcApi = DirectBitcoinApi(rpcClient, objectMapper, methods) } rpcApi?.let { api -> val upstream = BitcoinUpstream(config.id diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentUpstreams.kt index 8d0733cf..32eaaa7a 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/CurrentUpstreams.kt @@ -20,7 +20,7 @@ import io.emeraldpay.dshackle.BlockchainType import io.emeraldpay.dshackle.cache.CachesEnabled import io.emeraldpay.dshackle.cache.CachesFactory import io.emeraldpay.dshackle.startup.UpstreamChange -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.DirectBitcoinApi import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinChainUpstreams import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream import io.emeraldpay.dshackle.upstream.bitcoin.DefaultBitcoinMethods @@ -70,10 +70,10 @@ class CurrentUpstreams( } BlockchainType.BITCOIN -> { val up = change.upstream - .cast(BitcoinUpstream::class.java, BitcoinApi::class.java) - val current = chainMapping[chain] as ChainUpstreams? + .cast(BitcoinUpstream::class.java, DirectBitcoinApi::class.java) + val current = chainMapping[chain] as ChainUpstreams? val factory = Callable { - BitcoinChainUpstreams(chain, ArrayList(), cachesFactory.getCaches(chain), objectMapper) as ChainUpstreams + BitcoinChainUpstreams(chain, ArrayList(), cachesFactory.getCaches(chain), objectMapper) as ChainUpstreams } processUpdate(change, up, current, factory) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinChainUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinChainUpstreams.kt index 3f1fed75..4f882450 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinChainUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinChainUpstreams.kt @@ -14,7 +14,7 @@ class BitcoinChainUpstreams( val upstreams: MutableList, caches: Caches, objectMapper: ObjectMapper -) : ChainUpstreams(chain, upstreams as MutableList>, caches, objectMapper) { +) : ChainUpstreams(chain, upstreams as MutableList>, caches, objectMapper) { companion object { private val log = LoggerFactory.getLogger(BitcoinChainUpstreams::class.java) @@ -66,7 +66,7 @@ class BitcoinChainUpstreams( } override fun castApi(apiType: Class): Upstream { - if (!apiType.isAssignableFrom(BitcoinApi::class.java)) { + if (!apiType.isAssignableFrom(DirectBitcoinApi::class.java)) { throw ClassCastException("Cannot cast ${EthereumApi::class.java} to $apiType") } return this as Upstream diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinData.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinData.kt new file mode 100644 index 00000000..41be7173 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinData.kt @@ -0,0 +1,33 @@ +package io.emeraldpay.dshackle.upstream.bitcoin + +import io.emeraldpay.dshackle.upstream.Head +import org.slf4j.LoggerFactory +import org.springframework.context.Lifecycle + +open class BitcoinData( + api: DirectBitcoinApi, + head: Head +) : Lifecycle { + + companion object { + private val log = LoggerFactory.getLogger(BitcoinData::class.java) + } + + private val mempool = CachingMempoolData(api, head) + + open fun getMempool(): CachingMempoolData { + return mempool + } + + override fun isRunning(): Boolean { + return mempool.isRunning + } + + override fun start() { + mempool.start() + } + + override fun stop() { + mempool.stop() + } +} \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt index 837adc94..a6d6f192 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHead.kt @@ -3,7 +3,6 @@ package io.emeraldpay.dshackle.upstream.bitcoin import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.upstream.AbstractHead import io.emeraldpay.dshackle.upstream.Head -import io.emeraldpay.dshackle.upstream.ethereum.EthereumRpcHead import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle import org.springframework.scheduling.concurrent.CustomizableThreadFactory @@ -15,7 +14,7 @@ import java.time.Duration import java.util.concurrent.Executors class BitcoinRpcHead( - private val api: BitcoinApi, + private val api: DirectBitcoinApi, private val extractBlock: ExtractBlock, private val interval: Duration = Duration.ofSeconds(15) ) : Head, AbstractHead(), Lifecycle { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt index f801cd21..4c30c799 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstream.kt @@ -6,21 +6,20 @@ import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.* import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.grpc.Chain -import io.infinitape.etherjar.rpc.JacksonRpcConverter import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle import reactor.core.Disposable import reactor.core.publisher.Mono -class BitcoinUpstream( +open class BitcoinUpstream( id: String, val chain: Chain, - private val api: BitcoinApi, + private val api: DirectBitcoinApi, options: UpstreamsConfig.Options, val node: QuorumForLabels.QuorumItem, private val objectMapper: ObjectMapper, callMethods: CallMethods -) : DefaultUpstream(id, options, callMethods), Lifecycle { +) : DefaultUpstream(id, options, callMethods), Lifecycle { companion object { private val log = LoggerFactory.getLogger(BitcoinUpstream::class.java) @@ -28,6 +27,7 @@ class BitcoinUpstream( private val head: Head = createHead() private var validatorSubscription: Disposable? = null + private val data = BitcoinData(api, head) private fun createHead(): Head { return BitcoinRpcHead( @@ -36,11 +36,15 @@ class BitcoinUpstream( ) } + open fun getData(): BitcoinData { + return data + } + override fun getHead(): Head { return head } - override fun getApi(matcher: Selector.Matcher): Mono { + override fun getApi(matcher: Selector.Matcher): Mono { return Mono.just(api) } @@ -56,8 +60,8 @@ class BitcoinUpstream( } override fun castApi(apiType: Class): Upstream { - if (!apiType.isAssignableFrom(BitcoinApi::class.java)) { - throw ClassCastException("Cannot cast ${BitcoinApi::class.java} to $apiType") + if (!apiType.isAssignableFrom(DirectBitcoinApi::class.java)) { + throw ClassCastException("Cannot cast ${DirectBitcoinApi::class.java} to $apiType") } return this as Upstream } @@ -67,6 +71,7 @@ class BitcoinUpstream( if (head is Lifecycle) { runningAny = runningAny || head.isRunning } + runningAny = runningAny || data.isRunning return runningAny } @@ -77,6 +82,8 @@ class BitcoinUpstream( head.start() } } + data.start() + validatorSubscription?.dispose() if (getOptions().disableValidation != null && getOptions().disableValidation!!) { @@ -93,6 +100,7 @@ class BitcoinUpstream( if (head is Lifecycle) { head.stop() } + data.stop() validatorSubscription?.dispose() } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstreamValidator.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstreamValidator.kt index be7f3ef0..6d134a60 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstreamValidator.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinUpstreamValidator.kt @@ -1,10 +1,7 @@ package io.emeraldpay.dshackle.upstream.bitcoin -import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.config.UpstreamsConfig -import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.UpstreamAvailability -import io.infinitape.etherjar.rpc.Commands import org.slf4j.LoggerFactory import org.springframework.scheduling.concurrent.CustomizableThreadFactory import reactor.core.publisher.Flux @@ -14,7 +11,7 @@ import java.time.Duration import java.util.concurrent.Executors class BitcoinUpstreamValidator( - private val api: BitcoinApi, + private val api: DirectBitcoinApi, private val options: UpstreamsConfig.Options ) { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt new file mode 100644 index 00000000..65614c92 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/CachingMempoolData.kt @@ -0,0 +1,71 @@ +package io.emeraldpay.dshackle.upstream.bitcoin + +import io.emeraldpay.dshackle.upstream.Head +import org.slf4j.LoggerFactory +import org.springframework.context.Lifecycle +import reactor.core.Disposable +import reactor.core.publisher.Mono +import java.time.Duration +import java.time.Instant +import java.util.concurrent.atomic.AtomicReference +import java.util.concurrent.locks.ReentrantLock + +open class CachingMempoolData( + private val api: DirectBitcoinApi, + private val head: Head +) : Lifecycle { + + companion object { + private val log = LoggerFactory.getLogger(CachingMempoolData::class.java) + private val TTL = Duration.ofSeconds(15) + } + + private val current = AtomicReference(Container.empty()) + private val updateLock = ReentrantLock() + private var headListener: Disposable? = null + + open fun get(): Mono> { + val value = current.get() + return if (value.since < Instant.now().minus(TTL)) { + updateLock.lock() + fetchFromUpstream() + .timeout(Duration.ofSeconds(3), Mono.empty()) + .doOnNext { + current.set(Container(Instant.now(), it)) + }.doFinally { + updateLock.unlock() + } + } else { + Mono.just(value.value) + } + } + + fun fetchFromUpstream(): Mono> { + return api.executeAndResult(0, "getrawmempool", emptyList(), List::class.java) as Mono> + } + + class Container(val since: Instant, val value: List) { + companion object { + fun empty(): Container { + return Container(Instant.MIN, emptyList()) + } + } + } + + override fun isRunning(): Boolean { + return headListener != null + } + + override fun start() { + headListener?.dispose() + headListener = head.getFlux().doOnNext { + current.set(Container.empty()) + }.subscribe() + } + + override fun stop() { + headListener?.dispose() + headListener = null + } + +} \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinApi.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/DirectBitcoinApi.kt similarity index 94% rename from src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinApi.kt rename to src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/DirectBitcoinApi.kt index 195d7880..c128bd6c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinApi.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/bitcoin/DirectBitcoinApi.kt @@ -14,14 +14,14 @@ import io.infinitape.etherjar.rpc.json.ResponseJson import org.slf4j.LoggerFactory import reactor.core.publisher.Mono -open class BitcoinApi( +open class DirectBitcoinApi( val bitcoinRpcClient: BitcoinRpcClient, val objectMapper: ObjectMapper, val targets: CallMethods ) : UpstreamApi { companion object { - private val log = LoggerFactory.getLogger(BitcoinApi::class.java) + private val log = LoggerFactory.getLogger(DirectBitcoinApi::class.java) } open override fun execute(id: Int, method: String, params: List): Mono { @@ -104,7 +104,4 @@ open class BitcoinApi( return executeAndResult(0, "getrawtransaction", listOf(txid, true), Map::class.java) as Mono> } - open fun getMempool(): Mono> { - return executeAndResult(0, "getrawmempool", emptyList(), List::class.java) as Mono> - } } \ No newline at end of file diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy index ab17ef22..78689d28 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy @@ -9,7 +9,7 @@ import io.emeraldpay.dshackle.upstream.AggregatedUpstream import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstreams -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.DirectBitcoinApi import io.emeraldpay.grpc.Chain import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -201,7 +201,7 @@ class TrackBitcoinAddressSpec extends Specification { def "Get update for a balance"() { setup: - BitcoinApi api = Mock(BitcoinApi) { + DirectBitcoinApi api = Mock(DirectBitcoinApi) { 2 * executeAndResult(0, "listunspent", [], List) >>> [ Mono.just([]), Mono.just([[address: "1K7xkspJg7DDKNwzXgoRSDCUxiFsRegsSK", amount: 0.0123]]) ] diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinTxSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinTxSpec.groovy index 77a21f6b..05e4eafa 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinTxSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinTxSpec.groovy @@ -5,7 +5,10 @@ import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstreams -import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinData +import io.emeraldpay.dshackle.upstream.bitcoin.DirectBitcoinApi +import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream +import io.emeraldpay.dshackle.upstream.bitcoin.CachingMempoolData import io.emeraldpay.grpc.Chain import reactor.core.publisher.Flux import reactor.core.publisher.Mono @@ -20,14 +23,20 @@ class TrackBitcoinTxSpec extends Specification { def "loadMempool() returns not found when not found"() { setup: TrackBitcoinTx track = new TrackBitcoinTx(Stub(Upstreams)) - BitcoinApi api = Mock(BitcoinApi) { - 1 * getMempool() >> Mono.just([ + + CachingMempoolData mempoolAccess = Mock(CachingMempoolData) { + 1 * get() >> Mono.just([ "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9", "d296c6d47335a7f283574b06f1d6303b30ac75631e081ab128346a549ad93350" ]) } + BitcoinUpstream upstream = Mock(BitcoinUpstream) { + _ * getData() >> Mock(BitcoinData) { + _ * getMempool() >> mempoolAccess + } + } when: - def act = track.loadMempool(api, "65ce58db064bd105b14dc76a0bce0df14653cf5263d22d17e78864cf272ee367") + def act = track.loadMempool(upstream, "65ce58db064bd105b14dc76a0bce0df14653cf5263d22d17e78864cf272ee367") then: StepVerifier.create(act) @@ -41,14 +50,19 @@ class TrackBitcoinTxSpec extends Specification { def "loadMempool() returns ok when found"() { setup: TrackBitcoinTx track = new TrackBitcoinTx(Stub(Upstreams)) - BitcoinApi api = Mock(BitcoinApi) { - 1 * getMempool() >> Mono.just([ + CachingMempoolData mempoolAccess = Mock(CachingMempoolData) { + 1 * get() >> Mono.just([ "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9", "d296c6d47335a7f283574b06f1d6303b30ac75631e081ab128346a549ad93350" ]) } + BitcoinUpstream upstream = Mock(BitcoinUpstream) { + _ * getData() >> Mock(BitcoinData) { + _ * getMempool() >> mempoolAccess + } + } when: - def act = track.loadMempool(api, "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9") + def act = track.loadMempool(upstream, "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9") then: StepVerifier.create(act) @@ -63,7 +77,7 @@ class TrackBitcoinTxSpec extends Specification { setup: TrackBitcoinTx track = new TrackBitcoinTx(Stub(Upstreams)) def txid = "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9" - BitcoinApi api = Mock(BitcoinApi) { + DirectBitcoinApi api = Mock(DirectBitcoinApi) { 1 * getTx(txid) >> Mono.just([ txid: txid ]) @@ -84,7 +98,7 @@ class TrackBitcoinTxSpec extends Specification { setup: TrackBitcoinTx track = new TrackBitcoinTx(Stub(Upstreams)) def txid = "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9" - BitcoinApi api = Mock(BitcoinApi) { + DirectBitcoinApi api = Mock(DirectBitcoinApi) { 1 * getTx(txid) >> Mono.just([ txid : txid, blockhash: "0000000000000000000895d1b9d3898700e1deecc3b0e69f439aa77875e6042f", @@ -116,7 +130,7 @@ class TrackBitcoinTxSpec extends Specification { Head head = Mock(Head) { 1 * getFlux() >> next } - Upstream upstream = Mock(Upstream) { + Upstream upstream = Mock(BitcoinUpstream) { 1 * getHead() >> head } def status = new TrackBitcoinTx.TxStatus( @@ -147,7 +161,7 @@ class TrackBitcoinTxSpec extends Specification { Head head = Mock(Head) { 1 * getFlux() >> next } - BitcoinApi api = Mock(BitcoinApi) { + DirectBitcoinApi api = Mock(DirectBitcoinApi) { 3 * getTx(txid) >>> [ Mono.just([ txid: txid @@ -162,7 +176,7 @@ class TrackBitcoinTxSpec extends Specification { ]) ] } - Upstream upstream = Mock(Upstream) { + Upstream upstream = Mock(BitcoinUpstream) { 1 * getHead() >> head _ * getApi(_) >> Mono.just(api) } @@ -183,13 +197,7 @@ class TrackBitcoinTxSpec extends Specification { setup: TrackBitcoinTx track = new TrackBitcoinTx(Stub(Upstreams)) def txid = "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9" - BitcoinApi api = Mock(BitcoinApi) { - 4 * getMempool() >>> [ - Mono.just([]), - Mono.just(["4523c7ac0c5c1e5628f025474529c69cd44d7c641db82e6982f5ffe64527efc9"]), - Mono.just(["4523c7ac0c5c1e5628f025474529c69cd44d7c641db82e6982f5ffe64527efc9", txid]), - Mono.just(["4523c7ac0c5c1e5628f025474529c69cd44d7c641db82e6982f5ffe64527efc9", txid]) //second call when started over - ] + DirectBitcoinApi api = Mock(DirectBitcoinApi) { 1 * getTx(txid) >> Mono.just([ txid: txid ]) @@ -197,9 +205,20 @@ class TrackBitcoinTxSpec extends Specification { Head head = Mock(Head) { _ * getFlux() >> Flux.empty() } - Upstream upstream = Mock(Upstream) { + CachingMempoolData mempoolAccess = Mock(CachingMempoolData) { + 4 * get() >>> [ + Mono.just([]), + Mono.just(["4523c7ac0c5c1e5628f025474529c69cd44d7c641db82e6982f5ffe64527efc9"]), + Mono.just(["4523c7ac0c5c1e5628f025474529c69cd44d7c641db82e6982f5ffe64527efc9", txid]), + Mono.just(["4523c7ac0c5c1e5628f025474529c69cd44d7c641db82e6982f5ffe64527efc9", txid]) //second call when started over + ] + } + BitcoinUpstream upstream = Mock(BitcoinUpstream) { _ * getApi(_) >> Mono.just(api) _ * getHead() >> head + _ * getData() >> Mock(BitcoinData) { + _ * getMempool() >> mempoolAccess + } } when: @@ -220,7 +239,7 @@ class TrackBitcoinTxSpec extends Specification { setup: TrackBitcoinTx track = new TrackBitcoinTx(Stub(Upstreams)) def txid = "69cd44d7c641db82e69824523c7ac0c5c1e5628f025474529cf5ffe64527efc9" - BitcoinApi api = Mock(BitcoinApi) { + DirectBitcoinApi api = Mock(DirectBitcoinApi) { _ * getTx(txid) >> Mono.just([ txid : txid, blockhash: "0000000000000000000895d1b9d3898700e1deecc3b0e69f439aa77875e6042f", @@ -238,7 +257,7 @@ class TrackBitcoinTxSpec extends Specification { Head head = Mock(Head) { _ * getFlux() >> next } - Upstream upstream = Mock(Upstream) { + BitcoinUpstream upstream = Mock(BitcoinUpstream) { _ * getApi(_) >> Mono.just(api) _ * getHead() >> head } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHeadSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHeadSpec.groovy index 0aab0d19..687386dd 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHeadSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinRpcHeadSpec.groovy @@ -59,7 +59,7 @@ class BitcoinRpcHeadSpec extends Specification { } """ - BitcoinApi api = Mock(BitcoinApi) { + DirectBitcoinApi api = Mock(DirectBitcoinApi) { _ * executeAndResult(_, "getbestblockhash", _, String) >>> [ Mono.just(hash1), Mono.just(hash1), Mono.just(hash2) ] diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinApiSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/DirectBitcoinApiSpec.groovy similarity index 96% rename from src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinApiSpec.groovy rename to src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/DirectBitcoinApiSpec.groovy index d63a23cb..146b2443 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/BitcoinApiSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/bitcoin/DirectBitcoinApiSpec.groovy @@ -10,14 +10,14 @@ import spock.lang.Specification import java.time.Duration -class BitcoinApiSpec extends Specification { +class DirectBitcoinApiSpec extends Specification { ClientAndServer mockServer - BitcoinApi api + DirectBitcoinApi api def setup() { mockServer = ClientAndServer.startClientAndServer(18332); - api = new BitcoinApi( + api = new DirectBitcoinApi( new BitcoinRpcClient("localhost:18332", null), TestingCommons.objectMapper(), new DefaultBitcoinMethods(TestingCommons.objectMapper()) )