Update methods based on upstream availability (#216)
This commit is contained in:
@@ -54,7 +54,7 @@ abstract class DefaultUpstream(
|
||||
private val status = AtomicReference(Status(defaultLag, defaultAvail, statusByLag(defaultLag, defaultAvail)))
|
||||
private val statusStream = Sinks.many()
|
||||
.multicast()
|
||||
.directBestEffort<UpstreamAvailability>()
|
||||
.directBestEffort<UpstreamChangeState>()
|
||||
|
||||
init {
|
||||
if (id.length < 3 || !id.matches(Regex("[a-zA-Z][a-zA-Z0-9_-]+[a-zA-Z0-9]"))) {
|
||||
@@ -67,9 +67,14 @@ abstract class DefaultUpstream(
|
||||
}
|
||||
|
||||
fun onStatus(value: BlockchainOuterClass.ChainStatus) {
|
||||
this.onStatus(value, false)
|
||||
}
|
||||
|
||||
fun onStatus(value: BlockchainOuterClass.ChainStatus, stateChanged: Boolean = false) {
|
||||
val available = value.availability
|
||||
setStatus(
|
||||
if (available != null) UpstreamAvailability.fromGrpc(available.number) else UpstreamAvailability.UNAVAILABLE
|
||||
if (available != null) UpstreamAvailability.fromGrpc(available.number) else UpstreamAvailability.UNAVAILABLE,
|
||||
stateChanged
|
||||
)
|
||||
}
|
||||
|
||||
@@ -78,10 +83,16 @@ abstract class DefaultUpstream(
|
||||
}
|
||||
|
||||
open fun setStatus(avail: UpstreamAvailability) {
|
||||
this.setStatus(avail, false)
|
||||
}
|
||||
|
||||
open fun setStatus(avail: UpstreamAvailability, stateChanged: Boolean = false) {
|
||||
status.updateAndGet { curr ->
|
||||
Status(curr.lag, avail, statusByLag(curr.lag, avail))
|
||||
}.also {
|
||||
statusStream.emitNext(it.status) { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED }
|
||||
statusStream.emitNext(
|
||||
UpstreamChangeState(it.status, stateChanged)
|
||||
) { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED }
|
||||
log.trace("Status of upstream [$id] changed to [$it], requested change status to [$avail]")
|
||||
}
|
||||
}
|
||||
@@ -103,7 +114,17 @@ abstract class DefaultUpstream(
|
||||
|
||||
override fun observeStatus(): Flux<UpstreamAvailability> {
|
||||
return statusStream.asFlux()
|
||||
.distinctUntilChanged()
|
||||
.distinctUntilChanged(
|
||||
{ it },
|
||||
{ prev, current ->
|
||||
if (current.stateChanged) {
|
||||
false
|
||||
} else {
|
||||
prev.status == current.status
|
||||
}
|
||||
}
|
||||
)
|
||||
.map { it.status }
|
||||
}
|
||||
|
||||
override fun setLag(lag: Long) {
|
||||
@@ -111,7 +132,9 @@ abstract class DefaultUpstream(
|
||||
status.updateAndGet { curr ->
|
||||
Status(nLag, curr.avail, statusByLag(nLag, curr.avail))
|
||||
}.also {
|
||||
statusStream.emitNext(it.status) { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED }
|
||||
statusStream.emitNext(
|
||||
UpstreamChangeState(it.status, false)
|
||||
) { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED }
|
||||
log.trace("Status of upstream [$id] changed to [$it], requested change lag to [$lag]")
|
||||
}
|
||||
}
|
||||
@@ -147,4 +170,9 @@ abstract class DefaultUpstream(
|
||||
}
|
||||
|
||||
data class Status(val lag: Long, val avail: UpstreamAvailability, val status: UpstreamAvailability)
|
||||
|
||||
private data class UpstreamChangeState(
|
||||
val status: UpstreamAvailability,
|
||||
val stateChanged: Boolean
|
||||
)
|
||||
}
|
||||
|
||||
@@ -38,6 +38,7 @@ import reactor.core.Disposable
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.core.publisher.Mono
|
||||
import reactor.core.publisher.Sinks
|
||||
import reactor.util.function.Tuples
|
||||
import java.time.Duration
|
||||
import java.time.Instant
|
||||
import java.util.concurrent.atomic.AtomicReference
|
||||
@@ -81,6 +82,9 @@ abstract class Multistream(
|
||||
private val removedUpstreams = Sinks.many()
|
||||
.multicast()
|
||||
.directBestEffort<Upstream>()
|
||||
private val stateStream = Sinks.many()
|
||||
.multicast()
|
||||
.directBestEffort<UpstreamChangeState>()
|
||||
|
||||
init {
|
||||
UpstreamAvailability.values().forEach { status ->
|
||||
@@ -266,9 +270,41 @@ abstract class Multistream(
|
||||
// print status _change_ every 15 seconds, at most; otherwise prints it on interval of 30 seconds
|
||||
.sample(Duration.ofSeconds(15))
|
||||
.subscribe { printStatus() }
|
||||
|
||||
observeUpstreamsStatuses()
|
||||
|
||||
started = true
|
||||
}
|
||||
|
||||
private fun observeUpstreamsStatuses() {
|
||||
stateStream.asFlux()
|
||||
.distinctUntilChanged(
|
||||
{ it },
|
||||
{ prev, current ->
|
||||
prev.status == current.status || prev.equals(current)
|
||||
}
|
||||
).subscribe {
|
||||
upstreams.filter { it.isAvailable() }.map { it.getMethods() }.let {
|
||||
callMethods = AggregatedCallMethods(it)
|
||||
}
|
||||
}
|
||||
|
||||
subscribeAddedUpstreams()
|
||||
.filter { !it.isGrpc() }
|
||||
.distinctUntilChanged {
|
||||
it.getId()
|
||||
}.map {
|
||||
Tuples.of(it.getId(), it.observeStatus())
|
||||
}
|
||||
.subscribe { pair ->
|
||||
pair.t2.subscribe { status ->
|
||||
stateStream.emitNext(
|
||||
UpstreamChangeState(pair.t1, status)
|
||||
) { _, res -> res == Sinks.EmitResult.FAIL_NON_SERIALIZED }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override fun stop() {
|
||||
cacheSubscription?.dispose()
|
||||
cacheSubscription = null
|
||||
@@ -411,6 +447,9 @@ abstract class Multistream(
|
||||
fun subscribeRemovedUpstreams(): Flux<Upstream> =
|
||||
removedUpstreams.asFlux()
|
||||
|
||||
fun subscribeStateChanges(): Flux<UpstreamChangeState> =
|
||||
stateStream.asFlux()
|
||||
|
||||
abstract fun makeLagObserver(): HeadLagObserver
|
||||
|
||||
// --------------------------------------------------------------------------------------------------------
|
||||
@@ -435,4 +474,9 @@ abstract class Multistream(
|
||||
return curr == t
|
||||
}
|
||||
}
|
||||
|
||||
data class UpstreamChangeState(
|
||||
val upId: String,
|
||||
val status: UpstreamAvailability
|
||||
)
|
||||
}
|
||||
|
||||
@@ -166,7 +166,7 @@ class BitcoinGrpcUpstream(
|
||||
val upstreamStatusChanged = (upstreamStatus.update(conf) || (newCapabilities != capabilities)).also {
|
||||
capabilities = newCapabilities
|
||||
}
|
||||
conf.status?.let { status -> onStatus(status) }
|
||||
conf.status?.let { status -> onStatus(status, upstreamStatusChanged) }
|
||||
return buildInfoChanged || upstreamStatusChanged
|
||||
}
|
||||
}
|
||||
|
||||
@@ -151,7 +151,7 @@ open class EthereumGrpcUpstream(
|
||||
val upstreamStatusChanged = (upstreamStatus.update(conf) || (newCapabilities != capabilities)).also {
|
||||
capabilities = newCapabilities
|
||||
}
|
||||
conf.status?.let { status -> onStatus(status) }
|
||||
conf.status?.let { status -> onStatus(status, upstreamStatusChanged) }
|
||||
return buildInfoChanged || upstreamStatusChanged
|
||||
}
|
||||
|
||||
|
||||
@@ -118,7 +118,7 @@ open class EthereumPosGrpcUpstream(
|
||||
val upstreamStatusChanged = (upstreamStatus.update(conf) || (newCapabilities != capabilities)).also {
|
||||
capabilities = newCapabilities
|
||||
}
|
||||
conf.status?.let { status -> onStatus(status) }
|
||||
conf.status?.let { status -> onStatus(status, upstreamStatusChanged) }
|
||||
return buildInfoChanged || upstreamStatusChanged
|
||||
}
|
||||
|
||||
|
||||
@@ -55,7 +55,6 @@ import reactor.core.scheduler.Scheduler
|
||||
import java.io.IOException
|
||||
import java.time.Duration
|
||||
import java.util.concurrent.Executor
|
||||
import java.util.concurrent.atomic.AtomicReference
|
||||
import java.util.concurrent.locks.ReentrantLock
|
||||
import kotlin.concurrent.withLock
|
||||
|
||||
@@ -114,7 +113,7 @@ class GrpcUpstreams(
|
||||
}
|
||||
this.client = client
|
||||
|
||||
val statusSubscription = AtomicReference<Disposable>()
|
||||
val statusSubscriptions = mutableMapOf<Chain, Disposable>()
|
||||
|
||||
return Flux.interval(Duration.ZERO, Duration.ofSeconds(20))
|
||||
.flatMap {
|
||||
@@ -131,19 +130,19 @@ class GrpcUpstreams(
|
||||
}.flatMap { value ->
|
||||
processDescription(value)
|
||||
}.doOnNext {
|
||||
val subscription = client.subscribeStatus(
|
||||
StatusRequest.newBuilder()
|
||||
.addChains(Common.ChainRef.forNumber(it.chain.id)).build()
|
||||
).subscribeOn(chainStatusScheduler)
|
||||
.subscribe { value ->
|
||||
val chain = Chain.byId(value.chain.number)
|
||||
if (chain != Chain.UNSPECIFIED) {
|
||||
known[chain]?.onStatus(value)
|
||||
val sub = statusSubscriptions[it.chain]
|
||||
if (sub == null || sub.isDisposed) {
|
||||
val subscription = client.subscribeStatus(
|
||||
StatusRequest.newBuilder()
|
||||
.addChains(Common.ChainRef.forNumber(it.chain.id)).build()
|
||||
).subscribeOn(chainStatusScheduler)
|
||||
.subscribe { value ->
|
||||
val chain = Chain.byId(value.chain.number)
|
||||
if (chain != Chain.UNSPECIFIED) {
|
||||
known[chain]?.onStatus(value)
|
||||
}
|
||||
}
|
||||
}
|
||||
statusSubscription.updateAndGet { prev ->
|
||||
prev?.dispose()
|
||||
subscription
|
||||
statusSubscriptions[it.chain] = subscription
|
||||
}
|
||||
}.doOnError { t ->
|
||||
log.error("Failed to process update from gRPC upstream $id", t)
|
||||
|
||||
Reference in New Issue
Block a user