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) {