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 197fd056..cbf8c6cf 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt @@ -58,6 +58,7 @@ import java.time.Duration import java.time.Instant import java.util.Base64 import java.util.concurrent.Executors +import java.util.concurrent.ScheduledFuture import java.util.concurrent.TimeUnit import java.util.concurrent.TimeoutException import java.util.concurrent.atomic.AtomicBoolean @@ -85,6 +86,7 @@ open class WsConnection( // > io.netty.handler.codec.http.websocketx.CorruptedWebSocketFrameException: Max frame length of 65536 has been exceeded // It's unclear what is the right limit here, but 5mb seems to be working (1mb isn't always working) private const val DEFAULT_FRAME_SIZE = 5 * 1024 * 1024 + private const val RESET_BACKOFF_TIMEOUT = 10 * 1000 // The max size from multiple frames that may represent a single message // Accept up to 15Mb messages, because Geth is using 15mb, though it's not clear what should be a right value @@ -96,12 +98,15 @@ open class WsConnection( private var reconnectBackoff: BackOff = ExponentialBackOff().also { it.initialInterval = Duration.ofMillis(100).toMillis() - it.maxInterval = Duration.ofMinutes(1).toMillis() + it.maxInterval = Duration.ofMinutes(5).toMillis() } private var currentBackOff = reconnectBackoff.start() private val parser = ResponseWSParser() + private val resetBackoffExecutor = Executors.newScheduledThreadPool(2) + private var resetBackoffTask: ScheduledFuture? = null + private val blocks = Sinks .many() .multicast() @@ -147,6 +152,7 @@ open class WsConnection( if (alreadyReconnecting) { return } + // rpcSend is already CANCELLED, since the subscription owned by the previous connection is gone // so we need to create a new Sink. Emit Complete is probably useless, and just in case rpcSend.tryEmitComplete() @@ -154,7 +160,11 @@ open class WsConnection( .many() .unicast() .onBackpressureBuffer() + resetBackoffTask?.cancel(false) val retryInterval = currentBackOff.nextBackOff() + resetBackoffTask = resetBackoffExecutor.schedule({ + currentBackOff = reconnectBackoff.start() + }, RESET_BACKOFF_TIMEOUT + retryInterval, TimeUnit.MILLISECONDS) if (retryInterval == BackOffExecution.STOP) { log.warn("Reconnect backoff exhausted. Permanently closing the connection") return @@ -165,7 +175,8 @@ open class WsConnection( reconnecting.set(false) connectInternal() }, - retryInterval, TimeUnit.MILLISECONDS + retryInterval, + TimeUnit.MILLISECONDS ) } @@ -224,9 +235,6 @@ open class WsConnection( } fun handle(inbound: WebsocketInbound, outbound: WebsocketOutbound): Publisher { - // restart backoff after connection - currentBackOff = reconnectBackoff.start() - val consumer = inbound .aggregateFrames(msgSizeLimit) .receiveFrames() @@ -298,7 +306,10 @@ open class WsConnection( fun onRpc(msg: ResponseWSParser.WsResponse): Mono { return if (msg.id.isNumber()) { val resp = JsonRpcResponse( - msg.value, msg.error, msg.id, null + msg.value, + msg.error, + msg.id, + null ) Mono.fromCallable { val status = rpcReceive.tryEmitNext(resp)