From d140546d0c55e0baf9ccd3ba85401b6c0e96a8f7 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Tue, 29 Nov 2022 19:43:20 +0400 Subject: [PATCH 1/2] Added forwarding of selectors to grpc upstreams --- emerald-grpc | 2 +- .../io/emeraldpay/dshackle/rpc/NativeCall.kt | 16 ++- .../io/emeraldpay/dshackle/rpc/Selectors.kt | 43 +++++++ .../upstream/grpc/BitcoinGrpcUpstream.kt | 3 +- .../upstream/grpc/EthereumGrpcUpstream.kt | 3 +- .../upstream/grpc/EthereumPosGrpcUpstream.kt | 3 +- .../upstream/rpcclient/JsonRpcGrpcClient.kt | 12 +- .../upstream/rpcclient/JsonRpcRequest.kt | 13 +- .../dshackle/rpc/NativeCallSpec.groovy | 6 +- .../upstream/ethereum/WsConnectionSpec.groovy | 6 +- .../emeraldpay/dshackle/rpc/SelectorsTest.kt | 117 ++++++++++++++++++ 11 files changed, 193 insertions(+), 31 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/rpc/Selectors.kt create mode 100644 src/test/kotlin/io/emeraldpay/dshackle/rpc/SelectorsTest.kt diff --git a/emerald-grpc b/emerald-grpc index 441f828e..05543a33 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit 441f828eb6d4a0d5d91d6575911eec80271a24aa +Subproject commit 05543a33823e38448ffacca06c445a05fc863cb2 diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt index 40a180b0..b4d58cea 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt @@ -242,6 +242,8 @@ open class NativeCall( val requestDecorator = getRequestDecorator(requestItem.method) val resultDecorator = getResultDecorator(requestItem.method) + val selector = request.takeIf { it.hasSelector() }?.let { Selectors.keepForwarded(it.selector) } + ValidCallContext( requestItem.id, nonce, @@ -250,7 +252,8 @@ open class NativeCall( callQuorum, RawCallDetails(method, params), requestDecorator, - resultDecorator + resultDecorator, + selector ) } } @@ -267,7 +270,7 @@ open class NativeCall( fun fetch(ctx: ValidCallContext): Mono { return ctx.upstream.getRoutedApi(ctx.matcher) .flatMap { api -> - api.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce)) + api.read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce, ctx.forwardedSelector)) .flatMap(JsonRpcResponse::requireResult) .map { if (ctx.nonce != null) { @@ -300,7 +303,7 @@ open class NativeCall( AtomicInteger(-1) } return reader - .read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce)) + .read(JsonRpcRequest(ctx.payload.method, ctx.payload.params, ctx.nonce, ctx.forwardedSelector)) .map { val bytes = ctx.resultDecorator.processResult(it) CallResult(ctx.id, ctx.nonce, bytes, null, it.signature) @@ -418,7 +421,8 @@ open class NativeCall( val callQuorum: CallQuorum, val payload: T, val requestDecorator: RequestDecorator, - val resultDecorator: ResultDecorator + val resultDecorator: ResultDecorator, + val forwardedSelector: BlockchainOuterClass.Selector? ) : CallContext { constructor( @@ -428,7 +432,7 @@ open class NativeCall( matcher: Selector.Matcher, callQuorum: CallQuorum, payload: T - ) : this(id, nonce, upstream, matcher, callQuorum, payload, NoneRequestDecorator(), NoneResultDecorator()) + ) : this(id, nonce, upstream, matcher, callQuorum, payload, NoneRequestDecorator(), NoneResultDecorator(), null) override fun isValid(): Boolean { return true @@ -443,7 +447,7 @@ open class NativeCall( } fun withPayload(payload: X): ValidCallContext { - return ValidCallContext(id, nonce, upstream, matcher, callQuorum, payload, requestDecorator, resultDecorator) + return ValidCallContext(id, nonce, upstream, matcher, callQuorum, payload, requestDecorator, resultDecorator, forwardedSelector) } fun getApis(): ApiSource { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/Selectors.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Selectors.kt new file mode 100644 index 00000000..df3a8405 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/Selectors.kt @@ -0,0 +1,43 @@ +package io.emeraldpay.dshackle.rpc + +import io.emeraldpay.api.proto.BlockchainOuterClass.* + +object Selectors { + + private fun processMultiple(original: List, creator: (List) -> Selector.Builder): Selector? { + val newSelectors = original.mapNotNull { keepForwarded(it) } + return if (newSelectors.size > 1) { + return creator(newSelectors).setShouldBeForwarded(true).build() + } else if (newSelectors.size == 1) { + newSelectors.first() + } else { + null + } + } + + fun keepForwarded(selector: Selector): Selector? { + if (!selector.shouldBeForwarded) return null + + if (selector.hasOrSelector()) { + return processMultiple(selector.orSelector.selectorsList) { + Selector.newBuilder() + .setOrSelector(OrSelector.newBuilder().addAllSelectors(it)) + } + } else if (selector.hasAndSelector()) { + return processMultiple(selector.andSelector.selectorsList) { + Selector.newBuilder() + .setAndSelector(AndSelector.newBuilder().addAllSelectors(it)) + } + } else if (selector.hasNotSelector()) { + val inner = keepForwarded(selector.notSelector.selector) + return inner?.let { + Selector.newBuilder() + .setNotSelector(NotSelector.newBuilder().setSelector(it).build()) + .setShouldBeForwarded(true) + .build() + } + } else { + return selector + } + } +} diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt index 09ac0b5d..6237777c 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/BitcoinGrpcUpstream.kt @@ -26,7 +26,6 @@ import io.emeraldpay.dshackle.reader.Reader import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Lifecycle -import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.bitcoin.BitcoinUpstream @@ -67,7 +66,7 @@ class BitcoinGrpcUpstream( } private val extractBlock = ExtractBlock() - private val defaultReader: Reader = client.forSelector(Selector.empty) + private val defaultReader: Reader = client.getReader() private val blockConverter: Function = Function { value -> val block = BlockContainer( value.height, diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt index 6c40e2a0..4ea17f47 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumGrpcUpstream.kt @@ -28,7 +28,6 @@ import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Lifecycle -import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.calls.CallMethods @@ -106,7 +105,7 @@ open class EthereumGrpcUpstream( private val grpcHead = GrpcHead(chain, this, remote, blockConverter, reloadBlock, MostWorkForkChoice()) private var capabilities: Set = emptySet() - private val defaultReader: Reader = client.forSelector(Selector.empty) + private val defaultReader: Reader = client.getReader() var timeout = Defaults.timeout override fun getBlockchainApi(): ReactorBlockchainGrpc.ReactorBlockchainStub { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt index f180bb6a..1424f951 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/grpc/EthereumPosGrpcUpstream.kt @@ -28,7 +28,6 @@ import io.emeraldpay.dshackle.startup.QuorumForLabels import io.emeraldpay.dshackle.upstream.Capability import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Lifecycle -import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.UpstreamAvailability import io.emeraldpay.dshackle.upstream.calls.CallMethods @@ -106,7 +105,7 @@ open class EthereumPosGrpcUpstream( private val grpcHead = GrpcHead(chain, this, remote, blockConverter, reloadBlock, NoChoiceWithPriorityForkChoice(nodeRating)) private var capabilities: Set = emptySet() - private val defaultReader: Reader = client.forSelector(Selector.empty) + private val defaultReader: Reader = client.getReader() var timeout = Defaults.timeout override fun start() { diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt index 08f7e0f4..8acff9e2 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt @@ -22,7 +22,6 @@ import io.emeraldpay.api.proto.ReactorBlockchainGrpc import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.reader.Reader -import io.emeraldpay.dshackle.upstream.Selector import io.emeraldpay.dshackle.upstream.signature.ResponseSigner import io.emeraldpay.etherjar.rpc.RpcException import io.emeraldpay.etherjar.rpc.RpcResponseError @@ -40,14 +39,13 @@ class JsonRpcGrpcClient( private val log = LoggerFactory.getLogger(JsonRpcGrpcClient::class.java) } - fun forSelector(matcher: Selector.Matcher): Reader { - return Executor(stub, chain, matcher, metrics) + fun getReader(): Reader { + return Executor(stub, chain, metrics) } class Executor( private val stub: ReactorBlockchainGrpc.ReactorBlockchainStub, private val chain: Chain, - private val matcher: Selector.Matcher, private val metrics: RpcMetrics ) : Reader { @@ -56,11 +54,7 @@ class JsonRpcGrpcClient( val req = BlockchainOuterClass.NativeCallRequest.newBuilder() .setChainValue(chain.id) - if (matcher != Selector.empty) { - Selector.extractLabels(matcher)?.asProto().let { - req.setSelector(it) - } - } + key.selector?.let { req.selector = it } val reqItem = BlockchainOuterClass.NativeCallItem.newBuilder() .setId(1) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcRequest.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcRequest.kt index 86109d6a..be5faf47 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcRequest.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcRequest.kt @@ -19,16 +19,23 @@ import com.fasterxml.jackson.core.JsonParser import com.fasterxml.jackson.databind.DeserializationContext import com.fasterxml.jackson.databind.JsonDeserializer import com.fasterxml.jackson.databind.JsonNode +import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.dshackle.Global data class JsonRpcRequest( val method: String, val params: List, val id: Int, - val nonce: Long? + val nonce: Long?, + val selector: BlockchainOuterClass.Selector? ) { - @JvmOverloads constructor(method: String, params: List, nonce: Long? = null) : this(method, params, 1, nonce) + @JvmOverloads constructor( + method: String, + params: List, + nonce: Long? = null, + selectors: BlockchainOuterClass.Selector? = null + ) : this(method, params, 1, nonce, selectors) fun toJson(): ByteArray { val json = mapOf( @@ -63,7 +70,7 @@ data class JsonRpcRequest( throw IllegalStateException("Unsupported param type: ${it.asToken()}") } } - return JsonRpcRequest(method, params, id, null) + return JsonRpcRequest(method, params, id, null, null) } } } diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy index 2bff6234..26165cfa 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy @@ -543,7 +543,7 @@ class NativeCallSpec extends Specification { def nativeCall = nativeCall() def ctx = new NativeCall.ValidCallContext(1, null, Stub(Multistream), Selector.empty, new AlwaysQuorum(), new NativeCall.RawCallDetails("eth_getFilterUpdates", '["0xabcd"]'), - new NativeCall.WithFilterIdDecorator(), new NativeCall.NoneResultDecorator()) + new NativeCall.WithFilterIdDecorator(), new NativeCall.NoneResultDecorator(), null) when: def act = nativeCall.parseParams(ctx) then: @@ -564,7 +564,7 @@ class NativeCallSpec extends Specification { } def call = new NativeCall.ValidCallContext(1, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum, new NativeCall.ParsedCallDetails("eth_getFilterChanges", []), - new NativeCall.WithFilterIdDecorator(), new NativeCall.CreateFilterDecorator()) + new NativeCall.WithFilterIdDecorator(), new NativeCall.CreateFilterDecorator(), null) when: def resp = nativeCall.executeOnRemote(call).block(Duration.ofSeconds(1)) @@ -586,7 +586,7 @@ class NativeCallSpec extends Specification { } def call = new NativeCall.ValidCallContext(1, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum, new NativeCall.ParsedCallDetails("eth_getFilterChanges", []), - new NativeCall.WithFilterIdDecorator(), new NativeCall.CreateFilterDecorator()) + new NativeCall.WithFilterIdDecorator(), new NativeCall.CreateFilterDecorator(), null) when: def resp = nativeCall.executeOnRemote(call).block(Duration.ofSeconds(1)) diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy index 5fc78e37..18350845 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy @@ -84,7 +84,7 @@ class WsConnectionSpec extends Specification { when: Flux.from(ws.handle(wsApiMock.inbound, wsApiMock.outbound)).subscribe() - def act = ws.call(new JsonRpcRequest("eth_getTransactionByHash", ["0x3ec2ebf5d0ec474d0ac6bc50d2770d8409ad76e119968e7919f85d5ec8915200"], 15, null)) + def act = ws.call(new JsonRpcRequest("eth_getTransactionByHash", ["0x3ec2ebf5d0ec474d0ac6bc50d2770d8409ad76e119968e7919f85d5ec8915200"], 15, null, null)) then: StepVerifier.create(act) @@ -106,7 +106,7 @@ class WsConnectionSpec extends Specification { when: Flux.from(ws.handle(wsApiMock.inbound, wsApiMock.outbound)).subscribe() - def act = ws.call(new JsonRpcRequest("eth_getTransactionByHash", ["0x3ec2ebf5d0ec474d0ac6bc50d2770d8409ad76e119968e7919f85d5ec8915200"], 15, null)) + def act = ws.call(new JsonRpcRequest("eth_getTransactionByHash", ["0x3ec2ebf5d0ec474d0ac6bc50d2770d8409ad76e119968e7919f85d5ec8915200"], 15, null, null)) then: StepVerifier.create(act) @@ -130,7 +130,7 @@ class WsConnectionSpec extends Specification { when: Flux.from(ws.handle(wsApiMock.inbound, wsApiMock.outbound)).subscribe() - def act = ws.call(new JsonRpcRequest("eth_getTransactionByHash", ["0x3ec2ebf5d0ec474d0ac6bc50d2770d8409ad76e119968e7919f85d5ec8915200"], 15, null)) + def act = ws.call(new JsonRpcRequest("eth_getTransactionByHash", ["0x3ec2ebf5d0ec474d0ac6bc50d2770d8409ad76e119968e7919f85d5ec8915200"], 15, null, null)) then: StepVerifier.create(act) diff --git a/src/test/kotlin/io/emeraldpay/dshackle/rpc/SelectorsTest.kt b/src/test/kotlin/io/emeraldpay/dshackle/rpc/SelectorsTest.kt new file mode 100644 index 00000000..2a7e9a2f --- /dev/null +++ b/src/test/kotlin/io/emeraldpay/dshackle/rpc/SelectorsTest.kt @@ -0,0 +1,117 @@ +package io.emeraldpay.dshackle.rpc + +import io.emeraldpay.api.proto.BlockchainOuterClass.* +import org.junit.jupiter.api.Assertions.* +import org.junit.jupiter.params.ParameterizedTest +import org.junit.jupiter.params.provider.Arguments +import org.junit.jupiter.params.provider.MethodSource +import java.util.stream.Stream + +internal class SelectorsTest { + + companion object { + @JvmStatic + fun data(): Stream { + val leafLabelSelector = LabelSelector.newBuilder().build() + + return Stream.of( + Arguments.of( + Selector.newBuilder().setLabelSelector(leafLabelSelector).build(), + null + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setLabelSelector(leafLabelSelector).build(), + Selector.newBuilder().setShouldBeForwarded(true) + .setLabelSelector(leafLabelSelector).build() + ), + Arguments.of( + Selector.newBuilder().setOrSelector(OrSelector.newBuilder()).build(), + null + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setOrSelector(OrSelector.newBuilder()).build(), + null + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setOrSelector( + OrSelector.newBuilder().addSelectors( + Selector.newBuilder().setLabelSelector(leafLabelSelector) + ) + ).build(), + null + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setOrSelector( + OrSelector.newBuilder().addSelectors( + Selector.newBuilder().setShouldBeForwarded(true) + .setLabelSelector(leafLabelSelector) + ) + ).build(), + Selector.newBuilder().setShouldBeForwarded(true) + .setLabelSelector(leafLabelSelector).build() + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setOrSelector( + OrSelector.newBuilder() + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + .addSelectors(Selector.newBuilder().setLabelSelector(leafLabelSelector)) + ).build(), + Selector.newBuilder().setShouldBeForwarded(true) + .setOrSelector( + OrSelector.newBuilder() + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + ).build() + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setAndSelector( + AndSelector.newBuilder() + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + .addSelectors(Selector.newBuilder().setLabelSelector(leafLabelSelector)) + ).build(), + Selector.newBuilder().setShouldBeForwarded(true) + .setAndSelector( + AndSelector.newBuilder() + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + .addSelectors(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + ).build() + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setNotSelector( + NotSelector.newBuilder() + .setSelector(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + ).build(), + Selector.newBuilder().setShouldBeForwarded(true) + .setNotSelector( + NotSelector.newBuilder() + .setSelector(Selector.newBuilder().setShouldBeForwarded(true).setLabelSelector(leafLabelSelector)) + ).build() + ), + Arguments.of( + Selector.newBuilder().setShouldBeForwarded(true) + .setNotSelector( + NotSelector.newBuilder() + .setSelector(Selector.newBuilder().setLabelSelector(leafLabelSelector)) + ).build(), + null + ) + + ) + } + } + + @ParameterizedTest + @MethodSource("data") + fun testKeepForwarded(input: Selector, expected: Selector?) { + assertEquals(expected, Selectors.keepForwarded(input)) + } +} From c4669b8e05a5eb6dfd09661770f53b79257b1ac4 Mon Sep 17 00:00:00 2001 From: a10zn8 Date: Wed, 30 Nov 2022 13:59:12 +0400 Subject: [PATCH 2/2] submodule update --- emerald-grpc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/emerald-grpc b/emerald-grpc index 05543a33..7dea528b 160000 --- a/emerald-grpc +++ b/emerald-grpc @@ -1 +1 @@ -Subproject commit 05543a33823e38448ffacca06c445a05fc863cb2 +Subproject commit 7dea528b2be574104a26db1abcd71fbaa32dc9b0