From 90a8f8ccab09f9229cdc9da786603b6782235aec Mon Sep 17 00:00:00 2001 From: KirillPamPam Date: Fri, 12 Apr 2024 14:53:57 +0400 Subject: [PATCH] Change ws timeouts (#452) --- .../upstream/ethereum/GenericWsHead.kt | 13 ++- .../generic/connectors/GenericRpcConnector.kt | 2 + .../generic/connectors/GenericWsConnector.kt | 1 + .../ethereum/GenericWsHeadSpec.groovy | 96 +++++++++++++++++-- 4 files changed, 103 insertions(+), 9 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHead.kt index f68fb6d2..7d9c8f76 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHead.kt @@ -42,7 +42,18 @@ class GenericWsHead( headScheduler: Scheduler, upstream: DefaultUpstream, private val chainSpecific: ChainSpecific, + timeout: Duration, ) : GenericHead(upstream.getId(), forkChoice, blockValidator, headScheduler, chainSpecific), Lifecycle { + private val wsHeadTimeout = run { + val defaultTimeout = Duration.ofMinutes(1) + if (timeout >= defaultTimeout) { + timeout.plus(defaultTimeout) + } else { + defaultTimeout + } + }.also { + log.info("WS head timeout for ${upstream.getId()} is $it") + } private var connectionId: String? = null private var subscribed = false @@ -92,7 +103,7 @@ class GenericWsHead( .map { chainSpecific.parseHeader(it, "unknown") } - .timeout(Duration.ofSeconds(60), Mono.error(RuntimeException("No response from subscribe to newHeads"))) + .timeout(wsHeadTimeout, Mono.error(RuntimeException("No response from subscribe to newHeads"))) .onErrorResume { log.error("Error getting heads for $upstreamId", it) subscribed = false diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericRpcConnector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericRpcConnector.kt index ef5d53e5..93ffaecb 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericRpcConnector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericRpcConnector.kt @@ -93,6 +93,7 @@ class GenericRpcConnector( headScheduler, upstream, chainSpecific, + expectedBlockTime, ) // receive all new blocks through WebSockets, but also periodically verify with RPC in case if WS failed val rpcHead = @@ -118,6 +119,7 @@ class GenericRpcConnector( headScheduler, upstream, chainSpecific, + expectedBlockTime, ) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericWsConnector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericWsConnector.kt index 8ff1caa5..e7d25052 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericWsConnector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/connectors/GenericWsConnector.kt @@ -46,6 +46,7 @@ class GenericWsConnector( headScheduler, upstream, chainSpecific, + expectedBlockTime, ) liveness = HeadLivenessValidator(head, expectedBlockTime, headLivenessScheduler, upstream.getId()) subscriptions = chainSpecific.makeIngressSubscription(wsSubscriptions) diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHeadSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHeadSpec.groovy index 60c6c9d7..dc0955a9 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHeadSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/GenericWsHeadSpec.groovy @@ -73,7 +73,17 @@ class GenericWsHeadSpec extends Specification { Flux.fromIterable([headBlock]), "id", new AtomicReference("") ) - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, reader, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + reader, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) def res = BlockContainer.from(block) when: @@ -112,7 +122,17 @@ class GenericWsHeadSpec extends Specification { ] } - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, apiMock, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + apiMock, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: def act = head.getFlux() @@ -166,7 +186,17 @@ class GenericWsHeadSpec extends Specification { ] } - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, apiMock, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + apiMock, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: def act = head.getFlux() @@ -206,7 +236,17 @@ class GenericWsHeadSpec extends Specification { ] } - def head = new GenericWsHead( new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, apiMock, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + apiMock, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: def act = head.getFlux() @@ -245,7 +285,17 @@ class GenericWsHeadSpec extends Specification { ] } - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, apiMock, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + apiMock, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: def act = head.getFlux() @@ -298,7 +348,17 @@ class GenericWsHeadSpec extends Specification { ] } - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, apiMock, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + apiMock, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: def act = head.getFlux() @@ -351,7 +411,17 @@ class GenericWsHeadSpec extends Specification { Mono.just(new ChainResponse("".bytes, null)) } - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, reader, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + reader, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: def act = head.getFlux() @@ -388,7 +458,17 @@ class GenericWsHeadSpec extends Specification { ] } - def head = new GenericWsHead(new AlwaysForkChoice(), BlockValidator.ALWAYS_VALID, apiMock, ws, Schedulers.boundedElastic(), Schedulers.boundedElastic(), upstream, EthereumChainSpecific.INSTANCE) + def head = new GenericWsHead( + new AlwaysForkChoice(), + BlockValidator.ALWAYS_VALID, + apiMock, + ws, + Schedulers.boundedElastic(), + Schedulers.boundedElastic(), + upstream, + EthereumChainSpecific.INSTANCE, + Duration.ofSeconds(60), + ) when: head.start()