problem: no response for some requests, possible bug in Quorum selector
solution: refactor Quorum selector to make it easier to debug; no bug, issue on upstream side
This commit is contained in:
@@ -17,20 +17,25 @@ package io.emeraldpay.dshackle.quorum
|
|||||||
|
|
||||||
import io.emeraldpay.dshackle.reader.Reader
|
import io.emeraldpay.dshackle.reader.Reader
|
||||||
import io.emeraldpay.dshackle.upstream.ApiSource
|
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.JsonRpcException
|
||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
|
||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
|
||||||
import io.emeraldpay.etherjar.rpc.RpcException
|
import io.emeraldpay.etherjar.rpc.RpcException
|
||||||
|
import java.util.function.BiFunction
|
||||||
|
import java.util.function.Function
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
import reactor.core.publisher.Flux
|
import reactor.core.publisher.Flux
|
||||||
import reactor.core.publisher.Mono
|
import reactor.core.publisher.Mono
|
||||||
|
import reactor.util.function.Tuple2
|
||||||
import reactor.util.function.Tuples
|
import reactor.util.function.Tuples
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Makes request with applying Quorum
|
* Makes request with applying Quorum
|
||||||
*/
|
*/
|
||||||
class QuorumRpcReader(
|
class QuorumRpcReader(
|
||||||
private val apis: ApiSource,
|
private val apiControl: ApiSource,
|
||||||
private val quorum: CallQuorum
|
private val quorum: CallQuorum
|
||||||
) : Reader<JsonRpcRequest, QuorumRpcReader.Result> {
|
) : Reader<JsonRpcRequest, QuorumRpcReader.Result> {
|
||||||
|
|
||||||
@@ -38,8 +43,9 @@ class QuorumRpcReader(
|
|||||||
private val log = LoggerFactory.getLogger(QuorumRpcReader::class.java)
|
private val log = LoggerFactory.getLogger(QuorumRpcReader::class.java)
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun read(key: JsonRpcRequest): Mono<QuorumRpcReader.Result> {
|
override fun read(key: JsonRpcRequest): Mono<Result> {
|
||||||
apis.request(1)
|
// needs at least one response, so start a request
|
||||||
|
apiControl.request(1)
|
||||||
|
|
||||||
// uses a mix of retry strategy and managed Publisher for calls.
|
// uses a mix of retry strategy and managed Publisher for calls.
|
||||||
// retry is used when an error happened
|
// retry is used when an error happened
|
||||||
@@ -50,66 +56,17 @@ class QuorumRpcReader(
|
|||||||
signal.takeUntil {
|
signal.takeUntil {
|
||||||
it.totalRetries() >= 3 || quorum.isResolved() || quorum.isFailed()
|
it.totalRetries() >= 3 || quorum.isResolved() || quorum.isFailed()
|
||||||
}.doOnNext {
|
}.doOnNext {
|
||||||
// need one more API source if retried
|
// when retried it needs one more API source
|
||||||
apis.request(1)
|
apiControl.request(1)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
val defaultResult: Mono<Result> = Mono.just(quorum).flatMap { q ->
|
val defaultResult: Mono<Result> = setupDefaultResult(key)
|
||||||
if (q.isFailed()) {
|
|
||||||
Mono.error<Result>(
|
|
||||||
q.getError()?.asException(JsonRpcResponse.NumberId(1))
|
|
||||||
?: RpcException(-32000, "Unknown Upstream error")
|
|
||||||
)
|
|
||||||
} else {
|
|
||||||
log.warn("Did not get any result from upstream. Method [${key.method}] using [$q]")
|
|
||||||
Mono.empty<Result>()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return Flux.from(apis)
|
return Flux.from(apiControl)
|
||||||
.takeUntil {
|
.transform(execute(key, retrySpec))
|
||||||
quorum.isFailed() || quorum.isResolved()
|
.next()
|
||||||
}
|
// if last call resulted in error it's still possible that request was resolved correctly. ex. for BroadcastQuorum
|
||||||
.flatMap { api ->
|
|
||||||
api.getApi()
|
|
||||||
.read(key)
|
|
||||||
.flatMap { response ->
|
|
||||||
response.requireResult()
|
|
||||||
.onErrorResume { err ->
|
|
||||||
if (err is RpcException || err is JsonRpcException) {
|
|
||||||
// on error notify quorum, it may use error message or other details
|
|
||||||
val cleanErr: JsonRpcException = when (err) {
|
|
||||||
is RpcException -> JsonRpcException.from(err)
|
|
||||||
is JsonRpcException -> err
|
|
||||||
else -> throw IllegalStateException("Cannot convert from exception", err)
|
|
||||||
}
|
|
||||||
quorum.record(cleanErr, api)
|
|
||||||
// it it's failed after that, then we don't need more calls, stop api source
|
|
||||||
if (quorum.isFailed()) {
|
|
||||||
apis.resolve()
|
|
||||||
} else {
|
|
||||||
apis.request(1)
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
log.warn("Result processing error", err)
|
|
||||||
}
|
|
||||||
Mono.empty()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
.map { Tuples.of(it, api) }
|
|
||||||
}
|
|
||||||
.retryWhen(retrySpec)
|
|
||||||
// record all correct responses until quorum reached
|
|
||||||
.reduce(quorum, { res, a ->
|
|
||||||
if (res.record(a.t1, a.t2)) {
|
|
||||||
apis.resolve()
|
|
||||||
} else {
|
|
||||||
apis.request(1)
|
|
||||||
}
|
|
||||||
res
|
|
||||||
})
|
|
||||||
// if last call resulted in error it's still possible that request was resolved correctly. i.e. for BroadcastQuorum
|
|
||||||
.onErrorResume { err ->
|
.onErrorResume { err ->
|
||||||
if (quorum.isResolved()) {
|
if (quorum.isResolved()) {
|
||||||
Mono.just(quorum)
|
Mono.just(quorum)
|
||||||
@@ -122,13 +79,85 @@ class QuorumRpcReader(
|
|||||||
log.debug("No quorum for ${key.method} using [$quorum]. Error: ${it.getError()?.message ?: ""}")
|
log.debug("No quorum for ${key.method} using [$quorum]. Error: ${it.getError()?.message ?: ""}")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// return nothing if not resolved
|
.transform(processResult(defaultResult))
|
||||||
.filter { it.isResolved() }
|
}
|
||||||
.map {
|
|
||||||
// TODO find actual quorum number
|
fun execute(key: JsonRpcRequest, retrySpec: reactor.util.retry.Retry): Function<Flux<Upstream>, Mono<CallQuorum>> {
|
||||||
QuorumRpcReader.Result(it.getResult()!!, 1)
|
val quorumReduce = BiFunction<CallQuorum, Tuple2<ByteArray, Upstream>, CallQuorum> { res, a ->
|
||||||
|
if (res.record(a.t1, a.t2)) {
|
||||||
|
apiControl.resolve()
|
||||||
|
} else {
|
||||||
|
// quorum needs more responses, so ask api controller to make another
|
||||||
|
apiControl.request(1)
|
||||||
}
|
}
|
||||||
.switchIfEmpty(defaultResult)
|
res
|
||||||
|
}
|
||||||
|
return Function { apiFlux ->
|
||||||
|
apiFlux
|
||||||
|
.takeUntil {
|
||||||
|
quorum.isFailed() || quorum.isResolved()
|
||||||
|
}
|
||||||
|
.flatMap { api ->
|
||||||
|
callApi(api, key)
|
||||||
|
}
|
||||||
|
.retryWhen(retrySpec)
|
||||||
|
// record all correct responses until quorum reached
|
||||||
|
.reduce(quorum, quorumReduce)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun processResult(defaultResult: Mono<Result>): Function<Mono<CallQuorum>, Mono<Result>> {
|
||||||
|
return Function { quorumResult ->
|
||||||
|
quorumResult
|
||||||
|
.filter { it.isResolved() } // return nothing if not resolved
|
||||||
|
.map {
|
||||||
|
// TODO find actual quorum number
|
||||||
|
QuorumRpcReader.Result(it.getResult()!!, 1)
|
||||||
|
}
|
||||||
|
.switchIfEmpty(defaultResult)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun callApi(api: Upstream, key: JsonRpcRequest): Mono<Tuple2<ByteArray, Upstream>> {
|
||||||
|
return api.getApi()
|
||||||
|
.read(key)
|
||||||
|
.flatMap { response ->
|
||||||
|
response.requireResult()
|
||||||
|
.onErrorResume { err ->
|
||||||
|
// on error notify quorum, it may use 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}")
|
||||||
|
)
|
||||||
|
}
|
||||||
|
quorum.record(cleanErr, api)
|
||||||
|
// if it's failed after that, then we don't need more calls, stop api source
|
||||||
|
if (quorum.isFailed()) {
|
||||||
|
apiControl.resolve()
|
||||||
|
} else {
|
||||||
|
apiControl.request(1)
|
||||||
|
}
|
||||||
|
Mono.empty()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
.map { Tuples.of(it, api) }
|
||||||
|
}
|
||||||
|
|
||||||
|
fun setupDefaultResult(key: JsonRpcRequest): Mono<Result> {
|
||||||
|
return Mono.just(quorum).flatMap { q ->
|
||||||
|
if (q.isFailed()) {
|
||||||
|
Mono.error<Result>(
|
||||||
|
q.getError()?.asException(JsonRpcResponse.NumberId(key.id))
|
||||||
|
?: JsonRpcException(JsonRpcResponse.NumberId(key.id), JsonRpcError(-32603, "Unhandled Upstream error"))
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
log.warn("Did not get any result from upstream. Method [${key.method}] using [$q]")
|
||||||
|
Mono.empty<Result>()
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class Result(
|
class Result(
|
||||||
|
|||||||
@@ -232,10 +232,10 @@ open class NativeCall(
|
|||||||
CallResult(ctx.id, it.value, null)
|
CallResult(ctx.id, it.value, null)
|
||||||
}
|
}
|
||||||
.onErrorResume { t ->
|
.onErrorResume { t ->
|
||||||
val failure = if (t is CallFailure) {
|
val failure = when (t) {
|
||||||
CallResult.fail(t.id, t.reason)
|
is CallFailure -> CallResult.fail(t.id, t.reason)
|
||||||
} else {
|
is JsonRpcException -> CallResult.fail(ctx.id, t.error.code, t.error.message)
|
||||||
CallResult.fail(ctx.id, t)
|
else -> CallResult.fail(ctx.id, t)
|
||||||
}
|
}
|
||||||
Mono.just(failure)
|
Mono.just(failure)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -221,6 +221,39 @@ class QuorumRpcReaderSpec extends Specification {
|
|||||||
.verify(Duration.ofSeconds(2))
|
.verify(Duration.ofSeconds(2))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
def "always-quorum - error if failed"() {
|
||||||
|
setup:
|
||||||
|
def api = Mock(Reader) {
|
||||||
|
1 * read(new JsonRpcRequest("eth_test", [])) >>> [
|
||||||
|
Mono.just(JsonRpcResponse.error(1, "test error")),
|
||||||
|
]
|
||||||
|
}
|
||||||
|
def up = Mock(Upstream) {
|
||||||
|
_ * isAvailable() >> true
|
||||||
|
_ * getRole() >> UpstreamsConfig.UpstreamRole.PRIMARY
|
||||||
|
_ * getApi() >> api
|
||||||
|
}
|
||||||
|
def apis = new FilteredApis(
|
||||||
|
Chain.ETHEREUM,
|
||||||
|
[up], Selector.empty
|
||||||
|
)
|
||||||
|
def reader = new QuorumRpcReader(apis, new AlwaysQuorum())
|
||||||
|
|
||||||
|
when:
|
||||||
|
def act = reader.read(new JsonRpcRequest("eth_test", []))
|
||||||
|
.map {
|
||||||
|
new String(it.value)
|
||||||
|
}
|
||||||
|
|
||||||
|
then:
|
||||||
|
StepVerifier.create(act)
|
||||||
|
.expectErrorMatches { t ->
|
||||||
|
println("Error: $t.class / $t.message")
|
||||||
|
t instanceof JsonRpcException && t.message == "test error" && t.error.code == 1
|
||||||
|
}
|
||||||
|
.verify(Duration.ofSeconds(2))
|
||||||
|
}
|
||||||
|
|
||||||
def "Return error is upstream returned it"() {
|
def "Return error is upstream returned it"() {
|
||||||
setup:
|
setup:
|
||||||
def up = Mock(Upstream) {
|
def up = Mock(Upstream) {
|
||||||
|
|||||||
Reference in New Issue
Block a user