From 275a55c4eecfe0361bcc77161bf37073d47e6c9c Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Tue, 15 Sep 2020 23:32:23 -0400 Subject: [PATCH] problem: balance subscription misses recheck addresses for xpub --- .../dshackle/rpc/TrackBitcoinAddress.kt | 20 ++++++++++--------- .../rpc/TrackBitcoinAddressSpec.groovy | 7 ++++++- 2 files changed, 17 insertions(+), 10 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt index 4b914022..43167e8b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/TrackBitcoinAddress.kt @@ -131,19 +131,21 @@ class TrackBitcoinAddress( val chain = Chain.byId(request.asset.chainValue) val upstream = multistreamHolder.getUpstream(chain)?.cast(BitcoinMultistream::class.java) ?: return Flux.error(SilentException.UnsupportedBlockchain(request.asset.chainValue)) - val addresses = allAddresses(upstream, request) ?: return Flux.error(SilentException("Unsupported address")) - val initial = requestBalances(chain, upstream, addresses, request.includeUtxo) + val addresses = allAddresses(upstream, request).cache() val following = upstream.getHead().getFlux() .flatMap { block -> - requestBalances(chain, upstream, addresses, request.includeUtxo) + requestBalances(chain, upstream, Flux.from(addresses), request.includeUtxo) } - val last = HashMap() - val result = Flux.merge(initial, following) + val last = HashMap() + val result = following .filter { curr -> - val prev = last[curr.address] - val updated = prev == null || curr.balance != prev - last[curr.address] = curr.balance - updated + val prev = last[curr.address.address] + //TODO utxo can change without changing balance + val changed = prev == null || curr.balance != prev + if (changed) { + last[curr.address.address] = curr.balance + } + changed } return result.map(this@TrackBitcoinAddress::buildResponse) diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy index 34062d3e..be866781 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/TrackBitcoinAddressSpec.groovy @@ -235,7 +235,12 @@ class TrackBitcoinAddressSpec extends Specification { def blocks = TopicProcessor.create() Head head = Mock(Head) { - 1 * getFlux() >> Flux.from(blocks) + 1 * getFlux() >> Flux.concat( + Flux.just( + new BlockContainer(0L, BlockId.from(hash1), BigInteger.ZERO, Instant.now(), false, null, null, []) + ), + Flux.from(blocks) + ) } def upstream = null upstream = Mock(BitcoinMultistream) {