check chain_ids on all networks (#507)
* check chain_ids on all networks
This commit is contained in:
@@ -44,16 +44,21 @@ abstract class UpstreamValidator(
|
||||
val cp = Comparator { avail1: UpstreamAvailability, avail2: UpstreamAvailability -> if (avail1.isBetterTo(avail2)) -1 else 1 }
|
||||
return results.sortedWith(cp).last()
|
||||
}
|
||||
|
||||
fun resolve(results: Iterable<ValidateUpstreamSettingsResult>): ValidateUpstreamSettingsResult {
|
||||
val cp = Comparator { res1: ValidateUpstreamSettingsResult, res2: ValidateUpstreamSettingsResult -> if (res1.priority < res2.priority) -1 else 1 }
|
||||
return results.sortedWith(cp).last()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
enum class ValidateUpstreamSettingsResult {
|
||||
UPSTREAM_VALID,
|
||||
UPSTREAM_SETTINGS_ERROR,
|
||||
UPSTREAM_FATAL_SETTINGS_ERROR,
|
||||
enum class ValidateUpstreamSettingsResult(val priority: Int) {
|
||||
UPSTREAM_VALID(0),
|
||||
UPSTREAM_SETTINGS_ERROR(1),
|
||||
UPSTREAM_FATAL_SETTINGS_ERROR(2),
|
||||
}
|
||||
|
||||
data class SingleCallValidator(
|
||||
data class SingleCallValidator<T>(
|
||||
val method: ChainRequest,
|
||||
val check: (ByteArray) -> UpstreamAvailability,
|
||||
val check: (ByteArray) -> T,
|
||||
)
|
||||
|
||||
@@ -14,15 +14,19 @@ import io.emeraldpay.dshackle.upstream.SingleCallValidator
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability.OK
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamValidator
|
||||
import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult
|
||||
import io.emeraldpay.dshackle.upstream.generic.AbstractPollChainSpecific
|
||||
import io.emeraldpay.dshackle.upstream.generic.GenericUpstreamValidator
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.ObjectParams
|
||||
import org.slf4j.LoggerFactory
|
||||
import reactor.core.publisher.Mono
|
||||
import java.math.BigInteger
|
||||
import java.time.Instant
|
||||
|
||||
object CosmosChainSpecific : AbstractPollChainSpecific() {
|
||||
|
||||
val log = LoggerFactory.getLogger(this::class.java)
|
||||
override fun latestBlockRequest(): ChainRequest = ChainRequest("block", ObjectParams())
|
||||
|
||||
override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer {
|
||||
@@ -75,9 +79,24 @@ object CosmosChainSpecific : AbstractPollChainSpecific() {
|
||||
return GenericUpstreamValidator(
|
||||
upstream,
|
||||
options,
|
||||
SingleCallValidator(
|
||||
ChainRequest("health", ListParams()),
|
||||
) { _ -> OK },
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("health", ListParams()),
|
||||
) { _ -> OK },
|
||||
),
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("status", ListParams()),
|
||||
) { data ->
|
||||
val resp = Global.objectMapper.readValue(data, CosmosStatus::class.java)
|
||||
if (chain.chainId.isNotEmpty() && resp.nodeInfo.network.lowercase() != chain.chainId.lowercase()) {
|
||||
ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR
|
||||
} else {
|
||||
ValidateUpstreamSettingsResult.UPSTREAM_VALID
|
||||
}
|
||||
},
|
||||
),
|
||||
|
||||
)
|
||||
}
|
||||
|
||||
@@ -115,6 +134,7 @@ data class CosmosStatus(
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
data class CosmosNodeInfo(
|
||||
@JsonProperty("version") var version: String,
|
||||
@JsonProperty("network") var network: String,
|
||||
)
|
||||
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
|
||||
@@ -8,17 +8,28 @@ import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamValidator
|
||||
import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult
|
||||
import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult.UPSTREAM_VALID
|
||||
import reactor.core.publisher.Mono
|
||||
import java.util.concurrent.TimeoutException
|
||||
|
||||
class GenericUpstreamValidator(
|
||||
upstream: Upstream,
|
||||
options: ChainOptions.Options,
|
||||
private val validator: SingleCallValidator,
|
||||
private val validators: List<SingleCallValidator<UpstreamAvailability>>,
|
||||
private val startupValidators: List<SingleCallValidator<ValidateUpstreamSettingsResult>>,
|
||||
) : UpstreamValidator(upstream, options) {
|
||||
|
||||
override fun validate(): Mono<UpstreamAvailability> {
|
||||
return Mono.zip(
|
||||
validators.map { exec(it, UpstreamAvailability.UNAVAILABLE) },
|
||||
) { a -> a.map { it as UpstreamAvailability } }
|
||||
.map(::resolve)
|
||||
.defaultIfEmpty(UpstreamAvailability.UNAVAILABLE)
|
||||
.onErrorResume {
|
||||
log.error("Error during upstream validation for ${upstream.getId()}", it)
|
||||
Mono.just(UpstreamAvailability.UNAVAILABLE)
|
||||
}
|
||||
}
|
||||
fun <T : Any> exec(validator: SingleCallValidator<T>, onError: T): Mono<T> {
|
||||
return upstream.getIngressReader()
|
||||
.read(validator.method)
|
||||
.flatMap(ChainResponse::requireResult)
|
||||
@@ -29,10 +40,22 @@ class GenericUpstreamValidator(
|
||||
.then(Mono.error(TimeoutException("Validation timeout for ${validator.method.method}"))),
|
||||
)
|
||||
.doOnError { err -> log.error("Error during ${validator.method.method} validation for ${upstream.getId()}", err) }
|
||||
.onErrorReturn(UpstreamAvailability.UNAVAILABLE)
|
||||
.onErrorReturn(onError)
|
||||
}
|
||||
|
||||
override fun validateUpstreamSettings(): Mono<ValidateUpstreamSettingsResult> {
|
||||
return Mono.just(UPSTREAM_VALID)
|
||||
return Mono.zip(
|
||||
startupValidators.map { exec(it, ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR) },
|
||||
) { a -> a.map { it as ValidateUpstreamSettingsResult } }
|
||||
.map(::resolve)
|
||||
.defaultIfEmpty(ValidateUpstreamSettingsResult.UPSTREAM_VALID)
|
||||
.onErrorResume {
|
||||
log.error("Error during upstream validation for ${upstream.getId()}", it)
|
||||
Mono.just(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR)
|
||||
}
|
||||
}
|
||||
|
||||
override fun validateUpstreamSettingsOnStartup(): ValidateUpstreamSettingsResult {
|
||||
return validateUpstreamSettings().block() ?: ValidateUpstreamSettingsResult.UPSTREAM_VALID
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,6 +15,7 @@ import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamSettingsDetector
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamValidator
|
||||
import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult
|
||||
import io.emeraldpay.dshackle.upstream.generic.AbstractPollChainSpecific
|
||||
import io.emeraldpay.dshackle.upstream.generic.GenericUpstreamValidator
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundService
|
||||
@@ -64,11 +65,20 @@ object NearChainSpecific : AbstractPollChainSpecific() {
|
||||
return GenericUpstreamValidator(
|
||||
upstream,
|
||||
options,
|
||||
SingleCallValidator(
|
||||
ChainRequest("status", ListParams()),
|
||||
) { data ->
|
||||
validate(data)
|
||||
},
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("status", ListParams()),
|
||||
) { data ->
|
||||
validate(data)
|
||||
},
|
||||
),
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("status", ListParams()),
|
||||
) { data ->
|
||||
validateSettings(data, chain)
|
||||
},
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -85,6 +95,15 @@ object NearChainSpecific : AbstractPollChainSpecific() {
|
||||
}
|
||||
}
|
||||
|
||||
fun validateSettings(data: ByteArray, chain: Chain): ValidateUpstreamSettingsResult {
|
||||
val resp = Global.objectMapper.readValue(data, NearStatus::class.java)
|
||||
return if (chain.chainId.isNotEmpty() && resp.chainId.lowercase() != chain.chainId.lowercase()) {
|
||||
ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR
|
||||
} else {
|
||||
ValidateUpstreamSettingsResult.UPSTREAM_VALID
|
||||
}
|
||||
}
|
||||
|
||||
override fun upstreamSettingsDetector(chain: Chain, upstream: Upstream): UpstreamSettingsDetector {
|
||||
return NearUpstreamSettingsDetector(upstream)
|
||||
}
|
||||
@@ -108,6 +127,7 @@ data class NearHeader(
|
||||
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
data class NearStatus(
|
||||
@JsonProperty("chain_id") var chainId: String,
|
||||
@JsonProperty("sync_info") var syncInfo: NearSync,
|
||||
)
|
||||
|
||||
|
||||
@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.upstream.SingleCallValidator
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamValidator
|
||||
import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult
|
||||
import io.emeraldpay.dshackle.upstream.calls.CallMethods
|
||||
import io.emeraldpay.dshackle.upstream.calls.DefaultPolkadotMethods
|
||||
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
|
||||
@@ -97,11 +98,20 @@ object PolkadotChainSpecific : AbstractPollChainSpecific() {
|
||||
return GenericUpstreamValidator(
|
||||
upstream,
|
||||
options,
|
||||
SingleCallValidator(
|
||||
ChainRequest("system_health", ListParams()),
|
||||
) { data ->
|
||||
validate(data, options.minPeers, upstream.getId())
|
||||
},
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("system_health", ListParams()),
|
||||
) { data ->
|
||||
validate(data, options.minPeers, upstream.getId())
|
||||
},
|
||||
),
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("system_chain", ListParams()),
|
||||
) { data ->
|
||||
validateSettings(data, chain)
|
||||
},
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -123,6 +133,15 @@ object PolkadotChainSpecific : AbstractPollChainSpecific() {
|
||||
return UpstreamAvailability.OK
|
||||
}
|
||||
|
||||
fun validateSettings(data: ByteArray, chain: Chain): ValidateUpstreamSettingsResult {
|
||||
val id = Global.objectMapper.readValue(data, String::class.java)
|
||||
return if (chain.chainId.isNotEmpty() && id.lowercase() != chain.chainId.lowercase()) {
|
||||
ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR
|
||||
} else {
|
||||
ValidateUpstreamSettingsResult.UPSTREAM_VALID
|
||||
}
|
||||
}
|
||||
|
||||
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
|
||||
return GenericIngressSubscription(ws, DefaultPolkadotMethods.subs.map { it.first })
|
||||
}
|
||||
|
||||
@@ -124,17 +124,20 @@ object SolanaChainSpecific : AbstractChainSpecific() {
|
||||
return GenericUpstreamValidator(
|
||||
upstream,
|
||||
options,
|
||||
SingleCallValidator(
|
||||
ChainRequest("getHealth", ListParams()),
|
||||
) { data ->
|
||||
val resp = String(data)
|
||||
if (resp == "\"ok\"") {
|
||||
UpstreamAvailability.OK
|
||||
} else {
|
||||
log.warn("Upstream {} validation failed, solana status is {}", upstream.getId(), resp)
|
||||
UpstreamAvailability.UNAVAILABLE
|
||||
}
|
||||
},
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("getHealth", ListParams()),
|
||||
) { data ->
|
||||
val resp = String(data)
|
||||
if (resp == "\"ok\"") {
|
||||
UpstreamAvailability.OK
|
||||
} else {
|
||||
log.warn("Upstream {} validation failed, solana status is {}", upstream.getId(), resp)
|
||||
UpstreamAvailability.UNAVAILABLE
|
||||
}
|
||||
},
|
||||
),
|
||||
listOf(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -65,11 +65,14 @@ object StarknetChainSpecific : AbstractPollChainSpecific() {
|
||||
return GenericUpstreamValidator(
|
||||
upstream,
|
||||
options,
|
||||
SingleCallValidator(
|
||||
ChainRequest("starknet_syncing", ListParams()),
|
||||
) { data ->
|
||||
validate(data, config.laggingLagSize, upstream.getId())
|
||||
},
|
||||
listOf(
|
||||
SingleCallValidator(
|
||||
ChainRequest("starknet_syncing", ListParams()),
|
||||
) { data ->
|
||||
validate(data, config.laggingLagSize, upstream.getId())
|
||||
},
|
||||
),
|
||||
listOf(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user