Add BroadcastReader instead of BroadcastQuorum (#263)

This commit is contained in:
KirillPamPam
2023-08-01 17:26:31 +04:00
committed by GitHub
parent 6f576fc20c
commit 674c69f4d3
12 changed files with 497 additions and 207 deletions

View File

@@ -7,16 +7,15 @@ const val SPAN_STATUS_MESSAGE = "status.message"
const val SPAN_READER_RESULT = "reader.result"
const val SPAN_REQUEST_API_TYPE = "request.api.type"
const val SPAN_REQUEST_UPSTREAM_ID = "request.upstreamId"
const val SPAN_RESPONSE_UPSTREAM_ID = "response.upstreamId"
const val SPAN_REQUEST_ID = "request.id"
const val SPAN_NO_RESPONSE_MESSAGE = "no-response.message"
const val LOCAL_READER = "localReader"
const val REMOTE_QUORUM_RPC_READER = "remoteQuorumRpcReader"
const val RPC_READER = "rpcReader"
const val API_READER = "apiReader"
const val CACHE_BLOCK_BY_HASH_READER = "cacheBlockByHashReader"
const val DIRECT_QUORUM_RPC_READER = "directQuorumRpcReader"
const val CACHE_HEIGHT_BY_HASH_READER = "cacheHeightByHashReader"
const val CACHE_BLOCK_BY_HEIGHT_READER = "cacheBlockByHeightReader"
const val CACHE_TX_BY_HASH_READER = "cacheTxByHashReader"
const val CACHE_RECEIPTS_READER = "cacheReceiptsReader"
const val BROADCAST_READER = "broadcastReader"

View File

@@ -1,45 +0,0 @@
/**
* Copyright (c) 2020 EmeraldPay, Inc
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.emeraldpay.dshackle.quorum
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.ApiSource
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
import org.springframework.cloud.sleuth.Tracer
import java.util.concurrent.atomic.AtomicInteger
interface QuorumReader : Reader<JsonRpcRequest, QuorumRpcReader.Result> {
fun attempts(): AtomicInteger
}
// creates instance of a Quorum based reader
interface QuorumReaderFactory {
companion object {
fun default(): QuorumReaderFactory {
return Default()
}
}
fun create(apis: ApiSource, quorum: CallQuorum, signer: ResponseSigner?, tracer: Tracer): QuorumReader
class Default : QuorumReaderFactory {
override fun create(apis: ApiSource, quorum: CallQuorum, signer: ResponseSigner?, tracer: Tracer): QuorumReader {
return QuorumRpcReader(apis, quorum, signer, tracer)
}
}
}

View File

@@ -20,10 +20,10 @@ import io.emeraldpay.dshackle.commons.API_READER
import io.emeraldpay.dshackle.commons.SPAN_NO_RESPONSE_MESSAGE
import io.emeraldpay.dshackle.commons.SPAN_REQUEST_API_TYPE
import io.emeraldpay.dshackle.commons.SPAN_REQUEST_UPSTREAM_ID
import io.emeraldpay.dshackle.reader.RpcReader
import io.emeraldpay.dshackle.reader.SpannedReader
import io.emeraldpay.dshackle.upstream.ApiSource
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
@@ -47,9 +47,9 @@ import java.util.function.Function
class QuorumRpcReader(
private val apiControl: ApiSource,
private val quorum: CallQuorum,
private val signer: ResponseSigner?,
signer: ResponseSigner?,
private val tracer: Tracer
) : QuorumReader {
) : RpcReader(signer) {
companion object {
private val log = LoggerFactory.getLogger(QuorumRpcReader::class.java)
@@ -158,12 +158,7 @@ class QuorumRpcReader(
private fun withSignatureAndUpstream(api: Upstream, key: JsonRpcRequest, response: JsonRpcResponse): Function<Mono<ByteArray>, Mono<Tuple2<ByteArray, Optional<ResponseSigner.Signature>>>> {
return Function { src ->
src.map {
val signature = response.providedSignature
?: if (key.nonce != null) {
signer?.sign(key.nonce, response.getResult(), api.getId())
} else {
null
}
val signature = getSignature(key, response, api.getId())
Tuples.of(it, Optional.ofNullable(signature))
}
}
@@ -176,14 +171,7 @@ class QuorumRpcReader(
// when the call failed with an error we want to notify the quorum because
// it may use the error message or other details
//
val cleanErr: JsonRpcException = when (err) {
is RpcException -> JsonRpcException.from(err)
is JsonRpcException -> err
else -> JsonRpcException(
JsonRpcResponse.NumberId(key.id),
JsonRpcError(-32603, "Unhandled internal error: ${err.javaClass}: ${err.message}")
)
}
val cleanErr: JsonRpcException = getError(key, err)
quorum.record(cleanErr, null, api,)
// if it's failed after that, then we don't need more calls, stop api source
if (quorum.isFailed()) {
@@ -202,8 +190,7 @@ class QuorumRpcReader(
return Mono.just(quorum).flatMap { q ->
if (q.isFailed()) {
val resolvedBy = resolvedBy()?.getId()
val err = q.getError()?.asException(JsonRpcResponse.NumberId(key.id), resolvedBy)
?: JsonRpcException(JsonRpcResponse.NumberId(key.id), JsonRpcError(-32603, "Unhandled Upstream error"), resolvedBy)
val err = handleError(q.getError(), key.id, resolvedBy)
log.warn("Quorum is failed. Method ${key.method}, message ${err.message}")
Mono.error(err)
} else {
@@ -229,11 +216,4 @@ class QuorumRpcReader(
}
} ?: Mono.error(RpcException(1, "Quorum [$q] is not resolved [isResolved - ${q.isResolved()}]"))
}
class Result(
val value: ByteArray,
val signature: ResponseSigner.Signature?,
val quorum: Int,
val resolvedBy: Upstream?
)
}

View File

@@ -0,0 +1,118 @@
package io.emeraldpay.dshackle.reader
import io.emeraldpay.dshackle.commons.BROADCAST_READER
import io.emeraldpay.dshackle.commons.SPAN_REQUEST_UPSTREAM_ID
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
import org.slf4j.LoggerFactory
import org.springframework.cloud.sleuth.Tracer
import reactor.core.publisher.Mono
import java.util.concurrent.atomic.AtomicInteger
class BroadcastReader(
private val upstreams: List<Upstream>,
matcher: Selector.Matcher,
signer: ResponseSigner?,
private val tracer: Tracer
) : RpcReader(signer) {
private val internalMatcher = Selector.MultiMatcher(
listOf(Selector.AvailabilityMatcher(), matcher)
)
companion object {
private val log = LoggerFactory.getLogger(BroadcastReader::class.java)
}
override fun attempts(): AtomicInteger {
return AtomicInteger(1)
}
override fun read(key: JsonRpcRequest): Mono<Result> {
return Mono.just(upstreams)
.map { ups ->
ups.filter { internalMatcher.matches(it) }.map { execute(key, it) }
}.flatMap {
Mono.zip(it) { responses ->
analyzeResponses(
key,
getJsonRpcResponses(key.method, responses)
)
}.onErrorResume { err ->
log.error("Broadcast error: ${err.message}")
Mono.error(handleError(null, 0, null))
}.flatMap { broadcastResult ->
if (broadcastResult.result != null) {
Mono.just(
Result(broadcastResult.result, broadcastResult.signature, 0, null)
)
} else {
val err = handleError(broadcastResult.error, key.id, null)
Mono.error(err)
}
}
}
}
private fun analyzeResponses(key: JsonRpcRequest, jsonRpcResponses: List<BroadcastResponse>): BroadcastResult {
val errors = mutableListOf<JsonRpcError>()
jsonRpcResponses.forEach {
val response = it.jsonRpcResponse
if (response.hasResult()) {
val signature = getSignature(key, response, it.upstreamId)
return BroadcastResult(response.getResult(), null, signature)
} else if (response.hasError()) {
errors.add(response.error!!)
}
}
val error = errors.takeIf { it.isNotEmpty() }?.get(0)
return BroadcastResult(error)
}
private fun getJsonRpcResponses(method: String, responses: Array<Any>) =
responses
.map { response ->
(response as BroadcastResponse)
.also { r ->
if (r.jsonRpcResponse.hasResult()) {
log.info(
"Response for $method from upstream ${r.upstreamId}: ${String(r.jsonRpcResponse.getResult())}"
)
}
}
}
private fun execute(
key: JsonRpcRequest,
upstream: Upstream
): Mono<BroadcastResponse> =
SpannedReader(
upstream.getIngressReader(), tracer, BROADCAST_READER, mapOf(SPAN_REQUEST_UPSTREAM_ID to upstream.getId())
)
.read(key)
.map { BroadcastResponse(it, upstream.getId()) }
.onErrorResume {
log.warn("Error during execution ${key.method} from upstream ${upstream.getId()} with message - ${it.message}")
Mono.just(
BroadcastResponse(JsonRpcResponse(null, getError(key, it).error), upstream.getId())
)
}
private class BroadcastResponse(
val jsonRpcResponse: JsonRpcResponse,
val upstreamId: String
)
private class BroadcastResult(
val result: ByteArray?,
val error: JsonRpcError?,
val signature: ResponseSigner.Signature?
) {
constructor(error: JsonRpcError?) : this(null, error, null)
}
}

View File

@@ -0,0 +1,81 @@
package io.emeraldpay.dshackle.reader
import io.emeraldpay.dshackle.quorum.CallQuorum
import io.emeraldpay.dshackle.quorum.QuorumRpcReader
import io.emeraldpay.dshackle.reader.RpcReader.Result
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
import io.emeraldpay.etherjar.rpc.RpcException
import org.springframework.cloud.sleuth.Tracer
import java.util.concurrent.atomic.AtomicInteger
abstract class RpcReader(
private val signer: ResponseSigner?,
) : Reader<JsonRpcRequest, Result> {
abstract fun attempts(): AtomicInteger
protected fun getError(key: JsonRpcRequest, err: Throwable) =
when (err) {
is RpcException -> JsonRpcException.from(err)
is JsonRpcException -> err
else -> JsonRpcException(
JsonRpcResponse.NumberId(key.id),
JsonRpcError(-32603, "Unhandled internal error: ${err.javaClass}: ${err.message}")
)
}
protected fun handleError(error: JsonRpcError?, id: Int, resolvedBy: String?) =
error?.asException(JsonRpcResponse.NumberId(id), resolvedBy)
?: JsonRpcException(JsonRpcResponse.NumberId(id), JsonRpcError(-32603, "Unhandled Upstream error"), resolvedBy)
protected fun getSignature(key: JsonRpcRequest, response: JsonRpcResponse, upstreamId: String) =
response.providedSignature
?: if (key.nonce != null) {
signer?.sign(key.nonce, response.getResult(), upstreamId)
} else {
null
}
class Result(
val value: ByteArray,
val signature: ResponseSigner.Signature?,
val quorum: Int,
val resolvedBy: Upstream?
)
}
interface RpcReaderFactory {
companion object {
fun default(): RpcReaderFactory {
return Default()
}
}
fun create(data: RpcReaderData): RpcReader
class Default : RpcReaderFactory {
override fun create(data: RpcReaderData): RpcReader {
if (data.method == "eth_sendRawTransaction") {
return BroadcastReader(data.multistream.getAll(), data.matcher, data.signer, data.tracer)
}
val apis = data.multistream.getApiSource(data.matcher)
return QuorumRpcReader(apis, data.quorum, data.signer, data.tracer)
}
}
data class RpcReaderData(
val multistream: Multistream,
val method: String,
val matcher: Selector.Matcher,
val quorum: CallQuorum,
val signer: ResponseSigner?,
val tracer: Tracer
)
}

View File

@@ -25,15 +25,16 @@ import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.Global.Companion.nullValue
import io.emeraldpay.dshackle.SilentException
import io.emeraldpay.dshackle.commons.LOCAL_READER
import io.emeraldpay.dshackle.commons.REMOTE_QUORUM_RPC_READER
import io.emeraldpay.dshackle.commons.RPC_READER
import io.emeraldpay.dshackle.commons.SPAN_ERROR
import io.emeraldpay.dshackle.commons.SPAN_REQUEST_ID
import io.emeraldpay.dshackle.commons.SPAN_STATUS_MESSAGE
import io.emeraldpay.dshackle.config.MainConfig
import io.emeraldpay.dshackle.quorum.CallQuorum
import io.emeraldpay.dshackle.quorum.NotLaggingQuorum
import io.emeraldpay.dshackle.quorum.QuorumReaderFactory
import io.emeraldpay.dshackle.quorum.QuorumRpcReader
import io.emeraldpay.dshackle.reader.RpcReader
import io.emeraldpay.dshackle.reader.RpcReaderFactory
import io.emeraldpay.dshackle.reader.RpcReaderFactory.RpcReaderData
import io.emeraldpay.dshackle.reader.SpannedReader
import io.emeraldpay.dshackle.startup.UpstreamChangeEvent
import io.emeraldpay.dshackle.upstream.ApiSource
@@ -79,7 +80,7 @@ open class NativeCall(
private val localRouterEnabled = config.cache?.requestsCacheEnabled ?: true
private val passthrough = config.passthrough
var quorumReaderFactory: QuorumReaderFactory = QuorumReaderFactory.default()
var rpcReaderFactory: RpcReaderFactory = RpcReaderFactory.default()
private val ethereumCallSelectors = EnumMap<Chain, EthereumCallSelector>(Chain::class.java)
companion object {
@@ -391,15 +392,18 @@ open class NativeCall(
if (!ctx.upstream.getMethods().isCallable(ctx.payload.method)) {
return Mono.error(RpcException(RpcResponseError.CODE_METHOD_NOT_EXIST, "Unsupported method"))
}
val reader = quorumReaderFactory.create(ctx.getApis(), ctx.callQuorum, signer, tracer)
val reader = rpcReaderFactory.create(
RpcReaderData(ctx.upstream, ctx.payload.method, ctx.matcher, ctx.callQuorum, signer, tracer)
)
val counter = reader.attempts()
return SpannedReader(reader, tracer, REMOTE_QUORUM_RPC_READER)
return SpannedReader(reader, tracer, RPC_READER)
.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce, ctx.forwardedSelector))
.map {
val bytes = ctx.resultDecorator.processResult(it)
validateResult(bytes, "remote", ctx)
CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, it.resolvedBy?.getId(), ctx)
val upId = it.resolvedBy?.getId() ?: ctx.upstream.getId()
CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, upId, ctx)
}
.onErrorResume { t ->
Mono.just(CallResult.fail(ctx.id, ctx.nonce, t, ctx))
@@ -470,11 +474,11 @@ open class NativeCall(
}
interface ResultDecorator {
fun processResult(result: QuorumRpcReader.Result): ByteArray
fun processResult(result: RpcReader.Result): ByteArray
}
open class NoneResultDecorator : ResultDecorator {
override fun processResult(result: QuorumRpcReader.Result): ByteArray = result.value
override fun processResult(result: RpcReader.Result): ByteArray = result.value
}
open class CreateFilterDecorator : ResultDecorator {
@@ -482,7 +486,7 @@ open class NativeCall(
companion object {
const val quoteCode = '"'.code.toByte()
}
override fun processResult(result: QuorumRpcReader.Result): ByteArray {
override fun processResult(result: RpcReader.Result): ByteArray {
val bytes = result.value
if (bytes.last() == quoteCode && result.resolvedBy != null) {
val suffix = result.resolvedBy.nodeId()

View File

@@ -19,7 +19,6 @@ package io.emeraldpay.dshackle.upstream.calls
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.quorum.AlwaysQuorum
import io.emeraldpay.dshackle.quorum.BroadcastQuorum
import io.emeraldpay.dshackle.quorum.CallQuorum
import io.emeraldpay.dshackle.quorum.NotLaggingQuorum
import io.emeraldpay.dshackle.quorum.NotNullQuorum
@@ -198,7 +197,6 @@ class DefaultEthereumMethods(
when (method) {
"eth_getTransactionCount" -> AlwaysQuorum()
"eth_getBalance" -> AlwaysQuorum()
"eth_sendRawTransaction" -> BroadcastQuorum()
"eth_blockNumber" -> NotLaggingQuorum(0)
else -> AlwaysQuorum()
}

View File

@@ -10,8 +10,8 @@ import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.data.DefaultContainer
import io.emeraldpay.dshackle.data.TxContainer
import io.emeraldpay.dshackle.data.TxId
import io.emeraldpay.dshackle.quorum.QuorumReaderFactory
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.reader.RpcReaderFactory
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.calls.CallMethods
@@ -51,7 +51,7 @@ class EthereumDirectReader(
}
private val objectMapper: ObjectMapper = Global.objectMapper
var quorumReaderFactory: QuorumReaderFactory = QuorumReaderFactory.default()
var rpcReaderFactory: RpcReaderFactory = RpcReaderFactory.default()
val blockReader: Reader<BlockHash, Result<BlockContainer>>
val blockByHeightReader: Reader<Long, Result<BlockContainer>>
@@ -183,19 +183,17 @@ class EthereumDirectReader(
request: JsonRpcRequest,
matcher: Selector.Matcher = Selector.empty
): Mono<Result<ByteArray>> {
return Mono.just(quorumReaderFactory)
return Mono.just(rpcReaderFactory)
.map {
val requestMatcher = Selector.Builder()
.withMatcher(matcher)
.forMethod(request.method)
.build()
it.create(
up.getApiSource(
Selector.Builder()
.withMatcher(matcher)
.forMethod(request.method)
.build()
),
callMethodsFactory.create().createQuorumFor(request.method),
// we do not use Signer for internal requests because it doesn't make much sense
null,
tracer
RpcReaderFactory.RpcReaderData(
up, request.method, requestMatcher,
callMethodsFactory.create().createQuorumFor(request.method), null, tracer
)
)
}.flatMap {
it.read(request)