diff --git a/emerald-grpc b/emerald-grpc index 9ae95ae1..9a2bedc9 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit 9ae95ae1b5864c678fe920cd24c0d6986dc77119 +Subproject commit 9a2bedc9999ccbc34c9f08351ea3f33eb48322a3 diff --git a/src/main/kotlin/io/emeraldpay/dshackle/Defaults.kt b/src/main/kotlin/io/emeraldpay/dshackle/Defaults.kt index 571cb854..3ee70ffb 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/Defaults.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/Defaults.kt @@ -32,5 +32,6 @@ class Defaults { val grpcServerPermitKeepAliveTime: Long = 15 val grpcServerMaxConnectionIdle: Long = 3600 val multistreamUnavailableMethodDisableDuration: Long = 20 // minutes + val nativeSubscribeHeartbeatInterval: Duration = Duration.ofSeconds(30) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt index d7de94b6..86135c8b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeSubscribe.kt @@ -19,6 +19,7 @@ import com.google.protobuf.ByteString import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.api.proto.BlockchainOuterClass.NativeSubscribeReplyItem import io.emeraldpay.dshackle.Chain +import io.emeraldpay.dshackle.Defaults import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.SilentException import io.emeraldpay.dshackle.upstream.Multistream @@ -34,6 +35,7 @@ import org.springframework.beans.factory.annotation.Autowired import org.springframework.stereotype.Service import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import java.util.concurrent.atomic.AtomicLong @Service open class NativeSubscribe( @@ -51,14 +53,26 @@ open class NativeSubscribe( .flatMapMany { it -> val subscriptionId = it.subscriptionId - Mono.just(it) + val dataStream = Mono.just(it) .flatMapMany { start(it, subscriptionId) } .map(this@NativeSubscribe::convertToProto) + + val lastMessageTime = AtomicLong(System.currentTimeMillis()) + + val heartbeatStream = Flux.interval(Defaults.nativeSubscribeHeartbeatInterval) + .filter { System.currentTimeMillis() - lastMessageTime.get() >= Defaults.nativeSubscribeHeartbeatInterval.toMillis() } + .map { createHeartbeat() } + + Flux.merge( + dataStream.doOnNext { lastMessageTime.set(System.currentTimeMillis()) }, + heartbeatStream, + ) .doOnCancel { log.warn("Subscription $subscriptionId cancelled") - }.onErrorMap { + } + .onErrorMap { convertToStatus(it, subscriptionId) } } @@ -172,6 +186,12 @@ open class NativeSubscribe( return msg.build() } + private fun createHeartbeat(): NativeSubscribeReplyItem { + return NativeSubscribeReplyItem.newBuilder() + .setHeartbeat(true) + .build() + } + data class ResponseHolder( val response: Any, val nonce: Long?,