solution: update to latest etherjar with updated blockjson model
This commit is contained in:
@@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.cache
|
||||
import io.infinitape.etherjar.domain.BlockHash
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import reactor.core.publisher.Mono
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.ConcurrentLinkedQueue
|
||||
@@ -26,14 +27,14 @@ class BlocksMemCache(
|
||||
val maxSize: Int = 64
|
||||
) {
|
||||
|
||||
private val mapping = ConcurrentHashMap<BlockHash, BlockJson<TransactionId>>()
|
||||
private val mapping = ConcurrentHashMap<BlockHash, BlockJson<TransactionRefJson>>()
|
||||
private val queue = ConcurrentLinkedQueue<BlockHash>()
|
||||
|
||||
fun get(hash: BlockHash): Mono<BlockJson<TransactionId>> {
|
||||
fun get(hash: BlockHash): Mono<BlockJson<TransactionRefJson>> {
|
||||
return Mono.justOrEmpty(mapping[hash])
|
||||
}
|
||||
|
||||
fun add(block: BlockJson<TransactionId>) {
|
||||
fun add(block: BlockJson<TransactionRefJson>) {
|
||||
mapping.put(block.hash, block)
|
||||
queue.add(block.hash)
|
||||
|
||||
|
||||
@@ -19,13 +19,14 @@ import io.emeraldpay.dshackle.upstream.Head
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
|
||||
open class AlwaysQuorum: CallQuorum {
|
||||
|
||||
private var resolved = false
|
||||
private var result: ByteArray? = null
|
||||
|
||||
override fun init(head: Head<BlockJson<TransactionId>>) {
|
||||
override fun init(head: Head<BlockJson<TransactionRefJson>>) {
|
||||
}
|
||||
|
||||
override fun isResolved(): Boolean {
|
||||
|
||||
@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.JacksonRpcConverter
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
|
||||
open class BroadcastQuorum(
|
||||
jacksonRpcConverter: JacksonRpcConverter,
|
||||
@@ -30,7 +31,7 @@ open class BroadcastQuorum(
|
||||
private var txid: String? = null
|
||||
private var calls = 0
|
||||
|
||||
override fun init(head: Head<BlockJson<TransactionId>>) {
|
||||
override fun init(head: Head<BlockJson<TransactionRefJson>>) {
|
||||
}
|
||||
|
||||
override fun isResolved(): Boolean {
|
||||
|
||||
@@ -19,13 +19,14 @@ import io.emeraldpay.dshackle.upstream.Head
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import reactor.util.function.Tuple2
|
||||
import java.util.function.BiFunction
|
||||
import java.util.function.Predicate
|
||||
|
||||
interface CallQuorum {
|
||||
|
||||
fun init(head: Head<BlockJson<TransactionId>>)
|
||||
fun init(head: Head<BlockJson<TransactionRefJson>>)
|
||||
|
||||
fun isResolved(): Boolean
|
||||
fun record(response: ByteArray, upstream: Upstream): Boolean
|
||||
|
||||
@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.JacksonRpcConverter
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
|
||||
open class NonEmptyQuorum(
|
||||
jacksonRpcConverter: JacksonRpcConverter,
|
||||
@@ -29,7 +30,7 @@ open class NonEmptyQuorum(
|
||||
private var result: ByteArray? = null
|
||||
private var tries: Int = 0
|
||||
|
||||
override fun init(head: Head<BlockJson<TransactionId>>) {
|
||||
override fun init(head: Head<BlockJson<TransactionRefJson>>) {
|
||||
}
|
||||
|
||||
override fun isResolved(): Boolean {
|
||||
|
||||
@@ -21,6 +21,7 @@ import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.hex.HexQuantity
|
||||
import io.infinitape.etherjar.rpc.JacksonRpcConverter
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import java.util.concurrent.locks.ReentrantLock
|
||||
import kotlin.concurrent.withLock
|
||||
|
||||
@@ -35,7 +36,7 @@ open class NonceQuorum(
|
||||
private var receivedTimes = 0
|
||||
private var errors = 0
|
||||
|
||||
override fun init(head: Head<BlockJson<TransactionId>>) {
|
||||
override fun init(head: Head<BlockJson<TransactionRefJson>>) {
|
||||
}
|
||||
|
||||
override fun isResolved(): Boolean {
|
||||
|
||||
@@ -19,13 +19,14 @@ import io.emeraldpay.dshackle.upstream.Head
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import java.util.concurrent.atomic.AtomicReference
|
||||
|
||||
class NotLaggingQuorum(val maxLag: Long = 0): CallQuorum {
|
||||
|
||||
private val result: AtomicReference<ByteArray> = AtomicReference()
|
||||
|
||||
override fun init(head: Head<BlockJson<TransactionId>>) {
|
||||
override fun init(head: Head<BlockJson<TransactionRefJson>>) {
|
||||
}
|
||||
|
||||
override fun isResolved(): Boolean {
|
||||
|
||||
@@ -21,15 +21,16 @@ import io.infinitape.etherjar.domain.BlockHash
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.Commands
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import reactor.core.publisher.Mono
|
||||
import reactor.retry.Repeat
|
||||
import java.time.Duration
|
||||
|
||||
class BlockApiReader(
|
||||
val upstream: Upstream
|
||||
): Reader<BlockHash, BlockJson<TransactionId>> {
|
||||
): Reader<BlockHash, BlockJson<TransactionRefJson>> {
|
||||
|
||||
override fun read(key: BlockHash): Mono<BlockJson<TransactionId>> {
|
||||
override fun read(key: BlockHash): Mono<BlockJson<TransactionRefJson>> {
|
||||
return Mono.just(key)
|
||||
.flatMap {
|
||||
upstream.getApi(Selector.empty)
|
||||
|
||||
@@ -19,13 +19,14 @@ import io.emeraldpay.dshackle.cache.BlocksMemCache
|
||||
import io.infinitape.etherjar.domain.BlockHash
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import reactor.core.publisher.Mono
|
||||
|
||||
class BlockCacheReader(
|
||||
val cache: BlocksMemCache
|
||||
): Reader<BlockHash, BlockJson<TransactionId>> {
|
||||
): Reader<BlockHash, BlockJson<TransactionRefJson>> {
|
||||
|
||||
override fun read(key: BlockHash): Mono<BlockJson<TransactionId>> {
|
||||
override fun read(key: BlockHash): Mono<BlockJson<TransactionRefJson>> {
|
||||
return cache.get(key)
|
||||
}
|
||||
}
|
||||
@@ -22,6 +22,7 @@ import io.emeraldpay.dshackle.upstream.Upstreams
|
||||
import io.emeraldpay.grpc.Chain
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.stereotype.Service
|
||||
@@ -47,7 +48,7 @@ class StreamHead(
|
||||
}
|
||||
}
|
||||
|
||||
fun asProto(chain: Chain, block: BlockJson<TransactionId>): BlockchainOuterClass.ChainHead {
|
||||
fun asProto(chain: Chain, block: BlockJson<TransactionRefJson>): BlockchainOuterClass.ChainHead {
|
||||
return BlockchainOuterClass.ChainHead.newBuilder()
|
||||
.setChainValue(chain.id)
|
||||
.setHeight(block.number)
|
||||
|
||||
@@ -27,6 +27,7 @@ import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.Commands
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.scheduling.annotation.Scheduled
|
||||
@@ -191,7 +192,7 @@ class TrackTx(
|
||||
}
|
||||
}
|
||||
|
||||
fun setBlockDetails(tx: TxDetails, block: BlockJson<TransactionId>): TxDetails {
|
||||
fun setBlockDetails(tx: TxDetails, block: BlockJson<TransactionRefJson>): TxDetails {
|
||||
return if (block.number != null && block.totalDifficulty != null) {
|
||||
tx.withStatus(
|
||||
blockTotalDifficulty = block.totalDifficulty,
|
||||
|
||||
@@ -26,6 +26,7 @@ import io.emeraldpay.dshackle.upstream.ethereum.EthereumHead
|
||||
import io.infinitape.etherjar.domain.BlockHash
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.reactivestreams.Publisher
|
||||
import org.springframework.context.Lifecycle
|
||||
import reactor.core.Disposable
|
||||
@@ -43,7 +44,7 @@ abstract class AggregatedUpstream(
|
||||
|
||||
private val blocksCache = BlocksMemCache()
|
||||
private var cacheSubscription: Disposable? = null
|
||||
private val blockReader: Reader<BlockHash, BlockJson<TransactionId>> = CompoundReader(
|
||||
private val blockReader: Reader<BlockHash, BlockJson<TransactionRefJson>> = CompoundReader(
|
||||
listOf(BlockCacheReader(blocksCache))
|
||||
)
|
||||
var cache: CachingEthereumApi = CachingEthereumApi.empty()
|
||||
|
||||
@@ -26,12 +26,13 @@ import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.hex.HexQuantity
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.ResponseJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import reactor.core.publisher.Mono
|
||||
import java.util.function.Function
|
||||
|
||||
open class CachingEthereumApi(
|
||||
private val objectMapper: ObjectMapper,
|
||||
private val cache: Reader<BlockHash, BlockJson<TransactionId>>,
|
||||
private val cache: Reader<BlockHash, BlockJson<TransactionRefJson>>,
|
||||
private val head: EthereumHead
|
||||
): EthereumApi(objectMapper) {
|
||||
|
||||
|
||||
@@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.upstream
|
||||
import io.emeraldpay.dshackle.upstream.ethereum.EthereumHead
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.context.Lifecycle
|
||||
import reactor.core.Disposable
|
||||
@@ -57,7 +58,7 @@ class HeadLagObserver (
|
||||
}
|
||||
}
|
||||
|
||||
fun probeFollowers(top: BlockJson<TransactionId>): Flux<Tuple2<Long, Upstream>> {
|
||||
fun probeFollowers(top: BlockJson<TransactionRefJson>): Flux<Tuple2<Long, Upstream>> {
|
||||
return followers.toFlux()
|
||||
.parallel(followers.size)
|
||||
.flatMap { mapLagging(top, it, getCurrentBlocks(it)) }
|
||||
@@ -65,12 +66,12 @@ class HeadLagObserver (
|
||||
.onErrorContinue { t, _ -> log.warn("Failed to update lagging distance", t) }
|
||||
}
|
||||
|
||||
fun getCurrentBlocks(up: Upstream): Flux<BlockJson<TransactionId>> {
|
||||
fun getCurrentBlocks(up: Upstream): Flux<BlockJson<TransactionRefJson>> {
|
||||
val head = up.getHead()
|
||||
return head.getFlux().take(Duration.ofSeconds(1))
|
||||
}
|
||||
|
||||
fun mapLagging(top: BlockJson<TransactionId>, up: Upstream, blocks: Flux<BlockJson<TransactionId>>): Flux<Tuple2<Long, Upstream>> {
|
||||
fun mapLagging(top: BlockJson<TransactionRefJson>, up: Upstream, blocks: Flux<BlockJson<TransactionRefJson>>): Flux<Tuple2<Long, Upstream>> {
|
||||
return blocks
|
||||
.map { extractDistance(top, it) }
|
||||
.takeUntil{ lag -> lag <= 0L }
|
||||
@@ -80,7 +81,7 @@ class HeadLagObserver (
|
||||
}
|
||||
}
|
||||
|
||||
fun extractDistance(top: BlockJson<TransactionId>, curr: BlockJson<TransactionId>): Long {
|
||||
fun extractDistance(top: BlockJson<TransactionRefJson>, curr: BlockJson<TransactionRefJson>): Long {
|
||||
return when {
|
||||
curr.number > top.number -> if (curr.totalDifficulty >= top.totalDifficulty) 0 else forkDistance(top, curr)
|
||||
curr.number == top.number -> if (curr.totalDifficulty == top.totalDifficulty) 0 else forkDistance(top, curr)
|
||||
@@ -88,7 +89,7 @@ class HeadLagObserver (
|
||||
}
|
||||
}
|
||||
|
||||
fun forkDistance(top: BlockJson<TransactionId>, curr: BlockJson<TransactionId>): Long {
|
||||
fun forkDistance(top: BlockJson<TransactionRefJson>, curr: BlockJson<TransactionRefJson>): Long {
|
||||
//TODO look for common ancestor? though it may be a corruption
|
||||
return 6
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package io.emeraldpay.dshackle.upstream.ethereum
|
||||
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.slf4j.LoggerFactory
|
||||
import reactor.core.Disposable
|
||||
import reactor.core.publisher.Flux
|
||||
@@ -12,10 +13,10 @@ import java.util.concurrent.atomic.AtomicReference
|
||||
open class DefaultEthereumHead: EthereumHead {
|
||||
|
||||
private val log = LoggerFactory.getLogger(DefaultEthereumHead::class.java)
|
||||
private val head = AtomicReference<BlockJson<TransactionId>>(null)
|
||||
private val stream: TopicProcessor<BlockJson<TransactionId>> = TopicProcessor.create()
|
||||
private val head = AtomicReference<BlockJson<TransactionRefJson>>(null)
|
||||
private val stream: TopicProcessor<BlockJson<TransactionRefJson>> = TopicProcessor.create()
|
||||
|
||||
fun follow(source: Flux<BlockJson<TransactionId>>): Disposable {
|
||||
fun follow(source: Flux<BlockJson<TransactionRefJson>>): Disposable {
|
||||
return source.distinctUntilChanged {
|
||||
it.hash
|
||||
}.filter { block ->
|
||||
@@ -37,14 +38,14 @@ open class DefaultEthereumHead: EthereumHead {
|
||||
}
|
||||
}
|
||||
|
||||
override fun getFlux(): Flux<BlockJson<TransactionId>> {
|
||||
override fun getFlux(): Flux<BlockJson<TransactionRefJson>> {
|
||||
return Flux.merge(
|
||||
Mono.justOrEmpty(head.get()),
|
||||
Flux.from(stream)
|
||||
).onBackpressureLatest()
|
||||
}
|
||||
|
||||
fun getCurrent(): BlockJson<TransactionId>? {
|
||||
fun getCurrent(): BlockJson<TransactionRefJson>? {
|
||||
return head.get()
|
||||
}
|
||||
}
|
||||
@@ -17,11 +17,12 @@ package io.emeraldpay.dshackle.upstream.ethereum
|
||||
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import reactor.core.publisher.Flux
|
||||
|
||||
class EmptyEthereumHead : EthereumHead {
|
||||
|
||||
override fun getFlux(): Flux<BlockJson<TransactionId>> {
|
||||
override fun getFlux(): Flux<BlockJson<TransactionRefJson>> {
|
||||
return Flux.empty()
|
||||
}
|
||||
}
|
||||
@@ -18,6 +18,7 @@ package io.emeraldpay.dshackle.upstream.ethereum
|
||||
import io.emeraldpay.dshackle.upstream.Head
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
|
||||
interface EthereumHead: Head<BlockJson<TransactionId>> {
|
||||
interface EthereumHead: Head<BlockJson<TransactionRefJson>> {
|
||||
}
|
||||
@@ -17,6 +17,7 @@ package io.emeraldpay.dshackle.upstream.ethereum
|
||||
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.reactivestreams.Publisher
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.context.Lifecycle
|
||||
@@ -26,7 +27,7 @@ import reactor.core.publisher.Mono
|
||||
import java.util.concurrent.atomic.AtomicReference
|
||||
|
||||
class EthereumHeadMerge(
|
||||
private val fluxes: Iterable<Publisher<BlockJson<TransactionId>>>
|
||||
private val fluxes: Iterable<Publisher<BlockJson<TransactionRefJson>>>
|
||||
): DefaultEthereumHead(), Lifecycle {
|
||||
|
||||
private var subscription: Disposable? = null
|
||||
|
||||
@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.Commands
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import io.infinitape.etherjar.rpc.ws.WebsocketClient
|
||||
import org.slf4j.LoggerFactory
|
||||
import reactor.core.publisher.Flux
|
||||
@@ -37,7 +38,7 @@ class EthereumWs(
|
||||
|
||||
private val log = LoggerFactory.getLogger(EthereumWs::class.java)
|
||||
private val topic = TopicProcessor
|
||||
.builder<BlockJson<TransactionId>>()
|
||||
.builder<BlockJson<TransactionRefJson>>()
|
||||
.name("new-blocks")
|
||||
.build()
|
||||
var basicAuth: UpstreamsConfig.BasicAuth? = null
|
||||
@@ -71,7 +72,7 @@ class EthereumWs(
|
||||
}
|
||||
}
|
||||
|
||||
fun getFlux(): Flux<BlockJson<TransactionId>> {
|
||||
fun getFlux(): Flux<BlockJson<TransactionRefJson>> {
|
||||
return Flux.from(this.topic)
|
||||
.onBackpressureLatest()
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ import io.infinitape.etherjar.domain.TransactionId
|
||||
import io.infinitape.etherjar.rpc.*
|
||||
import io.infinitape.etherjar.rpc.emerald.ReactorEmeraldClient
|
||||
import io.infinitape.etherjar.rpc.json.BlockJson
|
||||
import io.infinitape.etherjar.rpc.json.TransactionRefJson
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.context.Lifecycle
|
||||
import reactor.core.Disposable
|
||||
@@ -110,7 +111,7 @@ open class GrpcUpstream(
|
||||
|
||||
internal fun observeHead(flux: Flux<BlockchainOuterClass.ChainHead>) {
|
||||
val base = flux.map { value ->
|
||||
val block = BlockJson<TransactionId>()
|
||||
val block = BlockJson<TransactionRefJson>()
|
||||
block.number = value.height
|
||||
block.totalDifficulty = BigInteger(1, value.weight.toByteArray())
|
||||
block.hash = BlockHash.from("0x"+value.blockId)
|
||||
|
||||
Reference in New Issue
Block a user