diff --git a/build.gradle b/build.gradle index bafce381..13da52b7 100644 --- a/build.gradle +++ b/build.gradle @@ -76,6 +76,7 @@ dependencies { compile "io.infinitape:etherjar-hex:$etherjarVersion" compile "io.infinitape:etherjar-rpc-http:$etherjarVersion" compile "io.infinitape:etherjar-rpc-ws:$etherjarVersion" + compile "io.infinitape:etherjar-rpc-emerald:$etherjarVersion" compile "io.infinitape:etherjar-tx:$etherjarVersion" compile 'org.apache.httpcomponents:httpmime:4.5.8' diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcTransport.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcTransport.kt deleted file mode 100644 index ac1735e6..00000000 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcTransport.kt +++ /dev/null @@ -1,133 +0,0 @@ -/** - * Copyright (c) 2019 ETCDEV GmbH - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package io.emeraldpay.dshackle.upstream.grpc - -import com.fasterxml.jackson.databind.ObjectMapper -import com.google.protobuf.ByteString -import io.emeraldpay.api.proto.BlockchainOuterClass -import io.emeraldpay.api.proto.Common -import io.emeraldpay.api.proto.ReactorBlockchainGrpc -import io.emeraldpay.dshackle.upstream.Selector -import io.emeraldpay.grpc.Chain -import io.infinitape.etherjar.rpc.Batch -import io.infinitape.etherjar.rpc.JacksonRpcConverter -import io.infinitape.etherjar.rpc.RpcException -import io.infinitape.etherjar.rpc.RpcResponseError -import io.infinitape.etherjar.rpc.transport.BatchStatus -import io.infinitape.etherjar.rpc.transport.RpcTransport -import reactor.util.function.Tuple3 -import reactor.util.function.Tuples -import java.time.Duration -import java.util.concurrent.CompletableFuture -import java.util.function.Function - -class EthereumGrpcTransport( - private val chainRef: Common.ChainRef, - private val labelSelector: BlockchainOuterClass.Selector?, - private val client: ReactorBlockchainGrpc.ReactorBlockchainStub, - private val objectMapper: ObjectMapper -): RpcTransport { - - private val jacksonRpcConverter = JacksonRpcConverter(objectMapper) - - constructor( - chain: Chain, - client: ReactorBlockchainGrpc.ReactorBlockchainStub, - objectMapper: ObjectMapper - ) : this(Common.ChainRef.forNumber(chain.id), null, client, objectMapper) - - fun withLabels(matcher: Selector.LabelSelectorMatcher?): EthereumGrpcTransport { - if ((matcher == null || matcher is Selector.AnyLabelMatcher) && labelSelector == null) { - return this - } - return EthereumGrpcTransport(chainRef, matcher?.asProto(), client, objectMapper) - } - - override fun close() { - } - - private fun replyProcessor(mapping: HashMap>): Function { - return Function { resp -> - val id = resp.id - val bi = mapping.remove(id) - if (bi != null) { - if (resp.succeed) { - try { - val rpcResp = jacksonRpcConverter.fromJson(resp.payload.toByteArray().inputStream(), bi.call.jsonType, Int::class.java) - bi.onComplete(rpcResp) - return@Function true - } catch (e: RpcException) { - bi.onError(e) - } - } else { - bi.onError(RpcException(RpcResponseError.CODE_INTERNAL_ERROR, resp.errorMessage)) - } - } - false - } - } - - private val sumStatus = { t: Tuple3, ok: Boolean -> - if (ok) Tuples.of(t.t1 + 1, t.t2, t.t3 + 1) - else Tuples.of(t.t1, t.t2 + 1, t.t3 + 1) - } - - private val asStatus = Function, BatchStatus> { - BatchStatus.newBuilder() - .withSucceed(it.t1) - .withFailed(it.t2) - .withTotal(it.t3) - .build() - } - - fun prepareMapping(items: List>, req: BlockchainOuterClass.NativeCallRequest.Builder): HashMap> { - val mapping = HashMap>() - var seq: Int = 0 - items.forEach { bi -> - val id = seq++ - mapping[id] = bi as Batch.BatchItem - val call = bi.call - val params = objectMapper.writeValueAsBytes(call.params) - val nativeCallItem = BlockchainOuterClass.NativeCallItem.newBuilder() - .setId(id) - .setMethod(call.method) - .setPayload(ByteString.copyFrom(params)) - .build() - req.addItems(nativeCallItem) - } - return mapping - } - - override fun execute(items: List>): CompletableFuture { - val req = BlockchainOuterClass.NativeCallRequest.newBuilder() - .setChain(chainRef) - if (labelSelector != null) { - req.setSelector(labelSelector) - } - val mapping = prepareMapping(items, req) - return client.nativeCall(req.build()) - .map(replyProcessor(mapping)) - .reduce(Tuples.of(0, 0, 0), sumStatus) - .map(asStatus) - .doFinally { - mapping.values.forEach { bi -> - bi.onError(RpcException(-32603, "RPC response not received for ${bi.call.method}")) - } - } - .timeout(Duration.ofSeconds(15)) - .toFuture() - } -} \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstream.kt index 45220ee4..f3657007 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstream.kt @@ -28,6 +28,7 @@ import io.emeraldpay.grpc.Chain import io.infinitape.etherjar.domain.BlockHash import io.infinitape.etherjar.domain.TransactionId import io.infinitape.etherjar.rpc.* +import io.infinitape.etherjar.rpc.emerald.EmeraldGrpcTransport import io.infinitape.etherjar.rpc.json.BlockJson import org.slf4j.LoggerFactory import org.springframework.context.Lifecycle @@ -48,7 +49,8 @@ open class GrpcUpstream( private val parentId: String, private val chain: Chain, private val client: ReactorBlockchainGrpc.ReactorBlockchainStub, - private val objectMapper: ObjectMapper + private val objectMapper: ObjectMapper, + private val grpcTransport: EmeraldGrpcTransport ): DefaultUpstream(), Lifecycle { private var allLabels: Collection = ArrayList() @@ -60,13 +62,15 @@ open class GrpcUpstream( private val nodes = AtomicReference(NodeDetailsList()) private val head = Head(this) private var targets: CallMethods? = null - private val grpcTransport = EthereumGrpcTransport(chain, client, objectMapper) private var headSubscription: Disposable? = null open fun createApi(matcher: Selector.Matcher): DirectEthereumApi { val targets = this.getMethods() - val rpcClient = DefaultRpcClient(grpcTransport.withLabels(Selector.extractLabels(matcher))) + val transport = Selector.extractLabels(matcher)?.let { selector -> + grpcTransport.copyWithSelector(selector.asProto()) + } ?: grpcTransport + val rpcClient = DefaultRpcClient(transport) return DirectEthereumApi(rpcClient, objectMapper, targets).let { it.upstream = this it diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt index 7b64fce2..0e479171 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreams.kt @@ -23,6 +23,7 @@ import io.emeraldpay.dshackle.upstream.UpstreamChange import io.emeraldpay.grpc.Chain import io.grpc.ManagedChannelBuilder import io.grpc.netty.NettyChannelBuilder +import io.infinitape.etherjar.rpc.emerald.EmeraldGrpcTransport import io.netty.handler.ssl.* import org.apache.commons.lang3.StringUtils import org.slf4j.LoggerFactory @@ -30,6 +31,7 @@ import reactor.core.publisher.Flux import java.io.File import java.time.Duration import java.util.* +import java.util.concurrent.Executors import java.util.concurrent.locks.ReentrantLock import kotlin.concurrent.withLock @@ -46,6 +48,8 @@ class GrpcUpstreams( private val known = HashMap() private val lock = ReentrantLock() + private var grpcTransport: EmeraldGrpcTransport? = null + fun start(): Flux { val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) { NettyChannelBuilder.forAddress(host, port) @@ -59,6 +63,13 @@ class GrpcUpstreams( val client = ReactorBlockchainGrpc.newReactorStub(channel.build()) this.client = client + var i = 0 + val grpcExecutor = Executors.newFixedThreadPool(8) { r -> Thread(r, "grpc-up-$id-${i++}") }; + this.grpcTransport = EmeraldGrpcTransport.newBuilder() + .forChannel(client.channel) + .setObjectMapper(objectMapper) + .setExecutorService(grpcExecutor) + .build() val updates = Flux.interval(Duration.ZERO, Duration.ofMinutes(1)) .flatMap { @@ -67,6 +78,8 @@ class GrpcUpstreams( processDescription(value) }.doOnError { t -> log.error("Failed to get description from $host:$port", t) + }.doFinally { + grpcExecutor.shutdown() } //TODO subscribe only after receiving details @@ -125,7 +138,7 @@ class GrpcUpstreams( lock.withLock { val current = known[chain] return if (current == null) { - val created = GrpcUpstream(id, chain, client!!, objectMapper) + val created = GrpcUpstream(id, chain, client!!, objectMapper, grpcTransport!!.copyForChain(chain)) known[chain] = created created.start() UpstreamChange(chain, created, UpstreamChange.ChangeType.ADDED) diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcTransportSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcTransportSpec.groovy deleted file mode 100644 index 13e5a5a7..00000000 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcTransportSpec.groovy +++ /dev/null @@ -1,142 +0,0 @@ -/** - * Copyright (c) 2019 ETCDEV GmbH - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package io.emeraldpay.dshackle.upstream.grpc - -import com.fasterxml.jackson.databind.ObjectMapper -import io.emeraldpay.api.proto.BlockchainOuterClass -import io.emeraldpay.api.proto.ReactorBlockchainGrpc -import io.emeraldpay.dshackle.rpc.NativeCall -import io.emeraldpay.dshackle.test.EthereumApiMock -import io.emeraldpay.dshackle.test.EthereumUpstreamMock -import io.emeraldpay.dshackle.test.MockServer -import io.emeraldpay.dshackle.test.TestingCommons -import io.emeraldpay.dshackle.upstream.DirectCallMethods -import io.emeraldpay.dshackle.upstream.QuorumBasedMethods -import io.emeraldpay.dshackle.upstream.Upstreams -import io.emeraldpay.dshackle.upstream.grpc.EthereumGrpcTransport -import io.emeraldpay.grpc.Chain -import io.infinitape.etherjar.rpc.Batch -import io.infinitape.etherjar.rpc.RpcCall -import io.infinitape.etherjar.rpc.RpcClient -import reactor.core.publisher.Flux -import reactor.core.publisher.Mono -import spock.lang.Specification - -class EthereumGrpcTransportSpec extends Specification { - - MockServer mockServer = new MockServer() - ObjectMapper objectMapper = TestingCommons.objectMapper() - def ethereumTargets = new QuorumBasedMethods(objectMapper, Chain.ETHEREUM) - - def "Make simple call"() { - setup: - def otherSideApi = new EthereumApiMock(Mock(RpcClient), objectMapper, Chain.ETHEREUM) - - def callData = [:] - def otherSideUpstreams = Mock(Upstreams) - def otherSideAggr = TestingCommons.aggregatedUpstream( - new EthereumUpstreamMock(Chain.ETHEREUM, otherSideApi, new DirectCallMethods(["eth_test"])) - ) - - def otherSideNativeCall = new NativeCall(otherSideUpstreams, objectMapper) - otherSideApi.upstream = otherSideAggr - - def client = mockServer.clientForServer(new ReactorBlockchainGrpc.BlockchainImplBase() { - @Override - Flux nativeCall(Mono request) { - callData["request"] = request.block() - return otherSideNativeCall.nativeCall(request) - } - }) - - EthereumGrpcTransport transport = new EthereumGrpcTransport(Chain.ETHEREUM, client, objectMapper) - when: - otherSideApi.answer("eth_test", [1], "bar") - def batch = new Batch() - def f = batch.add(RpcCall.create("eth_test", [1])) - def status = transport.execute(batch.items).get() - - then: - 1 * otherSideUpstreams.getUpstream(Chain.ETHEREUM) >> otherSideAggr - status.failed == 0 - status.succeed == 1 - status.total == 1 - callData.request != null - with((BlockchainOuterClass.NativeCallRequest)callData.request) { - chain.number == Chain.ETHEREUM.id - itemsCount == 1 - with(getItems(0)) { - method == "eth_test" - payload.toStringUtf8() == "[1]" - } - } - f.get() == "bar" - } - - def "Make few calls"() { - setup: - def otherSideApi = new EthereumApiMock(Mock(RpcClient), objectMapper, Chain.ETHEREUM) - - def callData = [:] - def otherSideUpstreams = Mock(Upstreams) - def otherSideAggr = TestingCommons.aggregatedUpstream( - new EthereumUpstreamMock(Chain.ETHEREUM, otherSideApi, new DirectCallMethods(["eth_test", "eth_test2"])) - ) - def otherSideNativeCall = new NativeCall(otherSideUpstreams, objectMapper) - otherSideApi.upstream = otherSideAggr - - def client = mockServer.clientForServer(new ReactorBlockchainGrpc.BlockchainImplBase() { - @Override - Flux nativeCall(Mono request) { - callData["request"] = request.block() - return otherSideNativeCall.nativeCall(request) - } - - }) - - EthereumGrpcTransport transport = new EthereumGrpcTransport(Chain.ETHEREUM, client, objectMapper) - when: - otherSideApi.answer("eth_test", [1], "bar") - otherSideApi.answer("eth_test2", [2, "3"], "baz") - - def batch = new Batch() - def f1 = batch.add(RpcCall.create("eth_test", [1])) - def f2 = batch.add(RpcCall.create("eth_test2", [2, "3"])) - def status = transport.execute(batch.items).get() - - then: - 1 * otherSideUpstreams.getUpstream(Chain.ETHEREUM) >> otherSideAggr - status.failed == 0 - status.succeed == 2 - status.total == 2 - callData.request != null - with((BlockchainOuterClass.NativeCallRequest)callData.request) { - chain.number == Chain.ETHEREUM.id - itemsCount == 2 - with(getItems(0)) { - method == "eth_test" - payload.toStringUtf8() == "[1]" - } - with(getItems(1)) { - method == "eth_test2" - payload.toStringUtf8() == "[2,\"3\"]" - } - } - f1.get() == "bar" - f2.get() == "baz" - } - -} diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamSpec.groovy index 27d6badb..5bfc4be6 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/grpc/GrpcUpstreamSpec.groovy @@ -28,12 +28,15 @@ import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.grpc.Chain import io.grpc.stub.StreamObserver import io.infinitape.etherjar.domain.BlockHash +import io.infinitape.etherjar.rpc.JacksonRpcConverter import io.infinitape.etherjar.rpc.RpcClient +import io.infinitape.etherjar.rpc.emerald.EmeraldGrpcTransport import io.infinitape.etherjar.rpc.json.BlockJson import spock.lang.Specification import java.time.Duration import java.util.concurrent.CompletableFuture +import java.util.concurrent.Executors class GrpcUpstreamSpec extends Specification { @@ -70,7 +73,8 @@ class GrpcUpstreamSpec extends Specification { ) } }) - def upstream = new GrpcUpstream("test", chain, client, objectMapper) + def transport = EmeraldGrpcTransport.newBuilder().forChannel(client.channel).build() + def upstream = new GrpcUpstream("test", chain, client, objectMapper, transport) upstream.setLag(0) upstream.init(BlockchainOuterClass.DescribeChain.newBuilder() .addAllSupportedMethods(["eth_getBlockByHash"]) @@ -129,7 +133,8 @@ class GrpcUpstreamSpec extends Specification { finished.complete(true) } }) - def upstream = new GrpcUpstream("test", chain, client, objectMapper) + def transport = EmeraldGrpcTransport.newBuilder().forChannel(client.channel).build() + def upstream = new GrpcUpstream("test", chain, client, objectMapper, transport) upstream.setLag(0) upstream.init(BlockchainOuterClass.DescribeChain.newBuilder() .addAllSupportedMethods(["eth_getBlockByHash"]) @@ -189,7 +194,8 @@ class GrpcUpstreamSpec extends Specification { finished.complete(true) } }) - def upstream = new GrpcUpstream("test", chain, client, objectMapper) + def transport = EmeraldGrpcTransport.newBuilder().forChannel(client.channel).build() + def upstream = new GrpcUpstream("test", chain, client, objectMapper, transport) upstream.setLag(0) upstream.init(BlockchainOuterClass.DescribeChain.newBuilder() .addAllSupportedMethods(["eth_getBlockByHash"])