Autodetect groups by rpc_modules, temporarily ban methods with "doesn't exist/is not available" error (#505)

This commit is contained in:
msizov
2024-06-07 22:54:56 +07:00
committed by GitHub
parent 2b94078325
commit 73e5030980
13 changed files with 206 additions and 12 deletions

View File

@@ -30,5 +30,6 @@ class Defaults {
val grpcServerKeepAliveTimeout: Long = 5
val grpcServerPermitKeepAliveTime: Long = 15
val grpcServerMaxConnectionIdle: Long = 3600
val multistreamUnavailableMethodDisableDuration: Long = 20 // minutes
}
}

View File

@@ -166,8 +166,8 @@ data class UpstreamsConfig(
)
data class MethodGroups(
val enabled: Set<String>,
val disabled: Set<String>,
var enabled: Set<String>,
var disabled: Set<String>,
)
data class Method(

View File

@@ -102,7 +102,20 @@ open class NativeCall(
CallResult.fail(id, 0, err, null),
)
}
.doOnNext { callRes -> completeSpan(callRes, requestCount) }
.doOnNext {
callRes ->
if (callRes.error?.message?.contains(Regex("method ([A-Za-z0-9_]+) does not exist/is not available")) == true) {
if (it is ValidCallContext<*>) {
if (it.payload is ParsedCallDetails) {
log.error("nativeCallResult method ${it.payload.method} of ${it.upstream.getId()} is not available, disabling")
val cm = (it.upstream.getMethods() as Multistream.DisabledCallMethods)
cm.disableMethodTemporarily(it.payload.method)
it.upstream.updateMethods(cm)
}
}
}
completeSpan(callRes, requestCount)
}
.doOnCancel {
tracer.currentSpan()?.tag(SPAN_STATUS_MESSAGE, SPAN_REQUEST_CANCELLED)?.end()
}
@@ -315,7 +328,6 @@ open class NativeCall(
""
}
val availableMethods = upstream.getMethods()
if (!availableMethods.isAvailable(method)) {
val errorMessage = "The method $method is not available"
return Mono.just(

View File

@@ -65,25 +65,24 @@ open class GenericUpstreamCreator(
chainConfig,
) ?: return UpstreamCreationData.default()
val methods = buildMethods(config, chain)
val hashUrl = connection.let {
if (it.connectorMode == GenericConnectorFactory.ConnectorMode.RPC_REQUESTS_WITH_MIXED_HEAD.name) it.rpc?.url ?: it.ws?.url else it.ws?.url ?: it.rpc?.url
}
val hash = getHash(nodeId, hashUrl!!, hashes)
val buildMethodsFun = { a: UpstreamsConfig.Upstream<*>, b: Chain -> this.buildMethods(a, b) }
val upstream = GenericUpstream(
config.id!!,
config,
chain,
hash,
options,
config.role,
methods,
QuorumForLabels.QuorumItem(1, UpstreamsConfig.Labels.fromMap(config.labels)),
chainConfig,
connectorFactory,
cs::validator,
cs::upstreamSettingsDetector,
cs::upstreamRpcModulesDetector,
buildMethodsFun,
cs::lowerBoundService,
)

View File

@@ -70,6 +70,15 @@ abstract class UpstreamCreator(
protected fun buildMethods(config: UpstreamsConfig.Upstream<*>, chain: Chain): CallMethods {
return if (config.methods != null || config.methodGroups != null) {
if (config.methodGroups == null) {
config.methodGroups = UpstreamsConfig.MethodGroups(setOf("filter"), setOf())
} else {
val disabled = config.methodGroups!!.disabled
if (!disabled.contains("filter")) {
config.methodGroups!!.enabled = config.methodGroups!!.enabled.plus("filter")
}
}
ManagedCallMethods(
delegate = callTargets.getDefaultMethods(chain, indexConfig.isChainEnabled(chain)),
enabled = config.methods?.enabled?.map { it.name }?.toSet() ?: emptySet(),

View File

@@ -37,7 +37,7 @@ abstract class DefaultUpstream(
defaultAvail: UpstreamAvailability,
private val options: ChainOptions.Options,
private val role: UpstreamsConfig.UpstreamRole,
private val targets: CallMethods?,
private var targets: CallMethods?,
private val node: QuorumForLabels.QuorumItem?,
private val chainConfig: ChainsConfig.ChainConfig,
private val chain: Chain,
@@ -147,6 +147,11 @@ abstract class DefaultUpstream(
return targets ?: throw IllegalStateException("Methods are not set")
}
override fun updateMethods(m: CallMethods) {
targets = m
sendUpstreamStateEvent(UpstreamChangeEvent.ChangeType.UPDATED)
}
override fun nodeId(): Byte = hash
open fun getQuorumByLabel(): QuorumForLabels {

View File

@@ -16,12 +16,16 @@
*/
package io.emeraldpay.dshackle.upstream
import com.github.benmanes.caffeine.cache.Cache
import com.github.benmanes.caffeine.cache.Caffeine
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Defaults.Companion.multistreamUnavailableMethodDisableDuration
import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.foundation.ChainOptions
import io.emeraldpay.dshackle.quorum.CallQuorum
import io.emeraldpay.dshackle.reader.ChainReader
import io.emeraldpay.dshackle.startup.QuorumForLabels
import io.emeraldpay.dshackle.startup.UpstreamChangeEvent
@@ -68,7 +72,7 @@ abstract class Multistream(
private var cacheSubscription: Disposable? = null
@Volatile
private var callMethods: CallMethods? = null
private var callMethods: DisabledCallMethods? = null
private var callMethodsFactory: Factory<CallMethods> = Factory {
return@Factory callMethods ?: throw FunctorException("Not initialized yet")
}
@@ -227,7 +231,16 @@ abstract class Multistream(
val upstreams = getAll()
val availableUpstreams = upstreams.filter { it.isAvailable() }
availableUpstreams.map { it.getMethods() }.let {
callMethods = AggregatedCallMethods(it)
if (callMethods == null) {
callMethods = DisabledCallMethods(this, multistreamUnavailableMethodDisableDuration, AggregatedCallMethods(it))
} else {
callMethods = DisabledCallMethods(
this,
multistreamUnavailableMethodDisableDuration,
AggregatedCallMethods(it),
callMethods!!.disabledMethods,
)
}
}
capabilities = if (upstreams.isEmpty()) {
emptySet()
@@ -320,6 +333,10 @@ abstract class Multistream(
return callMethods ?: throw IllegalStateException("Methods are not initialized yet")
}
override fun updateMethods(m: CallMethods) {
onUpstreamsUpdated()
}
fun getMethodsFactory(): Factory<CallMethods> {
return callMethodsFactory
}
@@ -552,6 +569,55 @@ abstract class Multistream(
// --------------------------------------------------------------------------------------------------------
class DisabledCallMethods(private val multistream: Multistream, private val defaultDisableTimeout: Long, private val callMethods: CallMethods) : CallMethods {
var disabledMethods: Cache<String, Boolean> = Caffeine.newBuilder()
.removalListener { key: String?, _: Boolean?, cause ->
if (cause.wasEvicted() && key != null) {
multistream.log.info("${multistream.getId()} restoring method $key")
multistream.onUpstreamsUpdated()
}
}
.expireAfterWrite(Duration.ofMinutes(defaultDisableTimeout))
.build<String, Boolean>()
constructor(
multistream: Multistream,
defaultDisableTimeout: Long,
callMethods: CallMethods,
disabledMethodsCopy: Cache<String, Boolean>,
) : this(multistream, defaultDisableTimeout, callMethods) {
disabledMethods = disabledMethodsCopy
}
override fun createQuorumFor(method: String): CallQuorum {
return callMethods.createQuorumFor(method)
}
override fun isCallable(method: String): Boolean {
return callMethods.isCallable(method) && disabledMethods.getIfPresent(method) == null
}
override fun getSupportedMethods(): Set<String> {
return callMethods.getSupportedMethods() - disabledMethods.asMap().keys
}
override fun isHardcoded(method: String): Boolean {
return callMethods.isHardcoded(method)
}
override fun executeHardcoded(method: String): ByteArray {
return callMethods.executeHardcoded(method)
}
override fun getGroupMethods(groupName: String): Set<String> {
return callMethods.getGroupMethods(groupName)
}
fun disableMethodTemporarily(method: String) {
disabledMethods.put(method, true)
}
}
class UpstreamStatus(val upstream: Upstream, val status: UpstreamAvailability)
class FilterBestAvailability : java.util.function.Function<UpstreamStatus, UpstreamAvailability> {

View File

@@ -43,6 +43,7 @@ interface Upstream : Lifecycle {
fun getLag(): Long?
fun getLabels(): Collection<UpstreamsConfig.Labels>
fun getMethods(): CallMethods
fun updateMethods(m: CallMethods)
fun getId(): String
fun getCapabilities(): Set<Capability>
fun isGrpc(): Boolean

View File

@@ -0,0 +1,41 @@
package io.emeraldpay.dshackle.upstream
import com.fasterxml.jackson.core.type.TypeReference
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono
typealias UpstreamRpcModulesDetectorBuilder = (Upstream) -> UpstreamRpcModulesDetector?
abstract class UpstreamRpcModulesDetector(
private val upstream: Upstream,
) {
protected val log: Logger = LoggerFactory.getLogger(this::class.java)
open fun detectRpcModules(): Mono<HashMap<String, String>> {
return upstream.getIngressReader()
.read(rpcModulesRequest())
.flatMap(ChainResponse::requireResult)
.map(::parseRpcModules)
.onErrorResume {
log.warn("Can't detect rpc_modules of upstream ${upstream.getId()}, reason - {}", it.message)
Mono.just(HashMap())
}
}
protected abstract fun rpcModulesRequest(): ChainRequest
protected abstract fun parseRpcModules(data: ByteArray): HashMap<String, String>
}
class BasicEthUpstreamRpcModulesDetector(
upstream: Upstream,
) : UpstreamRpcModulesDetector(upstream) {
override fun rpcModulesRequest(): ChainRequest = ChainRequest("rpc_modules", ListParams())
override fun parseRpcModules(data: ByteArray): HashMap<String, String> {
return Global.objectMapper.readValue(data, object : TypeReference<HashMap<String, String>>() {})
}
}

View File

@@ -6,6 +6,7 @@ import io.emeraldpay.dshackle.config.ChainsConfig.ChainConfig
import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.foundation.ChainOptions.Options
import io.emeraldpay.dshackle.reader.ChainReader
import io.emeraldpay.dshackle.upstream.BasicEthUpstreamRpcModulesDetector
import io.emeraldpay.dshackle.upstream.CachingReader
import io.emeraldpay.dshackle.upstream.ChainRequest
import io.emeraldpay.dshackle.upstream.EgressSubscription
@@ -14,6 +15,7 @@ import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.LogsOracle
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamRpcModulesDetector
import io.emeraldpay.dshackle.upstream.UpstreamSettingsDetector
import io.emeraldpay.dshackle.upstream.UpstreamValidator
import io.emeraldpay.dshackle.upstream.calls.CallMethods
@@ -91,6 +93,10 @@ object EthereumChainSpecific : AbstractPollChainSpecific() {
return EthereumUpstreamValidator(chain, upstream, options, config)
}
override fun upstreamRpcModulesDetector(upstream: Upstream): UpstreamRpcModulesDetector {
return BasicEthUpstreamRpcModulesDetector(upstream)
}
override fun lowerBoundService(chain: Chain, upstream: Upstream): LowerBoundService {
return EthereumLowerBoundService(chain, upstream)
}

View File

@@ -15,6 +15,7 @@ import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.NoIngressSubscription
import io.emeraldpay.dshackle.upstream.NoopCachingReader
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamRpcModulesDetector
import io.emeraldpay.dshackle.upstream.UpstreamSettingsDetector
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.CallSelector
@@ -42,6 +43,10 @@ abstract class AbstractChainSpecific : ChainSpecific {
return null
}
override fun upstreamRpcModulesDetector(upstream: Upstream): UpstreamRpcModulesDetector? {
return null
}
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
return NoIngressSubscription()
}

View File

@@ -23,6 +23,7 @@ import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.LogsOracle
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamRpcModulesDetector
import io.emeraldpay.dshackle.upstream.UpstreamSettingsDetector
import io.emeraldpay.dshackle.upstream.UpstreamValidator
import io.emeraldpay.dshackle.upstream.beaconchain.BeaconChainSpecific
@@ -64,6 +65,8 @@ interface ChainSpecific {
fun upstreamSettingsDetector(chain: Chain, upstream: Upstream): UpstreamSettingsDetector?
fun upstreamRpcModulesDetector(upstream: Upstream): UpstreamRpcModulesDetector?
fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription
fun callSelector(caches: Caches): CallSelector?

View File

@@ -16,6 +16,8 @@ import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.UNKNOWN_CLIENT_VERSION
import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
import io.emeraldpay.dshackle.upstream.UpstreamRpcModulesDetector
import io.emeraldpay.dshackle.upstream.UpstreamRpcModulesDetectorBuilder
import io.emeraldpay.dshackle.upstream.UpstreamSettingsDetectorBuilder
import io.emeraldpay.dshackle.upstream.UpstreamValidator
import io.emeraldpay.dshackle.upstream.UpstreamValidatorBuilder
@@ -48,6 +50,24 @@ open class GenericUpstream(
lowerBoundServiceBuilder: LowerBoundServiceBuilder,
) : DefaultUpstream(id, hash, null, UpstreamAvailability.OK, options, role, targets, node, chainConfig, chain), Lifecycle {
constructor(
config: UpstreamsConfig.Upstream<*>,
chain: Chain,
hash: Byte,
options: ChainOptions.Options,
node: QuorumForLabels.QuorumItem?,
chainConfig: ChainsConfig.ChainConfig,
connectorFactory: ConnectorFactory,
validatorBuilder: UpstreamValidatorBuilder,
upstreamSettingsDetectorBuilder: UpstreamSettingsDetectorBuilder,
upstreamRpcModulesDetectorBuilder: UpstreamRpcModulesDetectorBuilder,
buildMethods: (UpstreamsConfig.Upstream<*>, Chain) -> CallMethods,
lowerBoundServiceBuilder: LowerBoundServiceBuilder,
) : this(config.id!!, chain, hash, options, config.role, buildMethods(config, chain), node, chainConfig, connectorFactory, validatorBuilder, upstreamSettingsDetectorBuilder, lowerBoundServiceBuilder) {
rpcModulesDetector = upstreamRpcModulesDetectorBuilder(this)
detectRpcModules(config, buildMethods)
}
private val validator: UpstreamValidator? = validatorBuilder(chain, this, getOptions(), chainConfig)
private var validatorSubscription: Disposable? = null
private var validationSettingsSubscription: Disposable? = null
@@ -57,6 +77,7 @@ open class GenericUpstream(
protected val connector: GenericConnector = connectorFactory.create(this, chain)
private var livenessSubscription: Disposable? = null
private val settingsDetector = upstreamSettingsDetectorBuilder(chain, this)
private var rpcModulesDetector: UpstreamRpcModulesDetector? = null
private val lowerBoundService = lowerBoundServiceBuilder(chain, this)
@@ -190,6 +211,31 @@ open class GenericUpstream(
}
}
private fun detectRpcModules(config: UpstreamsConfig.Upstream<*>, buildMethods: (UpstreamsConfig.Upstream<*>, Chain) -> CallMethods) {
rpcModulesDetector?.detectRpcModules()
val rpcDetector = rpcModulesDetector?.detectRpcModules()?.block() ?: HashMap<String, String>()
log.info("Upstream rpc detector for ${getId()} returned $rpcDetector ")
if (rpcDetector.size != 0) {
var changed = false
for ((group, _) in rpcDetector) {
if (group == "trace" || group == "debug" || group == "filter") {
if (config.methodGroups == null) {
config.methodGroups = UpstreamsConfig.MethodGroups(setOf("filter"), setOf())
} else {
val disabled = config.methodGroups!!.disabled
val enabled = config.methodGroups!!.enabled
if (!disabled.contains(group) && !enabled.contains(group)) {
config.methodGroups!!.enabled = enabled.plus(group)
changed = true
}
}
}
}
if (changed) updateMethods(buildMethods(config, chain))
}
}
private fun upstreamStart() {
if (getOptions().disableValidation) {
log.warn("Disable validation for upstream ${this.getId()}")