problem: balance subscription misses recheck addresses for xpub
This commit is contained in:
@@ -131,19 +131,21 @@ class TrackBitcoinAddress(
|
|||||||
val chain = Chain.byId(request.asset.chainValue)
|
val chain = Chain.byId(request.asset.chainValue)
|
||||||
val upstream = multistreamHolder.getUpstream(chain)?.cast(BitcoinMultistream::class.java)
|
val upstream = multistreamHolder.getUpstream(chain)?.cast(BitcoinMultistream::class.java)
|
||||||
?: return Flux.error(SilentException.UnsupportedBlockchain(request.asset.chainValue))
|
?: return Flux.error(SilentException.UnsupportedBlockchain(request.asset.chainValue))
|
||||||
val addresses = allAddresses(upstream, request) ?: return Flux.error(SilentException("Unsupported address"))
|
val addresses = allAddresses(upstream, request).cache()
|
||||||
val initial = requestBalances(chain, upstream, addresses, request.includeUtxo)
|
|
||||||
val following = upstream.getHead().getFlux()
|
val following = upstream.getHead().getFlux()
|
||||||
.flatMap { block ->
|
.flatMap { block ->
|
||||||
requestBalances(chain, upstream, addresses, request.includeUtxo)
|
requestBalances(chain, upstream, Flux.from(addresses), request.includeUtxo)
|
||||||
}
|
}
|
||||||
val last = HashMap<Address, BigInteger>()
|
val last = HashMap<String, BigInteger>()
|
||||||
val result = Flux.merge(initial, following)
|
val result = following
|
||||||
.filter { curr ->
|
.filter { curr ->
|
||||||
val prev = last[curr.address]
|
val prev = last[curr.address.address]
|
||||||
val updated = prev == null || curr.balance != prev
|
//TODO utxo can change without changing balance
|
||||||
last[curr.address] = curr.balance
|
val changed = prev == null || curr.balance != prev
|
||||||
updated
|
if (changed) {
|
||||||
|
last[curr.address.address] = curr.balance
|
||||||
|
}
|
||||||
|
changed
|
||||||
}
|
}
|
||||||
|
|
||||||
return result.map(this@TrackBitcoinAddress::buildResponse)
|
return result.map(this@TrackBitcoinAddress::buildResponse)
|
||||||
|
|||||||
@@ -235,7 +235,12 @@ class TrackBitcoinAddressSpec extends Specification {
|
|||||||
|
|
||||||
def blocks = TopicProcessor.create()
|
def blocks = TopicProcessor.create()
|
||||||
Head head = Mock(Head) {
|
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
|
def upstream = null
|
||||||
upstream = Mock(BitcoinMultistream) {
|
upstream = Mock(BitcoinMultistream) {
|
||||||
|
|||||||
Reference in New Issue
Block a user