diff --git a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt index 264294be..33e05ebc 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/startup/ConfiguredUpstreams.kt @@ -212,20 +212,22 @@ open class ConfiguredUpstreams( } log.info("Using ${chain.chainName} upstream, at ${urls.joinToString()}") + + val directApi: Reader? = buildHttpClient(config) + if (directApi == null) { + log.warn("Upstream doesn't have API configuration") + return + } + val ethereumUpstream = if (wsFactoryApi != null && !conn.preferHttp) { EthereumWsUpstream( config.id!!, - chain, wsFactoryApi, + chain, directApi, wsFactoryApi, options, config.role, QuorumForLabels.QuorumItem(1, config.labels), methods ) } else { - val directApi: Reader? = buildHttpClient(config) - if (directApi == null) { - log.warn("Upstream doesn't have API configuration") - return - } EthereumRpcUpstream( config.id!!, chain, directApi, wsFactoryApi, diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsUpstream.kt index c9ddb32b..f2d85e74 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumWsUpstream.kt @@ -24,6 +24,7 @@ import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse +import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcSwitchClient import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcWsClient import io.emeraldpay.dshackle.upstream.rpcclient.RpcMetrics import io.emeraldpay.grpc.Chain @@ -38,6 +39,7 @@ import reactor.core.Disposable class EthereumWsUpstream( id: String, val chain: Chain, + httpConnection: Reader, ethereumWsFactory: EthereumWsFactory, options: UpstreamsConfig.Options, role: UpstreamsConfig.UpstreamRole, @@ -51,7 +53,7 @@ class EthereumWsUpstream( private val head: EthereumWsHead private val connection: WsConnection - private val api: JsonRpcWsClient + private val api: Reader private var validatorSubscription: Disposable? = null private val validator: EthereumUpstreamValidator @@ -78,7 +80,12 @@ class EthereumWsUpstream( connection = ethereumWsFactory.create(this, validator, metrics) head = EthereumWsHead(connection) - api = JsonRpcWsClient(connection) + // Sometimes the server may close the WebSocket connection during the execution of a call, for example if the response + // is too large for WebSockets Frame (and Geth is unable to split messages into separate frames) + // In this case the failed request must be rerouted to the HTTP connection, because otherwise it would always fail + api = JsonRpcSwitchClient( + JsonRpcWsClient(connection), httpConnection + ) } override fun getHead(): Head { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt index 858491aa..d09cda9e 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt @@ -23,6 +23,7 @@ import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.upstream.DefaultUpstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability 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.rpcclient.ResponseWSParser @@ -35,7 +36,6 @@ import io.netty.buffer.ByteBufInputStream import io.netty.buffer.Unpooled import io.netty.handler.codec.http.HttpHeaderNames import io.netty.resolver.DefaultAddressResolverGroup -import io.netty.resolver.DefaultNameResolver import org.reactivestreams.Publisher import org.slf4j.LoggerFactory import org.springframework.util.backoff.BackOff @@ -55,6 +55,7 @@ import reactor.retry.Repeat import reactor.util.function.Tuples import java.net.URI import java.time.Duration +import java.time.Instant import java.util.Base64 import java.util.concurrent.Executors import java.util.concurrent.TimeUnit @@ -62,7 +63,7 @@ import java.util.concurrent.TimeoutException import java.util.concurrent.atomic.AtomicBoolean import java.util.concurrent.atomic.AtomicInteger -class WsConnection( +open class WsConnection( private val uri: URI, private val origin: URI, private val basicAuth: AuthConfig.ClientBasicAuth?, @@ -113,12 +114,19 @@ class WsConnection( .many() .multicast() .directBestEffort() + private val disconnects = Sinks + .many() + .multicast() + .directBestEffort() private val sendIdSeq = AtomicInteger(IDS_START) private val sendExecutor = Executors.newSingleThreadExecutor() private var keepConnection = true private var connection: Disposable? = null private val reconnecting = AtomicBoolean(false) + open val isConnected: Boolean + get() = connection != null && !reconnecting.get() + fun setReconnectIntervalSeconds(value: Long) { reconnectBackoff = FixedBackOff(value * 1000, FixedBackOff.UNLIMITED_ATTEMPTS) currentBackOff = reconnectBackoff.start() @@ -165,6 +173,7 @@ class WsConnection( connection = HttpClient.create() .resolver(DefaultAddressResolverGroup.INSTANCE) .doOnDisconnected { + disconnects.tryEmitNext(Instant.now()) log.info("Disconnected from $uri") // mark upstream as UNAVAIL upstream?.setStatus(UpstreamAvailability.UNAVAILABLE) @@ -374,25 +383,39 @@ class WsConnection( fun waitForResponse(request: JsonRpcRequest, originalId: Int, startTime: Long): Mono { val expectedId = request.id.toLong() - val failResponse = JsonRpcResponse( - null, + val noResponse = JsonRpcException( + JsonRpcResponse.Id.from(originalId), JsonRpcError( RpcResponseError.CODE_INTERNAL_ERROR, "Response not received from WebSocket" - ), - JsonRpcResponse.Id.from(originalId), null + ) ) - return Flux.from(rpcReceive.asFlux()) + val response = Flux.from(rpcReceive.asFlux()) .doOnSubscribe { sendRpc(request) } .filter { resp -> resp.id.asNumber() == expectedId } .take(Defaults.timeout) .take(1) .singleOrEmpty() + + val failOnDisconnect = Mono.from(disconnects.asFlux()) + .flatMap { + Mono.error( + JsonRpcException( + JsonRpcResponse.Id.from(originalId), + JsonRpcError( + RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR, + "Disconnected from WebSocket" + ) + ) + ) + } + + return response.or(failOnDisconnect) .doOnNext { rpcMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) } .doOnError { rpcMetrics?.fails?.increment() } .map { it.copyWithId(JsonRpcResponse.Id.from(originalId)) } - .defaultIfEmpty(failResponse) + .switchIfEmpty(Mono.error(noResponse)) } fun getBlocksFlux(): Flux { 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 fb425ebf..bc56fc41 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt @@ -112,6 +112,7 @@ class JsonRpcHttpClient( return Mono.just(key) .map(JsonRpcRequest::toJson) .doOnNext { + println("rpc connection ${key.method}") startTime = System.nanoTime() } .flatMap(this@JsonRpcHttpClient::execute) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcSwitchClient.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcSwitchClient.kt new file mode 100644 index 00000000..82ab68cd --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcSwitchClient.kt @@ -0,0 +1,27 @@ +package io.emeraldpay.dshackle.upstream.rpcclient + +import io.emeraldpay.dshackle.reader.Reader +import org.slf4j.LoggerFactory +import reactor.core.publisher.Mono + +/** + * An aggregating JSON RPC Client that wraps two actual readers, a Primary and a Secondary. + * It always calls the Primary reader, and if it fails or produces an empty result, then it calls the Secondary reader. + */ +class JsonRpcSwitchClient( + private val primary: Reader, + private val secondary: Reader, +) : Reader { + + companion object { + private val log = LoggerFactory.getLogger(JsonRpcSwitchClient::class.java) + } + + override fun read(key: JsonRpcRequest): Mono { + return primary.read(key) + .switchIfEmpty(Mono.error(IllegalStateException("No response from Primary Connection"))) + .onErrorResume { + secondary.read(key) + } + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClient.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClient.kt index 0a10a868..6f65d0a2 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClient.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClient.kt @@ -17,6 +17,7 @@ package io.emeraldpay.dshackle.upstream.rpcclient import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.ethereum.WsConnection +import io.emeraldpay.etherjar.rpc.RpcResponseError import reactor.core.publisher.Mono class JsonRpcWsClient( @@ -24,6 +25,17 @@ class JsonRpcWsClient( ) : Reader { override fun read(key: JsonRpcRequest): Mono { + if (!ws.isConnected) { + return Mono.error( + JsonRpcException( + JsonRpcResponse.NumberId(key.id), + JsonRpcError( + RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR, + "WebSocket is not connected" + ) + ) + ) + } return ws.call(key) } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionRealSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionRealSpec.groovy index a84e6e00..9f16a8a2 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionRealSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionRealSpec.groovy @@ -93,6 +93,20 @@ class WsConnectionRealSpec extends Specification { act[0].value.contains("\"params\":[\"newHeads\"]") } + def "Error on request when server disconnects"() { + when: + conn.connect() + conn.reconnectIntervalSeconds = 2 + + def resp = conn.call(new JsonRpcRequest("foo_bar", [])) + + then: + StepVerifier.create(resp) + .then { server.stop() } + .expectError() + .verify(Duration.ofSeconds(1)) + } + def "Gets UNAVAIL status right after disconnect"() { setup: def up = Mock(DefaultUpstream) diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcSwitchClientSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcSwitchClientSpec.groovy new file mode 100644 index 00000000..883e400b --- /dev/null +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcSwitchClientSpec.groovy @@ -0,0 +1,70 @@ +package io.emeraldpay.dshackle.upstream.rpcclient + +import io.emeraldpay.dshackle.reader.Reader +import reactor.core.publisher.Mono +import spock.lang.Specification + +import java.time.Duration + +class JsonRpcSwitchClientSpec extends Specification { + + def "Uses primary response if it works"() { + setup: + def primaryCalled = false + def secondaryCalled = false + def request = new JsonRpcRequest("eth_test", []) + def response = JsonRpcResponse.ok("test".bytes, new JsonRpcResponse.NumberId(100)) + def primary = Mock(Reader) { + 1 * read(request) >> Mono.fromCallable { + primaryCalled = true + response + } + } + def secondary = Mock(Reader) { + _ * read(request) >> Mono.fromCallable { + secondaryCalled = true + response + } + } + + def client = new JsonRpcSwitchClient(primary, secondary) + + when: + def act = client.read(request).block(Duration.ofSeconds(1)) + + then: + act == response + primaryCalled + !secondaryCalled + } + + def "Uses secondary response if primary fails"() { + setup: + def primaryCalled = false + def secondaryCalled = false + def request = new JsonRpcRequest("eth_test", []) + def response = JsonRpcResponse.ok("test".bytes, new JsonRpcResponse.NumberId(100)) + def primary = Mock(Reader) { + 1 * read(request) >> Mono.fromCallable { + primaryCalled = true + throw new IllegalStateException("Primary Fail") + } + } + def secondary = Mock(Reader) { + 1 * read(request) >> Mono.fromCallable { + secondaryCalled = true + response + } + } + + def client = new JsonRpcSwitchClient(primary, secondary) + + when: + def act = client.read(request).block(Duration.ofSeconds(1)) + + then: + act == response + primaryCalled + secondaryCalled + } +} diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClientSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClientSpec.groovy new file mode 100644 index 00000000..e0ac8a1b --- /dev/null +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcWsClientSpec.groovy @@ -0,0 +1,23 @@ +package io.emeraldpay.dshackle.upstream.rpcclient + +import io.emeraldpay.dshackle.upstream.ethereum.WsConnection +import reactor.core.Exceptions +import spock.lang.Specification + +import java.time.Duration + +class JsonRpcWsClientSpec extends Specification { + + def "Produce error if WS is not connected"() { + setup: + def ws = Mock(WsConnection) + def client = new JsonRpcWsClient(ws) + when: + client.read(new JsonRpcRequest("foo_bar", [], 1)) + .block(Duration.ofSeconds(1)) + then: + def t = thrown(Exceptions.ReactiveException) + t.cause instanceof JsonRpcException + 1 * ws.isConnected() >> false + } +}