From 121dcf0ca7e6476ef89e8603a46eea1898775071 Mon Sep 17 00:00:00 2001 From: KirillPamPam Date: Thu, 20 Jul 2023 16:12:40 +0400 Subject: [PATCH] Return upstreamId in response for errors too (#254) --- .../dshackle/quorum/AlwaysQuorum.kt | 16 +++------ .../dshackle/quorum/BroadcastQuorum.kt | 14 ++------ .../emeraldpay/dshackle/quorum/CallQuorum.kt | 5 ++- .../dshackle/quorum/NotLaggingQuorum.kt | 17 +++------ .../dshackle/quorum/NotNullQuorum.kt | 11 ++---- .../dshackle/quorum/QuorumRpcReader.kt | 29 ++++++++------- .../dshackle/quorum/ValueAwareQuorum.kt | 23 +++++------- .../io/emeraldpay/dshackle/rpc/NativeCall.kt | 36 +++++++++++-------- .../upstream/ethereum/WsConnectionImpl.kt | 1 + .../upstream/rpcclient/JsonRpcError.kt | 6 +++- .../upstream/rpcclient/JsonRpcException.kt | 1 + .../upstream/rpcclient/JsonRpcHttpClient.kt | 2 +- .../dshackle/proxy/WriteRpcJsonSpec.groovy | 4 +-- .../dshackle/quorum/AlwaysQuorumSpec.groovy | 2 +- .../quorum/BroadcastQuorumSpec.groovy | 18 +++++----- .../quorum/NotLaggingQuorumSpec.groovy | 9 +++-- .../dshackle/quorum/NotNullQuorumSpec.groovy | 11 +++--- .../quorum/ValueAwareQuorumSpec.groovy | 9 ++--- .../dshackle/rpc/NativeCallSpec.groovy | 8 ++--- .../ethereum/EthereumCachingReaderSpec.groovy | 25 ++++++------- 20 files changed, 112 insertions(+), 135 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt index d9bd712b..5211765d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/AlwaysQuorum.kt @@ -21,7 +21,6 @@ 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.signature.ResponseSigner -import java.util.concurrent.ConcurrentLinkedQueue open class AlwaysQuorum : CallQuorum { @@ -29,8 +28,7 @@ open class AlwaysQuorum : CallQuorum { private var result: ByteArray? = null private var rpcError: JsonRpcError? = null private var sig: ResponseSigner.Signature? = null - private var providedUpstreamId: String? = null - private val resolvers: MutableCollection = ConcurrentLinkedQueue() + private val resolvers = ArrayList() override fun init(head: Head) { } @@ -47,21 +45,15 @@ open class AlwaysQuorum : CallQuorum { return sig } - override fun getProvidedUpstreamId(): String? { - return providedUpstreamId - } - override fun record( response: ByteArray, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ): Boolean { result = response resolved = true sig = signature resolvers.add(upstream) - this.providedUpstreamId = providedUpstreamId return true } @@ -72,6 +64,7 @@ open class AlwaysQuorum : CallQuorum { ) { this.rpcError = error.error sig = signature + resolvers.add(upstream) } override fun getResult(): ByteArray? { @@ -82,8 +75,7 @@ open class AlwaysQuorum : CallQuorum { return rpcError } - override fun getResolvedBy(): List = - resolvers.toList() + override fun getResolvedBy(): List = resolvers override fun toString(): String { return "Quorum: Accept Any" diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/BroadcastQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/BroadcastQuorum.kt index 1578ae0b..ec3faf3f 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/BroadcastQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/BroadcastQuorum.kt @@ -28,7 +28,6 @@ open class BroadcastQuorum( private var txid: String? = null private var calls = 0 private var sig: ResponseSigner.Signature? = null - private var providedUpstreamId: String? = null override fun init(head: Head) { } @@ -49,23 +48,17 @@ open class BroadcastQuorum( return sig } - override fun getProvidedUpstreamId(): String? { - return providedUpstreamId - } - override fun recordValue( response: ByteArray, responseValue: String?, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ) { calls++ if (txid == null && responseValue != null) { txid = responseValue sig = signature result = response - this.providedUpstreamId = providedUpstreamId } } @@ -73,16 +66,15 @@ open class BroadcastQuorum( response: ByteArray?, errorMessage: String?, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ) { // can be "message: known transaction: TXID", "Transaction with the same hash was already imported" or "message: Nonce too low" calls++ if (result == null) { result = response sig = signature - this.providedUpstreamId = providedUpstreamId } + resolvers.add(upstream) } override fun toString(): String { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt index c8bb15ca..5fc34f87 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/CallQuorum.kt @@ -32,9 +32,9 @@ interface CallQuorum { fun record( response: ByteArray, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ): Boolean + fun record( error: JsonRpcException, signature: ResponseSigner.Signature?, @@ -42,7 +42,6 @@ interface CallQuorum { ) fun getSignature(): ResponseSigner.Signature? - fun getProvidedUpstreamId(): String? fun getResult(): ByteArray? fun getError(): JsonRpcError? fun getResolvedBy(): Collection diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt index d553d5e9..1e71b774 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotLaggingQuorum.kt @@ -21,7 +21,6 @@ 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.signature.ResponseSigner -import java.util.concurrent.ConcurrentLinkedQueue import java.util.concurrent.atomic.AtomicReference /** @@ -35,8 +34,7 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { private val failed = AtomicReference(false) private var rpcError: JsonRpcError? = null private var sig: ResponseSigner.Signature? = null - private val resolvers: MutableCollection = ConcurrentLinkedQueue() - private var providedUpstreamId: String? = null + private val resolvers = ArrayList() override fun init(head: Head) { } @@ -52,14 +50,12 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { override fun record( response: ByteArray, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ): Boolean { val lagging = upstream.getLag()?.run { this > maxLag } ?: true if (!lagging) { result.set(response) sig = signature - this.providedUpstreamId = providedUpstreamId resolvers.add(upstream) return true } @@ -76,16 +72,13 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { if (!lagging && result.get() == null) { failed.set(true) } + resolvers.add(upstream) } override fun getSignature(): ResponseSigner.Signature? { return sig } - override fun getProvidedUpstreamId(): String? { - return providedUpstreamId - } - override fun getResult(): ByteArray { return result.get() } @@ -94,8 +87,8 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum { return rpcError } - override fun getResolvedBy(): Collection = - resolvers.toList() + override fun getResolvedBy(): Collection = resolvers + override fun toString(): String { return "Quorum: late <= $maxLag blocks" } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotNullQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotNullQuorum.kt index e1007d38..6f06c91d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotNullQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/NotNullQuorum.kt @@ -9,7 +9,6 @@ import io.emeraldpay.dshackle.upstream.signature.ResponseSigner class NotNullQuorum : CallQuorum { private var sig: ResponseSigner.Signature? = null - private var providedUpstreamId: String? = null private var result: ByteArray? = null private var rpcError: JsonRpcError? = null private val resolvers = ArrayList() @@ -26,8 +25,7 @@ class NotNullQuorum : CallQuorum { override fun record( response: ByteArray, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ): Boolean { allFailed = false val receivedNull = response.isEmpty() || Global.nullValue.contentEquals(response) @@ -35,7 +33,6 @@ class NotNullQuorum : CallQuorum { if (seenUpstreams.contains(upId) || !receivedNull) { sig = signature result = response - this.providedUpstreamId = providedUpstreamId resolvers.add(upstream) return true } @@ -50,22 +47,20 @@ class NotNullQuorum : CallQuorum { rpcError = error.error } else { result = Global.nullValue - resolvers.add(upstream) } sig = signature } + resolvers.add(upstream) seenUpstreams.add(upId) } override fun getSignature(): ResponseSigner.Signature? = sig - override fun getProvidedUpstreamId(): String? = providedUpstreamId - override fun getResult(): ByteArray? = result override fun getError(): JsonRpcError? = rpcError - override fun getResolvedBy(): Collection = resolvers.toList() + override fun getResolvedBy(): Collection = resolvers override fun toString(): String { return "Quorum: Not null" diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt index 546d9903..78ba2798 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/QuorumRpcReader.kt @@ -33,8 +33,8 @@ import org.slf4j.LoggerFactory import org.springframework.cloud.sleuth.Tracer import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import reactor.util.function.Tuple2 import reactor.util.function.Tuple3 -import reactor.util.function.Tuple4 import reactor.util.function.Tuples import java.util.Optional import java.util.concurrent.atomic.AtomicInteger @@ -99,8 +99,8 @@ class QuorumRpcReader( } private fun execute(key: JsonRpcRequest, retrySpec: reactor.util.retry.Retry): Function, Mono> { - val quorumReduce = BiFunction, Upstream, Optional>, CallQuorum> { res, a -> - if (res.record(a.t1, a.t2.orElse(null), a.t3, a.t4.orElse(null))) { + val quorumReduce = BiFunction, Upstream>, CallQuorum> { res, a -> + if (res.record(a.t1, a.t2.orElse(null), a.t3)) { log.trace("Quorum is resolved for method ${key.method}") apiControl.resolve() } else { @@ -131,13 +131,13 @@ class QuorumRpcReader( .filter { it.isResolved() } // return nothing if not resolved .map { quorum -> // TODO find actual quorum number - Result(quorum.getResult()!!, quorum.getSignature(), 1, quorum.getResolvedBy(), quorum.getProvidedUpstreamId()) + Result(quorum.getResult()!!, quorum.getSignature(), 1, resolvedBy()) } .switchIfEmpty(defaultResult) } } - private fun callApi(api: Upstream, key: JsonRpcRequest): Mono, Upstream, Optional>> { + private fun callApi(api: Upstream, key: JsonRpcRequest): Mono, Upstream>> { val apiReader = api.getIngressReader() val spanParams = mapOf( SPAN_REQUEST_API_TYPE to apiReader.javaClass.name, @@ -152,10 +152,10 @@ class QuorumRpcReader( } // must catch not only the processing of a response but also errors thrown from the .read() call .transform(withErrorResume(api, key)) - .map { Tuples.of(it.t1, it.t2, api, it.t3) } + .map { Tuples.of(it.t1, it.t2, api) } } - private fun withSignatureAndUpstream(api: Upstream, key: JsonRpcRequest, response: JsonRpcResponse): Function, Mono, Optional>>> { + private fun withSignatureAndUpstream(api: Upstream, key: JsonRpcRequest, response: JsonRpcResponse): Function, Mono>>> { return Function { src -> src.map { val signature = response.providedSignature @@ -164,7 +164,7 @@ class QuorumRpcReader( } else { null } - Tuples.of(it, Optional.ofNullable(signature), Optional.ofNullable(response.providedUpstreamId)) + Tuples.of(it, Optional.ofNullable(signature)) } } } @@ -201,8 +201,9 @@ class QuorumRpcReader( private fun setupDefaultResult(key: JsonRpcRequest): Mono { return Mono.just(quorum).flatMap { q -> if (q.isFailed()) { - val err = q.getError()?.asException(JsonRpcResponse.NumberId(key.id)) - ?: JsonRpcException(JsonRpcResponse.NumberId(key.id), JsonRpcError(-32603, "Unhandled Upstream error")) + 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) log.warn("Quorum is failed. Method ${key.method}, message ${err.message}") Mono.error(err) } else { @@ -212,13 +213,16 @@ class QuorumRpcReader( } } + private fun resolvedBy() = + if (quorum.getResolvedBy().isEmpty()) null else quorum.getResolvedBy().last() + private fun noResponse(method: String, q: CallQuorum): Mono { return apiControl.upstreamsMatchesResponse()?.run { tracer.currentSpan()?.tag(SPAN_NO_RESPONSE_MESSAGE, getFullCause()) val cause = getCause(method) ?: return Mono.empty() if (cause.shouldReturnNull) { Mono.just( - Result(Global.nullValue, null, 1, emptyList(), null) + Result(Global.nullValue, null, 1, null) ) } else { Mono.error(RpcException(1, "No response for method $method. Cause - ${cause.cause}")) @@ -230,7 +234,6 @@ class QuorumRpcReader( val value: ByteArray, val signature: ResponseSigner.Signature?, val quorum: Int, - val resolvers: Collection, - val providedUpstreamId: String? + val resolvedBy: Upstream? ) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt b/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt index 65bce368..b68e8004 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/quorum/ValueAwareQuorum.kt @@ -23,7 +23,6 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException import io.emeraldpay.dshackle.upstream.signature.ResponseSigner import io.emeraldpay.etherjar.rpc.RpcException import org.slf4j.LoggerFactory -import java.util.concurrent.ConcurrentLinkedQueue abstract class ValueAwareQuorum( val clazz: Class @@ -31,7 +30,7 @@ abstract class ValueAwareQuorum( private val log = LoggerFactory.getLogger(ValueAwareQuorum::class.java) private var rpcError: JsonRpcError? = null - private val resolvers: MutableCollection = ConcurrentLinkedQueue() + protected val resolvers = ArrayList() fun extractValue(response: ByteArray, clazz: Class): T? { return Global.objectMapper.readValue(response.inputStream(), clazz) @@ -40,17 +39,16 @@ abstract class ValueAwareQuorum( override fun record( response: ByteArray, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ): Boolean { try { val value = extractValue(response, clazz) - recordValue(response, value, signature, upstream, providedUpstreamId) + recordValue(response, value, signature, upstream) resolvers.add(upstream) } catch (e: RpcException) { - recordError(response, e.rpcMessage, signature, upstream, providedUpstreamId) + recordError(response, e.rpcMessage, signature, upstream) } catch (e: Exception) { - recordError(response, e.message, signature, upstream, providedUpstreamId) + recordError(response, e.message, signature, upstream) } return isResolved() } @@ -61,29 +59,26 @@ abstract class ValueAwareQuorum( upstream: Upstream ) { this.rpcError = error.error - recordError(null, error.error.message, signature, upstream, null) + recordError(null, error.error.message, signature, upstream) } abstract fun recordValue( response: ByteArray, responseValue: T?, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ) abstract fun recordError( response: ByteArray?, errorMessage: String?, signature: ResponseSigner.Signature?, - upstream: Upstream, - providedUpstreamId: String? + upstream: Upstream ) override fun getError(): JsonRpcError? { return rpcError } - override fun getResolvedBy(): Collection = - resolvers.toList() + override fun getResolvedBy(): Collection = resolvers } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt index e7e9b0d5..aa3bd758 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt @@ -28,7 +28,6 @@ import io.emeraldpay.dshackle.commons.LOCAL_READER import io.emeraldpay.dshackle.commons.REMOTE_QUORUM_RPC_READER import io.emeraldpay.dshackle.commons.SPAN_ERROR import io.emeraldpay.dshackle.commons.SPAN_REQUEST_ID -import io.emeraldpay.dshackle.commons.SPAN_RESPONSE_UPSTREAM_ID import io.emeraldpay.dshackle.commons.SPAN_STATUS_MESSAGE import io.emeraldpay.dshackle.config.MainConfig import io.emeraldpay.dshackle.quorum.CallQuorum @@ -131,9 +130,6 @@ open class NativeCall( private fun completeSpan(callResult: CallResult, requestCount: Int) { val span = tracer.currentSpan() - callResult.upstreamId?.let { - span?.tag(SPAN_RESPONSE_UPSTREAM_ID, it) - } if (callResult.isError()) { errorSpan(span, callResult.error?.message ?: "Internal error") } @@ -403,9 +399,7 @@ open class NativeCall( .map { val bytes = ctx.resultDecorator.processResult(it) validateResult(bytes, "remote", ctx) - val upId = it.providedUpstreamId - ?: if (it.resolvers.isEmpty()) ctx.upstream.getId() else it.resolvers.first().getId() - CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, upId, ctx) + CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, it.resolvedBy?.getId(), ctx) } .onErrorResume { t -> Mono.just(CallResult.fail(ctx.id, ctx.nonce, t, ctx)) @@ -490,10 +484,8 @@ open class NativeCall( } override fun processResult(result: QuorumRpcReader.Result): ByteArray { val bytes = result.value - if (bytes.last() == quoteCode) { - val suffix = result.resolvers - .map { it.nodeId() } - .first() + if (bytes.last() == quoteCode && result.resolvedBy != null) { + val suffix = result.resolvedBy.nodeId() .toUByte() .toString(16).padStart(2, padChar = '0').toByteArray() bytes[bytes.lastIndex] = suffix.first() @@ -598,7 +590,13 @@ open class NativeCall( open class CallFailure(val id: Int, val reason: Throwable) : Exception("Failed to call $id: ${reason.message}") - open class CallError(val id: Int, val message: String, val upstreamError: JsonRpcError?, val data: String?) { + open class CallError( + val id: Int, + val message: String, + val upstreamError: JsonRpcError?, + val data: String?, + val upstreamId: String? = null + ) { companion object { @@ -615,7 +613,7 @@ open class NativeCall( } fun from(t: Throwable): CallError { return when (t) { - is JsonRpcException -> CallError(t.error.code, t.error.message, t.error, getDataAsSting(t.error.details)) + is JsonRpcException -> CallError(t.error.code, t.error.message, t.error, getDataAsSting(t.error.details), t.upstreamId) is RpcException -> CallError(t.code, t.rpcMessage, null, getDataAsSting(t.details)) is CallFailure -> CallError(t.id, t.reason.message ?: "Upstream Error", null, null) else -> { @@ -641,6 +639,16 @@ open class NativeCall( val upstreamId: String?, val ctx: ValidCallContext? ) { + + constructor( + id: Int, + nonce: Long?, + result: ByteArray?, + callError: CallError?, + signature: ResponseSigner.Signature?, + ctx: ValidCallContext? + ) : this(id, nonce, result, callError, signature, callError?.upstreamId, ctx) + companion object { fun ok(id: Int, nonce: Long?, result: ByteArray, signature: ResponseSigner.Signature?, upstreamId: String?, ctx: ValidCallContext?): CallResult { return CallResult(id, nonce, result, null, signature, upstreamId, ctx) @@ -651,7 +659,7 @@ open class NativeCall( } fun fail(id: Int, nonce: Long?, error: Throwable, ctx: ValidCallContext?): CallResult { - return CallResult(id, nonce, null, CallError.from(error), null, null, ctx) + return CallResult(id, nonce, null, CallError.from(error), null, ctx) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt index 2aab647f..05ef39d2 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionImpl.kt @@ -377,6 +377,7 @@ open class WsConnectionImpl( RpcResponseError.CODE_INTERNAL_ERROR, "Response not received from WebSocket" ), + null, false ) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcError.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcError.kt index fe51daab..773fd80c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcError.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcError.kt @@ -31,6 +31,10 @@ data class JsonRpcError(val code: Int, val message: String, val details: Any?) { } fun asException(id: JsonRpcResponse.Id?): JsonRpcException { - return JsonRpcException(id ?: JsonRpcResponse.NumberId(-1), this, false) + return JsonRpcException(id ?: JsonRpcResponse.NumberId(-1), this, null, false) + } + + fun asException(id: JsonRpcResponse.Id?, upstreamId: String?): JsonRpcException { + return JsonRpcException(id ?: JsonRpcResponse.NumberId(-1), this, upstreamId, false) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcException.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcException.kt index b1a1b1d4..794291f9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcException.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcException.kt @@ -20,6 +20,7 @@ import io.emeraldpay.etherjar.rpc.RpcException class JsonRpcException( val id: JsonRpcResponse.Id, val error: JsonRpcError, + val upstreamId: String? = null, writableStackTrace: Boolean = true ) : Exception(error.message, null, true, writableStackTrace) { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt index 13d0808e..4ec3d81d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt @@ -131,7 +131,7 @@ class JsonRpcHttpClient( return Function { resp -> resp.flatMap { if (it.hasError()) { - Mono.error(JsonRpcException(it.id, it.error!!, false)) + Mono.error(JsonRpcException(it.id, it.error!!, null, false)) } else { Mono.just(it) } diff --git a/src/test/groovy/io/emeraldpay/dshackle/proxy/WriteRpcJsonSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/proxy/WriteRpcJsonSpec.groovy index 2214e68b..53e3bf1a 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/proxy/WriteRpcJsonSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/proxy/WriteRpcJsonSpec.groovy @@ -98,7 +98,7 @@ class WriteRpcJsonSpec extends Specification { def call = new ProxyCall(ProxyCall.RpcType.SINGLE) call.ids[1] = 1 def data = [ - new NativeCall.CallResult(1, null, null, new NativeCall.CallError(1, "Internal Error", null, null), null, null, null) + new NativeCall.CallResult(1, null, null, new NativeCall.CallError(1, "Internal Error", null, null, null), null, null, null) ] when: def act = writer.toJson(call, data[0]) @@ -127,7 +127,7 @@ class WriteRpcJsonSpec extends Specification { call.ids[3] = 15 def data = [ new NativeCall.CallResult(1, null, '"0x98dbb1"'.bytes, null, null, null, null), - new NativeCall.CallResult(2, null, null, new NativeCall.CallError(2, "oops", null, null), null, null, null), + new NativeCall.CallResult(2, null, null, new NativeCall.CallError(2, "oops", null, null, null), null, null, null), new NativeCall.CallResult(3, null, '{"hash": "0x2484f459dc"}'.bytes, null, null, null, null), ] when: diff --git a/src/test/groovy/io/emeraldpay/dshackle/quorum/AlwaysQuorumSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/quorum/AlwaysQuorumSpec.groovy index 1b2b403e..d34039c0 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/quorum/AlwaysQuorumSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/quorum/AlwaysQuorumSpec.groovy @@ -42,7 +42,7 @@ class AlwaysQuorumSpec extends Specification { def quorum = new AlwaysQuorum() def up = Stub(Upstream) when: - quorum.record("123".bytes, new ResponseSigner.Signature("sig1".bytes, "test", 100), up, null) + quorum.record("123".bytes, new ResponseSigner.Signature("sig1".bytes, "test", 100), up) then: quorum.isResolved() quorum.getResult() == "123".bytes diff --git a/src/test/groovy/io/emeraldpay/dshackle/quorum/BroadcastQuorumSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/quorum/BroadcastQuorumSpec.groovy index 14193481..ee770db4 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/quorum/BroadcastQuorumSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/quorum/BroadcastQuorumSpec.groovy @@ -40,21 +40,21 @@ class BroadcastQuorumSpec extends Specification { !q.isResolved() when: - q.record('"0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c"'.bytes, null, upstream1, null) + q.record('"0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c"'.bytes, null, upstream1) then: !q.isResolved() - 1 * q.recordValue(_, "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c", _, _, _) + 1 * q.recordValue(_, "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c", _, _) when: - q.record('"0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c"'.bytes, null, upstream2, null) + q.record('"0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c"'.bytes, null, upstream2) then: !q.isResolved() - 1 * q.recordValue(_, "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c", _, _, _) + 1 * q.recordValue(_, "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c", _, _) when: q.record(new JsonRpcException(1, "Nonce too low"), null, upstream3) then: - 1 * q.recordError(_, _, _, _, _) + 1 * q.recordError(_, _, _, _) q.isResolved() objectMapper.readValue(q.result, Object) == "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c" } @@ -75,18 +75,18 @@ class BroadcastQuorumSpec extends Specification { q.record(new JsonRpcException(1, "Internal error"), null, upstream1) then: !q.isResolved() - 1 * q.recordError(_, _, _, _, _) + 1 * q.recordError(_, _, _, _) when: - q.record('"0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c"'.bytes, null, upstream2, null) + q.record('"0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c"'.bytes, null, upstream2) then: !q.isResolved() - 1 * q.recordValue(_, "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c", _, _, _) + 1 * q.recordValue(_, "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c", _, _) when: q.record(new JsonRpcException(1, "Nonce too low"), null, upstream3) then: - 1 * q.recordError(_, _, _, _, _) + 1 * q.recordError(_, _, _, _) q.isResolved() objectMapper.readValue(q.result, Object) == "0xeaa972c0d8d1ecd3e34fbbef6d34e06670e745c788bdba31c4234a1762f0378c" } diff --git a/src/test/groovy/io/emeraldpay/dshackle/quorum/NotLaggingQuorumSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/quorum/NotLaggingQuorumSpec.groovy index c9930044..d3cc6021 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/quorum/NotLaggingQuorumSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/quorum/NotLaggingQuorumSpec.groovy @@ -30,7 +30,7 @@ class NotLaggingQuorumSpec extends Specification { def quorum = new NotLaggingQuorum(1) when: - quorum.record(value, null, up, null) + quorum.record(value, null, up) then: 1 * up.getLag() >> 0 quorum.isResolved() @@ -45,14 +45,13 @@ class NotLaggingQuorumSpec extends Specification { def quorum = new NotLaggingQuorum(1) when: - quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up, "test") + quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up) then: 1 * up.getLag() >> 0 quorum.isResolved() !quorum.isFailed() quorum.result == value quorum.signature == new ResponseSigner.Signature("sig1".bytes, "test", 100) - quorum.providedUpstreamId == "test" } def "Resolves if ok lag"() { @@ -62,7 +61,7 @@ class NotLaggingQuorumSpec extends Specification { def quorum = new NotLaggingQuorum(1) when: - quorum.record(value, null, up, null) + quorum.record(value, null, up) then: 1 * up.getLag() >> 1 quorum.isResolved() @@ -77,7 +76,7 @@ class NotLaggingQuorumSpec extends Specification { def quorum = new NotLaggingQuorum(1) when: - quorum.record(value, null, up, null) + quorum.record(value, null, up) then: 1 * up.getLag() >> 2 !quorum.isResolved() diff --git a/src/test/groovy/io/emeraldpay/dshackle/quorum/NotNullQuorumSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/quorum/NotNullQuorumSpec.groovy index a8b79e7c..e36d1f23 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/quorum/NotNullQuorumSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/quorum/NotNullQuorumSpec.groovy @@ -22,10 +22,10 @@ class NotNullQuorumSpec extends Specification { def quorum = new NotNullQuorum() when: - def res = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up, "id") - def res1 = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up1, "id1") - def res2 = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up2, "id2") - def res3 = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up, "id") + def res = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up) + def res1 = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up1) + def res2 = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up2) + def res3 = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up) then: !res !res1 @@ -35,7 +35,6 @@ class NotNullQuorumSpec extends Specification { !quorum.isFailed() quorum.isResolved() quorum.signature == new ResponseSigner.Signature("sig1".bytes, "test", 100) - quorum.providedUpstreamId == "id" } def "Failed if all upstreams respond with error"() { @@ -78,7 +77,7 @@ class NotNullQuorumSpec extends Specification { def quorum = new NotNullQuorum() when: - def res = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up, "id") + def res = quorum.record(value, new ResponseSigner.Signature("sig1".bytes, "test", 100), up) quorum.record(new JsonRpcException(10, "error"), new ResponseSigner.Signature("sig1".bytes, "test", 100), up1) quorum.record(new JsonRpcException(10, "error"), new ResponseSigner.Signature("sig1".bytes, "test", 100), up2) quorum.record(new JsonRpcException(10, "error"), new ResponseSigner.Signature("sig1".bytes, "test", 100), up) diff --git a/src/test/groovy/io/emeraldpay/dshackle/quorum/ValueAwareQuorumSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/quorum/ValueAwareQuorumSpec.groovy index 37ef0d1a..d008040a 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/quorum/ValueAwareQuorumSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/quorum/ValueAwareQuorumSpec.groovy @@ -66,14 +66,14 @@ class ValueAwareQuorumSpec extends Specification { } @Override - void recordValue(@NotNull byte[] response, @Nullable Object responseValue, @Nullable ResponseSigner.Signature signature, @NotNull Upstream upstream, @Nullable String providedUpstreamId) { + void recordValue(@NotNull byte[] response, @Nullable Object responseValue, @Nullable ResponseSigner.Signature signature, @NotNull Upstream upstream) { } @Override - void recordError(@Nullable byte[] response, @Nullable String errorMessage, @Nullable ResponseSigner.Signature signature, @NotNull Upstream upstream, @Nullable String providedUpstreamId) { + void recordError(@Nullable byte[] response, @Nullable String errorMessage, @Nullable ResponseSigner.Signature signature, @NotNull Upstream upstream) { } @@ -101,10 +101,5 @@ class ValueAwareQuorumSpec extends Specification { boolean isFailed() { return false } - - @Override - String getProvidedUpstreamId() { - return null - } } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy index 5f184403..0e245913 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy @@ -134,7 +134,7 @@ class NativeCallSpec extends Specification { def nativeCall = nativeCall() nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { 1 * create(_, _, _, _) >> Mock(QuorumReader) { - 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"foo\"".bytes, null, 1, Collections.singletonList(ups), null)) + 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"foo\"".bytes, null, 1, ups)) } } def call = new NativeCall.ValidCallContext(1, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum, @@ -181,7 +181,7 @@ class NativeCallSpec extends Specification { nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_test", [], 10)) >> Mono.error( - new JsonRpcException(JsonRpcResponse.Id.from(12), new JsonRpcError(-32123, "Foo Bar", "Foo Bar Baz"), true) + new JsonRpcException(JsonRpcResponse.Id.from(12), new JsonRpcError(-32123, "Foo Bar", "Foo Bar Baz"), null, true) ) } } @@ -617,7 +617,7 @@ class NativeCallSpec extends Specification { def nativeCall = nativeCall(multistreamHolder) nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { 1 * create(_, _, _, _) >> Mock(QuorumReader) { - 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"0xab\"".bytes, null, 1, Collections.singletonList(ups), null)) + 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"0xab\"".bytes, null, 1, ups)) } } def call = new NativeCall.ValidCallContext(1, 10, multistream, Selector.empty, quorum, @@ -652,7 +652,7 @@ class NativeCallSpec extends Specification { def nativeCall = nativeCall(multistreamHolder) nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) { 1 * create(_, _, _, _) >> Mock(QuorumReader) { - 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"0xab\"".bytes, null, 1, Collections.singletonList(ups), null)) + 1 * read(_) >> Mono.just(new QuorumRpcReader.Result("\"0xab\"".bytes, null, 1, ups)) } } def call = new NativeCall.ValidCallContext(1, 10, multistream, Selector.empty, quorum, diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumCachingReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumCachingReaderSpec.groovy index 0e6c1675..c50a0c01 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumCachingReaderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/EthereumCachingReaderSpec.groovy @@ -13,6 +13,7 @@ import io.emeraldpay.dshackle.upstream.ApiSource import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Multistream import io.emeraldpay.dshackle.upstream.Selector +import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.etherjar.domain.Address @@ -34,7 +35,7 @@ class EthereumDirectReaderSpec extends Specification { String hash1 = "0x40d15edaff9acdabd2a1c96fd5f683b3300aad34e7015f34def3c56ba8a7ffb5" String address1 = "0xe0aadb0a012dbcdc529c4c743d3e0385a0b54d3d" - List resolvers = Collections.singletonList((byte)1) + Upstream resolver = TestingCommons.upstream() def "Reads block by hash"() { setup: @@ -59,7 +60,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _,) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers, null) + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver) ) } } @@ -89,7 +90,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getBlockByHash", [hash1, false])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(null), null, 1, resolvers, null + Global.objectMapper.writeValueAsBytes(null), null, 1, resolver ) ) } @@ -127,7 +128,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getBlockByNumber", ["0x64", false])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers, null + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver ) ) } @@ -163,7 +164,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers, null + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver ) ) } @@ -199,7 +200,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getTransactionReceipt", [hash1])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, new ArrayList(), null + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver ) ) } @@ -236,7 +237,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getTransactionReceipt", [hash1])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, new ArrayList(), null + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver ) ) } @@ -264,7 +265,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getTransactionByHash", [hash1])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(null), null, 1, resolvers, null + Global.objectMapper.writeValueAsBytes(null), null, 1, resolver ) ) } @@ -295,7 +296,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getBalance", [address1, "latest"])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolvers, null + Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolver ) ) } @@ -327,7 +328,7 @@ class EthereumDirectReaderSpec extends Specification { 1 * create(_, _, _, _) >> Mock(QuorumReader) { 1 * read(new JsonRpcRequest("eth_getBalance", [address1, "0xa8c9bb"])) >> Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolvers, null + Global.objectMapper.writeValueAsBytes("0x100"), null, 1, resolver ) ) } @@ -359,7 +360,7 @@ class EthereumDirectReaderSpec extends Specification { } def result = Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers, null) + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver) ) EthereumDirectReader ethereumDirectReader = new EthereumDirectReader( up, Caches.default(), new CurrentBlockCache(), calls, TestingCommons.tracerMock() @@ -402,7 +403,7 @@ class EthereumDirectReaderSpec extends Specification { } def result = Mono.just( new QuorumRpcReader.Result( - Global.objectMapper.writeValueAsBytes(json), null, 1, resolvers, null) + Global.objectMapper.writeValueAsBytes(json), null, 1, resolver) ) EthereumDirectReader ethereumDirectReader = new EthereumDirectReader( up, Caches.default(), new CurrentBlockCache(), calls, TestingCommons.tracerMock()