add gRPC-client compression (#136)
* add upstream grpc-client compression * fix review issues
This commit is contained in:
@@ -16,7 +16,16 @@
|
||||
*/
|
||||
package io.emeraldpay.dshackle
|
||||
|
||||
import io.emeraldpay.dshackle.config.*
|
||||
import io.emeraldpay.dshackle.config.CacheConfig
|
||||
import io.emeraldpay.dshackle.config.ChainsConfig
|
||||
import io.emeraldpay.dshackle.config.CompressionConfig
|
||||
import io.emeraldpay.dshackle.config.HealthConfig
|
||||
import io.emeraldpay.dshackle.config.MainConfig
|
||||
import io.emeraldpay.dshackle.config.MainConfigReader
|
||||
import io.emeraldpay.dshackle.config.MonitoringConfig
|
||||
import io.emeraldpay.dshackle.config.SignatureConfig
|
||||
import io.emeraldpay.dshackle.config.TokensConfig
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import org.bouncycastle.jce.provider.BouncyCastleProvider
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
@@ -110,6 +119,11 @@ open class Config(
|
||||
return mainConfig.upstreams
|
||||
}
|
||||
|
||||
@Bean
|
||||
open fun compressionConfig(@Autowired mainConfig: MainConfig): CompressionConfig {
|
||||
return mainConfig.compression
|
||||
}
|
||||
|
||||
@Bean
|
||||
open fun cacheConfig(@Autowired mainConfig: MainConfig): CacheConfig {
|
||||
return mainConfig.cache ?: CacheConfig()
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
package io.emeraldpay.dshackle.config
|
||||
|
||||
class CompressionConfig(
|
||||
open class CompressionConfig(
|
||||
var grpc: GRPC = GRPC()
|
||||
) {
|
||||
/**
|
||||
|
||||
@@ -22,6 +22,7 @@ import io.emeraldpay.dshackle.Chain
|
||||
import io.emeraldpay.dshackle.FileResolver
|
||||
import io.emeraldpay.dshackle.Global
|
||||
import io.emeraldpay.dshackle.config.ChainsConfig
|
||||
import io.emeraldpay.dshackle.config.CompressionConfig
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import io.emeraldpay.dshackle.reader.JsonRpcReader
|
||||
import io.emeraldpay.dshackle.upstream.*
|
||||
@@ -62,6 +63,7 @@ import kotlin.math.abs
|
||||
open class ConfiguredUpstreams(
|
||||
private val fileResolver: FileResolver,
|
||||
private val config: UpstreamsConfig,
|
||||
private val compressionConfig: CompressionConfig,
|
||||
private val callTargets: CallTargetsHolder,
|
||||
private val eventPublisher: ApplicationEventPublisher,
|
||||
@Qualifier("grpcChannelExecutor")
|
||||
@@ -86,7 +88,7 @@ open class ConfiguredUpstreams(
|
||||
log.debug("Start upstream ${up.id}")
|
||||
if (up.connection is UpstreamsConfig.GrpcConnection) {
|
||||
val options = up.options ?: UpstreamsConfig.Options()
|
||||
buildGrpcUpstream(up.nodeId, up.cast(UpstreamsConfig.GrpcConnection::class.java), options)
|
||||
buildGrpcUpstream(up.nodeId, up.cast(UpstreamsConfig.GrpcConnection::class.java), options, compressionConfig.grpc.clientEnabled)
|
||||
} else {
|
||||
val chain = Global.chainById(up.chain)
|
||||
if (chain == Chain.UNSPECIFIED) {
|
||||
@@ -307,7 +309,8 @@ open class ConfiguredUpstreams(
|
||||
private fun buildGrpcUpstream(
|
||||
nodeId: Int?,
|
||||
config: UpstreamsConfig.Upstream<UpstreamsConfig.GrpcConnection>,
|
||||
options: UpstreamsConfig.Options
|
||||
options: UpstreamsConfig.Options,
|
||||
compression: Boolean
|
||||
) {
|
||||
if (!this::grpcUpstreamsScheduler.isInitialized) {
|
||||
grpcUpstreamsScheduler = Schedulers.fromExecutorService(
|
||||
@@ -324,6 +327,7 @@ open class ConfiguredUpstreams(
|
||||
endpoint.host!!,
|
||||
endpoint.port,
|
||||
endpoint.auth,
|
||||
compression,
|
||||
fileResolver,
|
||||
endpoint.upstreamRating,
|
||||
config.labels,
|
||||
|
||||
@@ -34,6 +34,7 @@ import io.emeraldpay.dshackle.upstream.Lifecycle
|
||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.RpcMetrics
|
||||
import io.grpc.Codec
|
||||
import io.grpc.netty.NettyChannelBuilder
|
||||
import io.micrometer.core.instrument.Counter
|
||||
import io.micrometer.core.instrument.Metrics
|
||||
@@ -63,6 +64,7 @@ class GrpcUpstreams(
|
||||
private val host: String,
|
||||
private val port: Int,
|
||||
private val auth: AuthConfig.ClientTlsAuth? = null,
|
||||
private val compression: Boolean,
|
||||
private val fileResolver: FileResolver,
|
||||
private val nodeRating: Int,
|
||||
private val labels: UpstreamsConfig.Labels,
|
||||
@@ -94,7 +96,10 @@ class GrpcUpstreams(
|
||||
chanelBuilder.usePlaintext()
|
||||
}
|
||||
|
||||
val client = ReactorBlockchainGrpc.newReactorStub(chanelBuilder.build())
|
||||
var client = ReactorBlockchainGrpc.newReactorStub(chanelBuilder.build())
|
||||
if (compression) {
|
||||
client = client.withCompression(Codec.Gzip().messageEncoding)
|
||||
}
|
||||
this.client = client
|
||||
|
||||
val statusSubscription = AtomicReference<Disposable>()
|
||||
|
||||
Reference in New Issue
Block a user