propagate grpc provided upstream id
This commit is contained in:
@@ -29,6 +29,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<Upstream> = ConcurrentLinkedQueue()
|
||||
|
||||
override fun init(head: Head) {
|
||||
@@ -46,15 +47,29 @@ open class AlwaysQuorum : CallQuorum {
|
||||
return sig
|
||||
}
|
||||
|
||||
override fun record(response: ByteArray, signature: ResponseSigner.Signature?, upstream: Upstream): Boolean {
|
||||
override fun getProvidedUpstreamId(): String? {
|
||||
return providedUpstreamId
|
||||
}
|
||||
|
||||
override fun record(
|
||||
response: ByteArray,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
): Boolean {
|
||||
result = response
|
||||
resolved = true
|
||||
sig = signature
|
||||
resolvers.add(upstream)
|
||||
this.providedUpstreamId = providedUpstreamId
|
||||
return true
|
||||
}
|
||||
|
||||
override fun record(error: JsonRpcException, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun record(
|
||||
error: JsonRpcException,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream
|
||||
) {
|
||||
this.rpcError = error.error
|
||||
sig = signature
|
||||
}
|
||||
|
||||
@@ -28,6 +28,7 @@ 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) {
|
||||
}
|
||||
@@ -48,21 +49,39 @@ open class BroadcastQuorum(
|
||||
return sig
|
||||
}
|
||||
|
||||
override fun recordValue(response: ByteArray, responseValue: String?, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun getProvidedUpstreamId(): String? {
|
||||
return providedUpstreamId
|
||||
}
|
||||
|
||||
override fun recordValue(
|
||||
response: ByteArray,
|
||||
responseValue: String?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
) {
|
||||
calls++
|
||||
if (txid == null && responseValue != null) {
|
||||
txid = responseValue
|
||||
sig = signature
|
||||
result = response
|
||||
this.providedUpstreamId = providedUpstreamId
|
||||
}
|
||||
}
|
||||
|
||||
override fun recordError(response: ByteArray?, errorMessage: String?, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun recordError(
|
||||
response: ByteArray?,
|
||||
errorMessage: String?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
) {
|
||||
// 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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -29,9 +29,20 @@ interface CallQuorum {
|
||||
fun isResolved(): Boolean
|
||||
fun isFailed(): Boolean
|
||||
|
||||
fun record(response: ByteArray, signature: ResponseSigner.Signature?, upstream: Upstream): Boolean
|
||||
fun record(error: JsonRpcException, signature: ResponseSigner.Signature?, upstream: Upstream)
|
||||
fun record(
|
||||
response: ByteArray,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
): Boolean
|
||||
fun record(
|
||||
error: JsonRpcException,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream
|
||||
)
|
||||
fun getSignature(): ResponseSigner.Signature?
|
||||
|
||||
fun getProvidedUpstreamId(): String?
|
||||
fun getResult(): ByteArray?
|
||||
fun getError(): JsonRpcError?
|
||||
fun getResolvedBy(): Collection<Upstream>
|
||||
|
||||
@@ -27,6 +27,7 @@ open class NonEmptyQuorum(
|
||||
private var result: ByteArray? = null
|
||||
private var tries: Int = 0
|
||||
private var sig: ResponseSigner.Signature? = null
|
||||
private var providedUpstreamId: String? = null
|
||||
override fun init(head: Head) {
|
||||
}
|
||||
|
||||
@@ -42,11 +43,22 @@ open class NonEmptyQuorum(
|
||||
return sig
|
||||
}
|
||||
|
||||
override fun recordValue(response: ByteArray, responseValue: Any?, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun getProvidedUpstreamId(): String? {
|
||||
return providedUpstreamId
|
||||
}
|
||||
|
||||
override fun recordValue(
|
||||
response: ByteArray,
|
||||
responseValue: Any?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
) {
|
||||
tries++
|
||||
if (responseValue != null) {
|
||||
result = response
|
||||
sig = signature
|
||||
this.providedUpstreamId = providedUpstreamId
|
||||
}
|
||||
}
|
||||
|
||||
@@ -54,7 +66,13 @@ open class NonEmptyQuorum(
|
||||
return result
|
||||
}
|
||||
|
||||
override fun recordError(response: ByteArray?, errorMessage: String?, sig: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun recordError(
|
||||
response: ByteArray?,
|
||||
errorMessage: String?,
|
||||
sig: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
) {
|
||||
tries++
|
||||
}
|
||||
|
||||
|
||||
@@ -33,6 +33,7 @@ open class NonceQuorum(
|
||||
private var receivedTimes = 0
|
||||
private var errors = 0
|
||||
private var sig: ResponseSigner.Signature? = null
|
||||
private var providedUpstreamId: String? = null
|
||||
|
||||
override fun init(head: Head) {
|
||||
}
|
||||
@@ -50,7 +51,17 @@ open class NonceQuorum(
|
||||
return sig
|
||||
}
|
||||
|
||||
override fun recordValue(response: ByteArray, responseValue: String?, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun getProvidedUpstreamId(): String? {
|
||||
return providedUpstreamId
|
||||
}
|
||||
|
||||
override fun recordValue(
|
||||
response: ByteArray,
|
||||
responseValue: String?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
) {
|
||||
val value = responseValue?.let { str ->
|
||||
HexQuantity.from(str).value.toLong()
|
||||
}
|
||||
@@ -60,9 +71,11 @@ open class NonceQuorum(
|
||||
resultValue = value
|
||||
result = response
|
||||
sig = signature
|
||||
this.providedUpstreamId = providedUpstreamId
|
||||
} else if (result == null) {
|
||||
result = response
|
||||
sig = signature
|
||||
this.providedUpstreamId = providedUpstreamId
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -71,7 +84,13 @@ open class NonceQuorum(
|
||||
return result
|
||||
}
|
||||
|
||||
override fun recordError(response: ByteArray?, errorMessage: String?, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun recordError(
|
||||
response: ByteArray?,
|
||||
errorMessage: String?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
) {
|
||||
errors++
|
||||
}
|
||||
|
||||
|
||||
@@ -36,6 +36,7 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum {
|
||||
private var rpcError: JsonRpcError? = null
|
||||
private var sig: ResponseSigner.Signature? = null
|
||||
private val resolvers: MutableCollection<Upstream> = ConcurrentLinkedQueue()
|
||||
private var providedUpstreamId: String? = null
|
||||
|
||||
override fun init(head: Head) {
|
||||
}
|
||||
@@ -48,18 +49,28 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum {
|
||||
return failed.get()
|
||||
}
|
||||
|
||||
override fun record(response: ByteArray, signature: ResponseSigner.Signature?, upstream: Upstream): Boolean {
|
||||
override fun record(
|
||||
response: ByteArray,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
): Boolean {
|
||||
val lagging = upstream.getLag() > maxLag
|
||||
if (!lagging) {
|
||||
result.set(response)
|
||||
sig = signature
|
||||
this.providedUpstreamId = providedUpstreamId
|
||||
resolvers.add(upstream)
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
override fun record(error: JsonRpcException, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun record(
|
||||
error: JsonRpcException,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream
|
||||
) {
|
||||
this.rpcError = error.error
|
||||
val lagging = upstream.getLag() > maxLag
|
||||
if (!lagging && result.get() == null) {
|
||||
@@ -70,6 +81,11 @@ class NotLaggingQuorum(val maxLag: Long = 0) : CallQuorum {
|
||||
override fun getSignature(): ResponseSigner.Signature? {
|
||||
return sig
|
||||
}
|
||||
|
||||
override fun getProvidedUpstreamId(): String? {
|
||||
return providedUpstreamId
|
||||
}
|
||||
|
||||
override fun getResult(): ByteArray {
|
||||
return result.get()
|
||||
}
|
||||
|
||||
@@ -27,8 +27,8 @@ import io.emeraldpay.etherjar.rpc.RpcException
|
||||
import org.slf4j.LoggerFactory
|
||||
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
|
||||
@@ -90,8 +90,8 @@ class QuorumRpcReader(
|
||||
}
|
||||
|
||||
fun execute(key: JsonRpcRequest, retrySpec: reactor.util.retry.Retry): Function<Flux<Upstream>, Mono<CallQuorum>> {
|
||||
val quorumReduce = BiFunction<CallQuorum, Tuple3<ByteArray, Optional<ResponseSigner.Signature>, Upstream>, CallQuorum> { res, a ->
|
||||
if (res.record(a.t1, a.t2.orElse(null), a.t3)) {
|
||||
val quorumReduce = BiFunction<CallQuorum, Tuple4<ByteArray, Optional<ResponseSigner.Signature>, Upstream, Optional<String>>, CallQuorum> { res, a ->
|
||||
if (res.record(a.t1, a.t2.orElse(null), a.t3, a.t4.orElse(null))) {
|
||||
apiControl.resolve()
|
||||
} else {
|
||||
// quorum needs more responses, so ask api controller to make another
|
||||
@@ -119,25 +119,25 @@ 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())
|
||||
Result(quorum.getResult()!!, quorum.getSignature(), 1, quorum.getResolvedBy(), quorum.getProvidedUpstreamId())
|
||||
}
|
||||
.switchIfEmpty(defaultResult)
|
||||
}
|
||||
}
|
||||
|
||||
fun callApi(api: Upstream, key: JsonRpcRequest): Mono<Tuple3<ByteArray, Optional<ResponseSigner.Signature>, Upstream>> {
|
||||
fun callApi(api: Upstream, key: JsonRpcRequest): Mono<Tuple4<ByteArray, Optional<ResponseSigner.Signature>, Upstream, Optional<String>>> {
|
||||
return api.getIngressReader()
|
||||
.read(key)
|
||||
.flatMap { response ->
|
||||
response.requireResult()
|
||||
.transform(withSignature(api, key, response))
|
||||
.transform(withSignatureAndUpstream(api, key, response))
|
||||
}
|
||||
// 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) }
|
||||
.map { Tuples.of(it.t1, it.t2, api, it.t3) }
|
||||
}
|
||||
|
||||
fun withSignature(api: Upstream, key: JsonRpcRequest, response: JsonRpcResponse): Function<Mono<ByteArray>, Mono<Tuple2<ByteArray, Optional<ResponseSigner.Signature>>>> {
|
||||
fun withSignatureAndUpstream(api: Upstream, key: JsonRpcRequest, response: JsonRpcResponse): Function<Mono<ByteArray>, Mono<Tuple3<ByteArray, Optional<ResponseSigner.Signature>, Optional<String>>>> {
|
||||
return Function { src ->
|
||||
src.map {
|
||||
val signature = response.providedSignature
|
||||
@@ -146,7 +146,7 @@ class QuorumRpcReader(
|
||||
} else {
|
||||
null
|
||||
}
|
||||
Tuples.of(it, Optional.ofNullable(signature))
|
||||
Tuples.of(it, Optional.ofNullable(signature), Optional.ofNullable(response.providedUpstreamId))
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -165,7 +165,7 @@ class QuorumRpcReader(
|
||||
JsonRpcError(-32603, "Unhandled internal error: ${err.javaClass}: ${err.message}")
|
||||
)
|
||||
}
|
||||
quorum.record(cleanErr, null, api)
|
||||
quorum.record(cleanErr, null, api,)
|
||||
// if it's failed after that, then we don't need more calls, stop api source
|
||||
if (quorum.isFailed()) {
|
||||
apiControl.resolve()
|
||||
@@ -198,6 +198,7 @@ class QuorumRpcReader(
|
||||
val value: ByteArray,
|
||||
val signature: ResponseSigner.Signature?,
|
||||
val quorum: Int,
|
||||
val resolvers: Collection<Upstream>
|
||||
val resolvers: Collection<Upstream>,
|
||||
val providedUpstreamId: String?
|
||||
)
|
||||
}
|
||||
|
||||
@@ -37,27 +37,48 @@ abstract class ValueAwareQuorum<T>(
|
||||
return Global.objectMapper.readValue(response.inputStream(), clazz)
|
||||
}
|
||||
|
||||
override fun record(response: ByteArray, signature: ResponseSigner.Signature?, upstream: Upstream): Boolean {
|
||||
override fun record(
|
||||
response: ByteArray,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
): Boolean {
|
||||
try {
|
||||
val value = extractValue(response, clazz)
|
||||
recordValue(response, value, signature, upstream)
|
||||
recordValue(response, value, signature, upstream, providedUpstreamId)
|
||||
resolvers.add(upstream)
|
||||
} catch (e: RpcException) {
|
||||
recordError(response, e.rpcMessage, signature, upstream)
|
||||
recordError(response, e.rpcMessage, signature, upstream, providedUpstreamId)
|
||||
} catch (e: Exception) {
|
||||
recordError(response, e.message, signature, upstream)
|
||||
recordError(response, e.message, signature, upstream, providedUpstreamId)
|
||||
}
|
||||
return isResolved()
|
||||
}
|
||||
|
||||
override fun record(error: JsonRpcException, signature: ResponseSigner.Signature?, upstream: Upstream) {
|
||||
override fun record(
|
||||
error: JsonRpcException,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream
|
||||
) {
|
||||
this.rpcError = error.error
|
||||
recordError(null, error.error.message, signature, upstream)
|
||||
recordError(null, error.error.message, signature, upstream, null)
|
||||
}
|
||||
|
||||
abstract fun recordValue(response: ByteArray, responseValue: T?, signature: ResponseSigner.Signature?, upstream: Upstream)
|
||||
abstract fun recordValue(
|
||||
response: ByteArray,
|
||||
responseValue: T?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
)
|
||||
|
||||
abstract fun recordError(response: ByteArray?, errorMessage: String?, signature: ResponseSigner.Signature?, upstream: Upstream)
|
||||
abstract fun recordError(
|
||||
response: ByteArray?,
|
||||
errorMessage: String?,
|
||||
signature: ResponseSigner.Signature?,
|
||||
upstream: Upstream,
|
||||
providedUpstreamId: String?
|
||||
)
|
||||
|
||||
override fun getError(): JsonRpcError? {
|
||||
return rpcError
|
||||
|
||||
@@ -312,7 +312,7 @@ open class NativeCall(
|
||||
.map {
|
||||
val bytes = ctx.resultDecorator.processResult(it)
|
||||
validateResult(bytes, "remote", ctx)
|
||||
CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, it.resolvers.first().getId(), ctx)
|
||||
CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, it.providedUpstreamId ?: it.resolvers.first().getId(), ctx)
|
||||
}
|
||||
.onErrorResume { t ->
|
||||
Mono.just(CallResult.fail(ctx.id, ctx.nonce, t, ctx))
|
||||
|
||||
@@ -90,7 +90,7 @@ class JsonRpcGrpcClient(
|
||||
} else {
|
||||
null
|
||||
}
|
||||
Mono.just(JsonRpcResponse(bytes, null, JsonRpcResponse.NumberId(0), signature))
|
||||
Mono.just(JsonRpcResponse(bytes, null, JsonRpcResponse.NumberId(0), signature, resp.upstreamId))
|
||||
} else {
|
||||
metrics?.fails?.increment()
|
||||
Mono.error(
|
||||
|
||||
@@ -29,7 +29,8 @@ class JsonRpcResponse(
|
||||
/**
|
||||
* When making a request through Dshackle protocol a remote may provide its signature with the response, which we keep here
|
||||
*/
|
||||
val providedSignature: ResponseSigner.Signature? = null
|
||||
val providedSignature: ResponseSigner.Signature? = null,
|
||||
val providedUpstreamId: String? = null
|
||||
) {
|
||||
|
||||
constructor(result: ByteArray?, error: JsonRpcError?) : this(result, error, NumberId(0))
|
||||
@@ -113,11 +114,7 @@ class JsonRpcResponse(
|
||||
}
|
||||
|
||||
fun copyWithId(id: Id): JsonRpcResponse {
|
||||
return JsonRpcResponse(result, error, id, providedSignature)
|
||||
}
|
||||
|
||||
fun copyWithSignature(signature: ResponseSigner.Signature): JsonRpcResponse {
|
||||
return JsonRpcResponse(result, error, id, signature)
|
||||
return JsonRpcResponse(result, error, id, providedSignature, providedUpstreamId)
|
||||
}
|
||||
|
||||
override fun equals(other: Any?): Boolean {
|
||||
|
||||
Reference in New Issue
Block a user