imcrease gRPC message size limit
This commit is contained in:
@@ -29,7 +29,6 @@ import io.emeraldpay.dshackle.upstream.DefaultUpstream
|
|||||||
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
import io.emeraldpay.dshackle.upstream.UpstreamAvailability
|
||||||
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
|
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcGrpcClient
|
||||||
import io.emeraldpay.dshackle.upstream.rpcclient.RpcMetrics
|
import io.emeraldpay.dshackle.upstream.rpcclient.RpcMetrics
|
||||||
import io.grpc.ManagedChannelBuilder
|
|
||||||
import io.grpc.netty.NettyChannelBuilder
|
import io.grpc.netty.NettyChannelBuilder
|
||||||
import io.micrometer.core.instrument.Counter
|
import io.micrometer.core.instrument.Counter
|
||||||
import io.micrometer.core.instrument.Metrics
|
import io.micrometer.core.instrument.Metrics
|
||||||
@@ -70,21 +69,21 @@ class GrpcUpstreams(
|
|||||||
private val lock = ReentrantLock()
|
private val lock = ReentrantLock()
|
||||||
|
|
||||||
fun start(): Flux<UpstreamChangeEvent> {
|
fun start(): Flux<UpstreamChangeEvent> {
|
||||||
val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) {
|
val chanelBuilder = NettyChannelBuilder.forAddress(host, port)
|
||||||
NettyChannelBuilder.forAddress(host, port)
|
// some messages are very large. many of them in megabytes, some even in gigabytes (ex. ETH Traces)
|
||||||
// some messages are very large. many of them in megabytes, some even in gigabytes (ex. ETH Traces)
|
.maxInboundMessageSize(Int.MAX_VALUE)
|
||||||
.maxInboundMessageSize(Int.MAX_VALUE)
|
.enableRetry()
|
||||||
|
.maxRetryAttempts(3)
|
||||||
|
if (auth != null && StringUtils.isNotEmpty(auth.ca)) {
|
||||||
|
chanelBuilder
|
||||||
.useTransportSecurity()
|
.useTransportSecurity()
|
||||||
.enableRetry()
|
|
||||||
.maxRetryAttempts(3)
|
|
||||||
.sslContext(withTls(auth))
|
.sslContext(withTls(auth))
|
||||||
} else {
|
} else {
|
||||||
log.warn("Using insecure connection to $host:$port")
|
log.warn("Using insecure connection to $host:$port")
|
||||||
ManagedChannelBuilder.forAddress(host, port)
|
chanelBuilder.usePlaintext()
|
||||||
.usePlaintext()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
val client = ReactorBlockchainGrpc.newReactorStub(channel.build())
|
val client = ReactorBlockchainGrpc.newReactorStub(chanelBuilder.build())
|
||||||
this.client = client
|
this.client = client
|
||||||
|
|
||||||
val statusSubscription = AtomicReference<Disposable>()
|
val statusSubscription = AtomicReference<Disposable>()
|
||||||
|
|||||||
Reference in New Issue
Block a user