From 045c2fb3faee13596cc8602ada5101c0b084b0b7 Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Wed, 31 Jul 2019 23:50:38 -0400 Subject: [PATCH] problem: not subscribing to tx updates --- .../io/emeraldpay/dshackle/rpc/TrackTx.kt | 212 ++++++++++------- .../dshackle/upstream/EthereumApi.kt | 6 +- .../dshackle/rpc/TrackTxSpec.groovy | 217 ++++++++++++++++++ 3 files changed, 347 insertions(+), 88 deletions(-) create mode 100644 src/test/groovy/io/emeraldpay/dshackle/rpc/TrackTxSpec.groovy diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackTx.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackTx.kt index 860c8004..f1ad51c9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackTx.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackTx.kt @@ -4,7 +4,6 @@ import com.google.protobuf.ByteString import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.Common import io.emeraldpay.dshackle.upstream.AvailableChains -import io.emeraldpay.dshackle.upstream.ConfiguredUpstreams import io.emeraldpay.dshackle.upstream.Upstreams import io.emeraldpay.grpc.Chain import io.infinitape.etherjar.domain.BlockHash @@ -17,11 +16,14 @@ import reactor.core.publisher.Flux import reactor.core.publisher.Mono import reactor.core.publisher.TopicProcessor import reactor.core.publisher.toFlux +import reactor.core.scheduler.Scheduler +import reactor.util.function.Tuples import java.lang.Exception import java.math.BigInteger import java.time.Duration import java.time.Instant import java.util.concurrent.ConcurrentLinkedQueue +import java.util.concurrent.atomic.AtomicLong import javax.annotation.PostConstruct import kotlin.math.max import kotlin.math.min @@ -29,13 +31,15 @@ import kotlin.math.min @Service class TrackTx( @Autowired private val upstreams: Upstreams, - @Autowired private val availableChains: AvailableChains + @Autowired private val availableChains: AvailableChains, + @Autowired private val upstreamScheduler: Scheduler ) { private val ZERO_BLOCK = BlockHash.from("0x0000000000000000000000000000000000000000000000000000000000000000") private val log = LoggerFactory.getLogger(TrackTx::class.java) private val clients = HashMap>() + private val seq = AtomicLong(0) @PostConstruct fun init() { @@ -53,40 +57,42 @@ class TrackTx( fun add(requestMono: Mono): Flux { return requestMono.map { request -> - val sender = TopicProcessor.create() + val bus = TopicProcessor.create() TrackTx.TrackedTx( Chain.byId(request.chainValue), - sender, + bus, Instant.now(), TransactionId.from(request.txId), - min(max(1, request.confirmationLimit), 100) + min(max(1, request.confirmationLimit), 100), + seq.incrementAndGet() ) }.filter { clients.containsKey(it.chain) - }.map { tx -> - currentList(tx.chain)!!.let { list -> - list.add(tx) - tx.stream.doOnError { - list.remove(tx) - tx.stream.dispose() + }.flatMapMany { tx -> + val current = checkForUpdate(tx).doOnNext{ + currentList(tx.chain)?.add(tx) + }.map(this::asProto) + + val next = Flux.from(tx.bus) + Flux.merge(current, next).doFinally { + currentList(tx.chain)?.removeIf { x -> x.id == tx.id } + }.doOnNext { txp -> + if (txp.confirmations >= tx.maxConfirmations) { + tx.bus.onComplete() } } - tx - }.map { tx -> - verify(tx) - notify(tx) - tx - }.flatMapMany { tx -> - tx.stream } } private fun verifyAll(chain: Chain) { currentList(chain)!! .toFlux() - .filter(this::verify) - .subscribe { - notify(it) + .parallel(8).runOn(upstreamScheduler) + .flatMap { checkForUpdate(it) } + .sequential() + .map { Tuples.of(it.bus, asProto(it)) } + .subscribe { t -> + notify(t.t1, t.t2) } } @@ -106,113 +112,147 @@ class TrackTx( mined = false ) } + }.doOnError { t -> + log.warn("Failed to update weight", t) } } - private fun verify(tx: TrackedTx): Boolean { - val found = tx.status.found - val mined = tx.status.mined - val upstream = upstreams.getUpstream(tx.chain) ?: return false + private fun checkForUpdate(tx: TrackedTx): Mono { + val upstream = upstreams.getUpstream(tx.chain) ?: return Mono.error(Exception("Unsupported blockchain: ${tx.chain}")) val execution = upstream.getApi() .executeAndConvert(Commands.eth().getTransaction(tx.txid)) - val update = execution.flatMap { - if (it.blockNumber != null - && it.blockHash != null && it.blockHash != ZERO_BLOCK) { - tx.withStatus( + return execution.flatMap { + if (it.blockNumber != null && it.blockHash != null && it.blockHash != ZERO_BLOCK) { + val updated = tx.withStatus( blockHash = it.blockHash, height = it.blockNumber, found = true, mined = true, - confirmation = 1 + confirmations = 1 ) - return@flatMap upstream.getHead().getHead().map { head -> - tx.withStatus( - confirmation = head.number - tx.status.height!! + 1 - ) + upstream.getHead().getHead().map { head -> + if (updated.status.height == null || head.number < updated.status.height) { + updated + } else { + updated.withStatus( + confirmations = head.number - updated.status.height + 1 + ) + } }.flatMap(this::loadWeight) } else { - tx.withStatus( - found = true, - mined = false - ) + Mono.just(tx.withStatus( + found = true, + mined = false + )) } - return@flatMap Mono.just(tx) - }.block() - if (update == null) { - tx.withStatus( - found = false, - mined = false - ) + }.switchIfEmpty(Mono.just(tx.withStatus(found = false))).filter { current -> + current.status != tx.status } - if (!found) { - return tx.status.found != found - } - if (!mined) { - return tx.status.mined != mined - } - return true } - private fun notify(tx: TrackedTx) { - val client = tx.stream + private fun asProto(tx: TrackedTx): BlockchainOuterClass.TxStatus { val data = BlockchainOuterClass.TxStatus.newBuilder() .setTxId(tx.txid.toHex()) - .setConfirmations(tx.status.confirmation.toInt()) - .setMined(tx.status.mined) - .setBroadcasted(tx.status.found) + .setConfirmations(tx.status.confirmations.toInt()) - if (tx.status.mined) { - data.setBlock( - Common.BlockInfo.newBuilder() - .setBlockId(tx.status.blockHash!!.toHex().substring(2)) - .setTimestamp(tx.status.blockTime!!.toEpochMilli()) - .setWeight(ByteString.copyFrom(tx.status.blockTotalDifficulty!!.toByteArray())) - .setHeight(tx.status.height!!) - .setTimestamp(tx.status.blockTime!!.toEpochMilli()) - ) + if (tx.status.found != null) { + data.broadcasted = tx.status.found } - client.onNext(data.build()) + if (tx.status.mined != null) { + data.mined = tx.status.mined + if (tx.status.mined) { + data.setBlock( + Common.BlockInfo.newBuilder() + .setBlockId(tx.status.blockHash!!.toHex().substring(2)) + .setTimestamp(tx.status.blockTime!!.toEpochMilli()) + .setWeight(ByteString.copyFrom(tx.status.blockTotalDifficulty!!.toByteArray())) + .setHeight(tx.status.height!!) + .setTimestamp(tx.status.blockTime!!.toEpochMilli()) + ) + } + } + return data.build() + } + + private fun notify(client: TopicProcessor, data: BlockchainOuterClass.TxStatus) { + client.onNext(data) } class TrackedTx(val chain: Chain, - val stream: TopicProcessor, + val bus: TopicProcessor, val since: Instant, val txid: TransactionId, val maxConfirmations: Int, - var status: TxStatus = TxStatus()) { + val id: Long, + val status: TxStatus = TxStatus()) { - fun withStatus(found: Boolean = this.status.found, + fun withStatus(found: Boolean? = this.status.found, height: Long? = this.status.height, - mined: Boolean = this.status.mined, + mined: Boolean? = this.status.mined, blockHash: BlockHash? = this.status.blockHash, blockTime: Instant? = this.status.blockTime, blockTotalDifficulty: BigInteger? = this.status.blockTotalDifficulty, - confirmation: Long = this.status.confirmation): TrackedTx { - this.status = this.status.copy(found, height, mined, blockHash, blockTime, blockTotalDifficulty, confirmation) - return this - } + confirmations: Long = this.status.confirmations) + = TrackedTx( + chain, bus, since, txid, maxConfirmations, id, + this.status.copy(found, height, mined, blockHash, blockTime, blockTotalDifficulty, confirmations) + ) + + fun withCleanStatus() + = TrackedTx( + chain, bus, since, txid, maxConfirmations, id, this.status.clean() + ) fun shouldClose(): Boolean { - return maxConfirmations <= this.status.confirmation + return maxConfirmations <= this.status.confirmations || since.isBefore(Instant.now().minus(Duration.ofHours(1))) } } - class TxStatus(var found: Boolean = false, - var height: Long? = null, - var mined: Boolean = false, - var blockHash: BlockHash? = null, - var blockTime: Instant? = null, - var blockTotalDifficulty: BigInteger? = null, - var confirmation: Long = 0) { - fun copy(found: Boolean = this.found, + class TxStatus(val found: Boolean? = null, + val height: Long? = null, + val mined: Boolean? = null, + val blockHash: BlockHash? = null, + val blockTime: Instant? = null, + val blockTotalDifficulty: BigInteger? = null, + val confirmations: Long = 0) { + + fun copy(found: Boolean? = this.found, height: Long? = this.height, - mined: Boolean = this.mined, + mined: Boolean? = this.mined, blockHash: BlockHash? = this.blockHash, blockTime: Instant? = this.blockTime, blockTotalDifficulty: BigInteger? = this.blockTotalDifficulty, - confirmation: Long = this.confirmation) + confirmation: Long = this.confirmations) = TxStatus(found, height, mined, blockHash, blockTime, blockTotalDifficulty, confirmation) + + fun clean() = TxStatus(false, null, false, null, null, null, 0) + + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (javaClass != other?.javaClass) return false + + other as TxStatus + + if (found != other.found) return false + if (height != other.height) return false + if (mined != other.mined) return false + if (blockHash != other.blockHash) return false + if (blockTime != other.blockTime) return false + if (blockTotalDifficulty != other.blockTotalDifficulty) return false + if (confirmations != other.confirmations) return false + + return true + } + + override fun hashCode(): Int { + var result = found.hashCode() + result = 31 * result + (height?.hashCode() ?: 0) + result = 31 * result + (blockHash?.hashCode() ?: 0) + return result + } + + } } \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/EthereumApi.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/EthereumApi.kt index df8edc32..d3468629 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/EthereumApi.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/EthereumApi.kt @@ -68,8 +68,10 @@ open class EthereumApi( open fun executeAndConvert(rpcCall: RpcCall): Mono { return execute(0, rpcCall.method, rpcCall.params as List) - .map { - jacksonRpcConverter.fromJson(it.inputStream(), rpcCall.jsonType, Int::class.java) + .flatMap { + val jsonValue: JS? = jacksonRpcConverter.fromJson(it.inputStream(), rpcCall.jsonType, Int::class.java); + if (jsonValue == null) Mono.empty() + else Mono.just(jsonValue) }.map { rpcCall.converter.apply(it) } diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackTxSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackTxSpec.groovy new file mode 100644 index 00000000..03b62eda --- /dev/null +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackTxSpec.groovy @@ -0,0 +1,217 @@ +package io.emeraldpay.dshackle.rpc + +import com.google.protobuf.ByteString +import io.emeraldpay.api.proto.BlockchainOuterClass +import io.emeraldpay.api.proto.Common +import io.emeraldpay.dshackle.test.EthereumApiMock +import io.emeraldpay.dshackle.test.TestingCommons +import io.emeraldpay.dshackle.upstream.AggregatedUpstreams +import io.emeraldpay.dshackle.upstream.AvailableChains +import io.emeraldpay.dshackle.upstream.EthereumHead +import io.emeraldpay.dshackle.upstream.Upstreams +import io.emeraldpay.grpc.Chain +import io.infinitape.etherjar.domain.BlockHash +import io.infinitape.etherjar.domain.TransactionId +import io.infinitape.etherjar.rpc.RpcClient +import io.infinitape.etherjar.rpc.json.BlockJson +import io.infinitape.etherjar.rpc.json.TransactionJson +import reactor.core.publisher.Mono +import reactor.core.publisher.TopicProcessor +import reactor.core.scheduler.Schedulers +import reactor.test.StepVerifier +import spock.lang.Specification + +import java.time.Duration + +class TrackTxSpec extends Specification { + + AvailableChains availableChains = new AvailableChains() + Upstreams upstreams + TrackTx trackTx + + def chain = Common.ChainRef.CHAIN_ETHEREUM + def txId = "0xba61ce4672751fd6086a9ac2b55547a5555af17535b6c0334ede2ecb6d64070a" + + def setup() { + upstreams = Mock(Upstreams) + trackTx = new TrackTx(upstreams, availableChains, Schedulers.immediate()) + } + + def start() { + trackTx.init() + availableChains.add(Chain.ETHEREUM) + availableChains.add(Chain.TESTNET_KOVAN) + } + + def "Gives details for an old transaction"() { + setup: + def req = BlockchainOuterClass.TxStatusRequest.newBuilder() + .setChain(chain) + .setConfirmationLimit(6) + .setTxId(txId) + .build() + + def blockJson = new BlockJson().with { + it.hash = BlockHash.from("0xa0e65cbc1b52a8ca60562112c6060552d882f16f34a9dba2ccdc05c0a6a27c22") + it.timestamp = new Date(156400000000) + it.number = 100 + it.totalDifficulty = BigInteger.valueOf(500) + it + } + + def blockHeadJson = new BlockJson().with { + it.hash = BlockHash.from("0xa0e65cbc1b52a8ca60562112c6060552d882f16f34a9dba2ccdc05c0a6a27c22") + it.timestamp = new Date(156400200000) + it.number = 108 + it.totalDifficulty = BigInteger.valueOf(800) + it + } + + + def txJson = new TransactionJson().with { + it.hash = TransactionId.from("0xba61ce4672751fd6086a9ac2b55547a5555af17535b6c0334ede2ecb6d64070a") + it.blockHash = blockJson.hash + it.blockNumber = blockJson.number + it.nonce = 1 + it + } + + def exp1 = BlockchainOuterClass.TxStatus.newBuilder() + .setTxId(txId) + .setBroadcasted(true) + .setMined(true) + .setConfirmations(8 + 1) + .setBlock( + Common.BlockInfo.newBuilder() + .setHeight(blockJson.number) + .setWeight(ByteString.copyFrom(blockJson.totalDifficulty.toByteArray())) + .setBlockId(blockJson.hash.toHex().substring(2)) + .setTimestamp(blockJson.timestamp.getTime()) + ).build() + + + def upstreamMock = Mock(AggregatedUpstreams) + def blocksBus = TopicProcessor.create() + def headMock = Mock(EthereumHead) + + def apiMock = new EthereumApiMock(Mock(RpcClient), TestingCommons.objectMapper(), Chain.ETHEREUM) + apiMock.answer("eth_getTransactionByHash", [txId], txJson) + apiMock.answer("eth_getBlockByHash", [blockJson.hash.toHex(), false], blockJson) + + _ * upstreams.getUpstream(Chain.ETHEREUM) >> upstreamMock + _ * upstreamMock.getApi() >> apiMock + _ * upstreamMock.getHead() >> headMock + _ * headMock.getFlux() >> blocksBus + _ * headMock.getHead() >> Mono.just(blockHeadJson) + start() + + when: + def flux = trackTx.add(Mono.just(req)) + then: + StepVerifier.create(flux) + .expectNext(exp1) + .expectComplete() + .verify(Duration.ofSeconds(4)) + } + + def "Starts to follow new transaction"() { + setup: + def req = BlockchainOuterClass.TxStatusRequest.newBuilder() + .setChain(chain) + .setConfirmationLimit(4) + .setTxId(txId) + .build() + + List> blocks = (0..9).collect { i -> + return new BlockJson().with { + it.hash = BlockHash.from("0xa0e65cbc1b52a8ca60562112c6060552d882f16f34a9dba2ccdc05c0a6a2000${i}") + it.timestamp = new Date(156400000000 + i * 10000) + it.setNumber(100L + i.longValue()) + it.totalDifficulty = BigInteger.valueOf(500 + i) + it + } + } + + def txJsonBroadcasted = new TransactionJson().with { + it.hash = TransactionId.from("0xba61ce4672751fd6086a9ac2b55547a5555af17535b6c0334ede2ecb6d64070a") + it.blockHash = null + it.blockNumber = null + it.nonce = 1 + it + } + + def txJsonMined = new TransactionJson().with { + it.hash = TransactionId.from("0xba61ce4672751fd6086a9ac2b55547a5555af17535b6c0334ede2ecb6d64070a") + it.blockHash = blocks[2].hash + it.blockNumber = blocks[2].number + it.nonce = 1 + it + } + + def exp1 = BlockchainOuterClass.TxStatus.newBuilder() + .setTxId(txId) + .setBroadcasted(false) + .setMined(false) + .setConfirmations(0) + + def exp2 = BlockchainOuterClass.TxStatus.newBuilder() + .setTxId(txId) + .setBroadcasted(true) + .setMined(true) + .setBlock( + Common.BlockInfo.newBuilder() + .setHeight(blocks[2].number) + .setWeight(ByteString.copyFrom(blocks[2].totalDifficulty.toByteArray())) + .setBlockId(blocks[2].hash.toHex().substring(2)) + .setTimestamp(blocks[2].timestamp.getTime()) + ) + + + def upstreamMock = Mock(AggregatedUpstreams) + def blocksBus = TopicProcessor.create() + def headMock = Mock(EthereumHead) + + def apiMock = new EthereumApiMock(Mock(RpcClient), TestingCommons.objectMapper(), Chain.ETHEREUM) + apiMock.answerOnce("eth_getTransactionByHash", [txId], null) + apiMock.answerOnce("eth_getTransactionByHash", [txId], txJsonBroadcasted) + apiMock.answer("eth_getTransactionByHash", [txId], txJsonMined) + blocks.forEach { block -> + apiMock.answer("eth_getBlockByHash", [block.hash.toHex(), false], block) + } + + def headBlock = blocks[0] + + _ * upstreams.getUpstream(Chain.ETHEREUM) >> upstreamMock + _ * upstreamMock.getApi() >> apiMock + _ * upstreamMock.getHead() >> headMock + _ * headMock.getFlux() >> blocksBus + _ * headMock.getHead() >> { return Mono.just(headBlock) } + start() + + def nextBlock = { int i -> + return { + println("block $i"); + headBlock = blocks[i]; + blocksBus.onNext(blocks[i]) + } as Runnable + } + + when: + def flux = trackTx.add(Mono.just(req)) + then: + StepVerifier.create(flux) + .expectNext(exp1.build()) + .then(nextBlock(1)) + .expectNext(exp1.setBroadcasted(true).build()) + .then(nextBlock(2)) + .expectNext(exp2.setConfirmations(1).build()) + .then(nextBlock(3)) + .expectNext(exp2.setConfirmations(2).build()) + .then(nextBlock(4)) + .expectNext(exp2.setConfirmations(3).build()) + .then(nextBlock(5)) + .expectNext(exp2.setConfirmations(4).build()) + .expectComplete() + .verify(Duration.ofSeconds(4)) + } +}