solana support (#339)

This commit is contained in:
a10zn8
2023-11-13 17:01:58 +03:00
committed by GitHub
parent e5f844b091
commit 2e80da58df
39 changed files with 513 additions and 249 deletions

View File

@@ -117,6 +117,7 @@ open class CodeGen(private val config: ChainsConfig) {
"bitcoin" -> "BlockchainType.BITCOIN" "bitcoin" -> "BlockchainType.BITCOIN"
"starknet" -> "BlockchainType.STARKNET" "starknet" -> "BlockchainType.STARKNET"
"polkadot" -> "BlockchainType.POLKADOT" "polkadot" -> "BlockchainType.POLKADOT"
"solana" -> "BlockchainType.SOLANA"
else -> throw IllegalArgumentException("unknown blockchain type $type") else -> throw IllegalArgumentException("unknown blockchain type $type")
} }
} }

View File

@@ -1,5 +1,5 @@
package io.emeraldpay.dshackle package io.emeraldpay.dshackle
enum class BlockchainType { enum class BlockchainType {
UNKNOWN, BITCOIN, ETHEREUM, STARKNET, POLKADOT; UNKNOWN, BITCOIN, ETHEREUM, STARKNET, POLKADOT, SOLANA;
} }

View File

@@ -664,3 +664,26 @@ chain-settings:
short-names: [ vara-testnet ] short-names: [ vara-testnet ]
chain-id: 0x0 chain-id: 0x0
grpcId: 10036 grpcId: 10036
- id: solana
label: solana
type: solana
settings:
expected-block-time: 1s
options:
validate-peers: false
lags:
syncing: 20
lagging: 10
chains:
- id: Mainnet
priority: 1
code: SOLANA_MAINNET
short-names: [ solana ]
chain-id: 0x0
grpcId: 1028
- id: Testnet
priority: 1
code: SOLANA_TESTMET
short-names: [ solana-testnet ]
chain-id: 0x0
grpcId: 10037

View File

@@ -19,6 +19,7 @@ package io.emeraldpay.dshackle.data
import io.emeraldpay.dshackle.upstream.ethereum.json.BlockJson import io.emeraldpay.dshackle.upstream.ethereum.json.BlockJson
import io.emeraldpay.etherjar.domain.BlockHash import io.emeraldpay.etherjar.domain.BlockHash
import org.bouncycastle.util.encoders.Hex import org.bouncycastle.util.encoders.Hex
import java.util.Base64
class BlockId( class BlockId(
value: ByteArray, value: ByteArray,
@@ -55,5 +56,10 @@ class BlockId(
val bytes = Hex.decode(even) val bytes = Hex.decode(even)
return BlockId(bytes) return BlockId(bytes)
} }
fun fromBase64(id: String): BlockId {
val bytes = Base64.getDecoder().decode(id)
return BlockId(bytes)
}
} }
} }

View File

@@ -74,7 +74,12 @@ open class NativeSubscribe(
val publisher = multistream.tryProxySubscribe(matcher, request) ?: run { val publisher = multistream.tryProxySubscribe(matcher, request) ?: run {
val method = request.method val method = request.method
val params: Any? = request.payload?.takeIf { !it.isEmpty }?.let { val params: Any? = request.payload?.takeIf { !it.isEmpty }?.let {
objectMapper.readValue(it.newInput(), Map::class.java) val raw = it.toStringUtf8()
if (raw.startsWith("{")) {
objectMapper.readValue(it.newInput(), Map::class.java)
} else {
objectMapper.readValue(it.newInput(), List::class.java)
}
} }
subscribe(chain, method, params, matcher) subscribe(chain, method, params, matcher)
} }

View File

@@ -164,7 +164,6 @@ class SubscribeNodeStatus(
.build() .build()
}, },
) )
.addAllSupportedSubscriptions(up.getSubscriptionTopics())
.addAllSupportedMethods(up.getMethods().getSupportedMethods()) .addAllSupportedMethods(up.getMethods().getSupportedMethods())
(up as? GrpcUpstream)?.let { (up as? GrpcUpstream)?.let {
it.getBuildInfo().version?.let { version -> it.getBuildInfo().version?.let { version ->

View File

@@ -85,7 +85,6 @@ open class GenericUpstreamCreator(
connectorFactory, connectorFactory,
cs::validator, cs::validator,
cs::labelDetector, cs::labelDetector,
cs::subscriptionTopics,
) )
upstream.start() upstream.start()

View File

@@ -3,6 +3,7 @@ package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.BlockchainType.BITCOIN import io.emeraldpay.dshackle.BlockchainType.BITCOIN
import io.emeraldpay.dshackle.BlockchainType.ETHEREUM import io.emeraldpay.dshackle.BlockchainType.ETHEREUM
import io.emeraldpay.dshackle.BlockchainType.POLKADOT import io.emeraldpay.dshackle.BlockchainType.POLKADOT
import io.emeraldpay.dshackle.BlockchainType.SOLANA
import io.emeraldpay.dshackle.BlockchainType.STARKNET import io.emeraldpay.dshackle.BlockchainType.STARKNET
import io.emeraldpay.dshackle.BlockchainType.UNKNOWN import io.emeraldpay.dshackle.BlockchainType.UNKNOWN
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
@@ -27,6 +28,7 @@ class CallTargetsHolder {
ETHEREUM -> DefaultEthereumMethods(chain) ETHEREUM -> DefaultEthereumMethods(chain)
STARKNET -> DefaultStarknetMethods(chain) STARKNET -> DefaultStarknetMethods(chain)
POLKADOT -> DefaultPolkadotMethods() POLKADOT -> DefaultPolkadotMethods()
SOLANA -> DefaultSolanaMethods()
UNKNOWN -> throw IllegalArgumentException("unknown chain") UNKNOWN -> throw IllegalArgumentException("unknown chain")
} }
callTargets[chain] = created callTargets[chain] = created

View File

@@ -0,0 +1,110 @@
package io.emeraldpay.dshackle.upstream
import io.emeraldpay.dshackle.quorum.AlwaysQuorum
import io.emeraldpay.dshackle.quorum.BroadcastQuorum
import io.emeraldpay.dshackle.quorum.CallQuorum
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.etherjar.rpc.RpcException
class DefaultSolanaMethods : CallMethods {
companion object {
val subs = setOf(
"accountSubscribe" to "accountUnsubscribe",
"blockSubscribe" to "blockUnsubscribe",
"logsSubscribe" to "logsUnsubscribe",
"programSubscribe" to "programUnsubscribe",
"signatureSubscribe" to "signatureUnsubscribe",
"slotSubscribe" to "slotUnsubscribe",
)
}
private val all = setOf(
"getAccountInfo",
"getBalance",
"getBlock",
"getBlockHeight",
"getBlockProduction",
"getBlockCommitment",
"getBlocks",
"getBlocksWithLimit",
"getBlockTime",
"getClusterNodes",
"getEpochInfo",
"getEpochSchedule",
"getFeeForMessage",
"getFirstAvailableBlock",
"getGenesisHash",
"getHealth",
"getHighestSnapshotSlot",
"getInflationGovernor",
"getInflationRate",
"getInflationReward",
"getLargestAccounts",
"getLatestBlockhash",
"getLeaderSchedule",
"getMaxRetransmitSlot",
"getMaxShredInsertSlot",
"getMinimumBalanceForRentExemption",
"getMultipleAccounts",
"getProgramAccounts",
"getRecentPerformanceSamples",
"getRecentPrioritizationFees",
"getSignaturesForAddress",
"getSignatureStatuses",
"getSlot",
"getSlotLeader",
"getSlotLeaders",
"getStakeActivation",
"getStakeMinimumDelegation",
"getSupply",
"getTokenAccountBalance",
"getTokenAccountsByDelegate",
"getTokenAccountsByOwner",
"getTokenLargestAccounts",
"getTokenSupply",
"getTransaction",
"getTransactionCount",
"getVersion",
"getVoteAccounts",
"isBlockhashValid",
"minimumLedgerSlot",
"requestAirdrop",
"simulateTransaction",
)
private val add = setOf(
"sendTransaction",
)
private val allowedMethods: Set<String> = all + add
override fun createQuorumFor(method: String): CallQuorum {
return when {
add.contains(method) -> BroadcastQuorum()
all.contains(method) -> AlwaysQuorum()
else -> AlwaysQuorum()
}
}
override fun isCallable(method: String): Boolean {
return allowedMethods.contains(method)
}
override fun isHardcoded(method: String): Boolean {
return false
}
override fun executeHardcoded(method: String): ByteArray {
throw RpcException(-32601, "Method not found")
}
override fun getGroupMethods(groupName: String): Set<String> =
when (groupName) {
"default" -> getSupportedMethods()
else -> emptyList()
}.toSet()
override fun getSupportedMethods(): Set<String> {
return allowedMethods.toSortedSet()
}
}

View File

@@ -21,5 +21,5 @@ package io.emeraldpay.dshackle.upstream
interface IngressSubscription { interface IngressSubscription {
fun getAvailableTopics(): List<String> fun getAvailableTopics(): List<String>
fun <T> get(topic: String): SubscriptionConnect<T>? fun <T> get(topic: String, params: Any?): SubscriptionConnect<T>?
} }

View File

@@ -94,7 +94,7 @@ abstract class Multistream(
.multicast() .multicast()
.directBestEffort<UpstreamChangeEvent>() .directBestEffort<UpstreamChangeEvent>()
override fun getSubscriptionTopics(): List<String> { fun getSubscriptionTopics(): List<String> {
return getEgressSubscription().getAvailableTopics() return getEgressSubscription().getAvailableTopics()
} }

View File

@@ -25,7 +25,7 @@ open class NoIngressSubscription : IngressSubscription {
return listOf() return listOf()
} }
override fun <T> get(topic: String): SubscriptionConnect<T>? { override fun <T> get(topic: String, params: Any?): SubscriptionConnect<T>? {
return null return null
} }
} }

View File

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

View File

@@ -59,10 +59,6 @@ open class BitcoinRpcUpstream(
return directApi return directApi
} }
override fun getSubscriptionTopics(): List<String> {
return listOf()
}
override fun getLabels(): Collection<UpstreamsConfig.Labels> { override fun getLabels(): Collection<UpstreamsConfig.Labels> {
return listOf(UpstreamsConfig.Labels()) return listOf(UpstreamsConfig.Labels())
} }

View File

@@ -29,12 +29,11 @@ class DefaultPolkadotMethods : CallMethods {
companion object { companion object {
val subs = setOf( val subs = setOf(
Pair("subscribe_newHead", "unsubscribe_newHead"), "subscribe_newHead" to "unsubscribe_newHead",
Pair("chain_subscribeAllHeads", "chain_unsubscribeAllHeads"), "chain_subscribeAllHeads" to "chain_unsubscribeAllHeads",
Pair("chain_subscribeFinalizedHeads", "chain_unsubscribeFinalizedHeads"), "chain_subscribeFinalizedHeads" to "chain_unsubscribeFinalizedHeads",
Pair("chain_subscribeNewHeads", "chain_unsubscribeNewHeads"), "chain_subscribeNewHeads" to "chain_unsubscribeNewHeads",
Pair("chain_subscribeRuntimeVersion", "chain_unsubscribeNewHeads"), "chain_subscribeRuntimeVersion" to "chain_unsubscribeRuntimeVersion",
Pair("chain_subscribeRuntimeVersion", "chain_unsubscribeRuntimeVersion"),
) )
} }

View File

@@ -7,7 +7,6 @@ import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.foundation.ChainOptions.Options import io.emeraldpay.dshackle.foundation.ChainOptions.Options
import io.emeraldpay.dshackle.reader.JsonRpcReader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.CachingReader import io.emeraldpay.dshackle.upstream.CachingReader
import io.emeraldpay.dshackle.upstream.Capability
import io.emeraldpay.dshackle.upstream.EgressSubscription import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.Head import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.IngressSubscription import io.emeraldpay.dshackle.upstream.IngressSubscription
@@ -23,15 +22,15 @@ import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumLabelsDetector
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumWsIngressSubscription import io.emeraldpay.dshackle.upstream.ethereum.subscribe.EthereumWsIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes import io.emeraldpay.dshackle.upstream.ethereum.subscribe.NoPendingTxes
import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource import io.emeraldpay.dshackle.upstream.ethereum.subscribe.PendingTxesSource
import io.emeraldpay.dshackle.upstream.generic.AbstractPollChainSpecific
import io.emeraldpay.dshackle.upstream.generic.CachingReaderBuilder import io.emeraldpay.dshackle.upstream.generic.CachingReaderBuilder
import io.emeraldpay.dshackle.upstream.generic.ChainSpecific
import io.emeraldpay.dshackle.upstream.generic.GenericUpstream import io.emeraldpay.dshackle.upstream.generic.GenericUpstream
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import org.springframework.cloud.sleuth.Tracer import org.springframework.cloud.sleuth.Tracer
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import reactor.core.scheduler.Scheduler import reactor.core.scheduler.Scheduler
object EthereumChainSpecific : ChainSpecific { object EthereumChainSpecific : AbstractPollChainSpecific() {
override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer { override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer {
return BlockContainer.fromEthereumJson(data, upstreamId) return BlockContainer.fromEthereumJson(data, upstreamId)
} }
@@ -58,6 +57,7 @@ object EthereumChainSpecific : ChainSpecific {
val pendingTxes: PendingTxesSource = (ms.getAll()) val pendingTxes: PendingTxesSource = (ms.getAll())
.filter { it is GenericUpstream } .filter { it is GenericUpstream }
.map { it as GenericUpstream } .map { it as GenericUpstream }
.filter { it.getIngressSubscription() is EthereumIngressSubscription }
.mapNotNull { .mapNotNull {
(it.getIngressSubscription() as EthereumIngressSubscription).getPendingTxes() (it.getIngressSubscription() as EthereumIngressSubscription).getPendingTxes()
}.let { }.let {
@@ -90,15 +90,6 @@ object EthereumChainSpecific : ChainSpecific {
return EthereumLabelsDetector(reader, chain) return EthereumLabelsDetector(reader, chain)
} }
override fun subscriptionTopics(upstream: GenericUpstream): List<String> {
val subs = if (upstream.getCapabilities().contains(Capability.WS_HEAD)) {
listOf(EthereumEgressSubscription.METHOD_NEW_HEADS, EthereumEgressSubscription.METHOD_LOGS)
} else {
listOf()
}
return upstream.getIngressSubscription().getAvailableTopics().plus(subs).toSet().toList()
}
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription { override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
return EthereumWsIngressSubscription(ws) return EthereumWsIngressSubscription(ws)
} }

View File

@@ -93,7 +93,7 @@ class GenericWsHead(
} }
.timeout(Duration.ofSeconds(60), Mono.error(RuntimeException("No response from subscribe to newHeads"))) .timeout(Duration.ofSeconds(60), Mono.error(RuntimeException("No response from subscribe to newHeads")))
.onErrorResume { .onErrorResume {
log.error("Error getting heads for $upstreamId - ${it.message}") log.error("Error getting heads for $upstreamId", it)
subscribed = false subscribed = false
unsubscribe() unsubscribe()
} }

View File

@@ -1,56 +0,0 @@
/**
* Copyright (c) 2022 EmeraldPay, Inc
*
* 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.ethereum.subscribe
import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.api.proto.ReactorBlockchainGrpc
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.SubscriptionConnect
import io.emeraldpay.dshackle.upstream.ethereum.EthereumEgressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.EthereumIngressSubscription
import org.slf4j.LoggerFactory
class EthereumDshackleIngressSubscription(
private val blockchain: Chain,
private val conn: ReactorBlockchainGrpc.ReactorBlockchainStub,
) : IngressSubscription, EthereumIngressSubscription {
companion object {
private val log = LoggerFactory.getLogger(EthereumDshackleIngressSubscription::class.java)
}
private val pendingTxes = DshacklePendingTxesSource(blockchain, conn)
override fun getAvailableTopics(): List<String> {
return listOf(EthereumEgressSubscription.METHOD_PENDING_TXES)
}
override fun <T> get(topic: String): SubscriptionConnect<T>? {
if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) {
return pendingTxes as SubscriptionConnect<T>
}
return null
}
fun update(conf: BlockchainOuterClass.DescribeChain) {
pendingTxes.update(conf)
}
override fun getPendingTxes(): PendingTxesSource? {
return pendingTxes
}
}

View File

@@ -32,7 +32,7 @@ class EthereumWsIngressSubscription(
} }
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
override fun <T> get(topic: String): SubscriptionConnect<T>? { override fun <T> get(topic: String, params: Any?): SubscriptionConnect<T>? {
if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) { if (topic == EthereumEgressSubscription.METHOD_PENDING_TXES) {
return pendingTxes as SubscriptionConnect<T> return pendingTxes as SubscriptionConnect<T>
} }

View File

@@ -0,0 +1,79 @@
package io.emeraldpay.dshackle.upstream.generic
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches
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.JsonRpcReader
import io.emeraldpay.dshackle.upstream.CachingReader
import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.EmptyEgressSubscription
import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.LabelsDetector
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.UpstreamValidator
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.CallSelector
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import org.springframework.cloud.sleuth.Tracer
import reactor.core.publisher.Mono
import reactor.core.scheduler.Scheduler
abstract class AbstractChainSpecific : ChainSpecific {
override fun localReaderBuilder(
cachingReader: CachingReader,
methods: CallMethods,
head: Head,
): Mono<JsonRpcReader> {
return Mono.just(LocalReader(methods))
}
override fun makeCachingReaderBuilder(tracer: Tracer): CachingReaderBuilder {
return { _, _, _ -> NoopCachingReader }
}
override fun validator(
chain: Chain,
upstream: Upstream,
options: Options,
config: ChainConfig,
): UpstreamValidator? {
return null
}
override fun labelDetector(chain: Chain, reader: JsonRpcReader): LabelsDetector? {
return null
}
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
return NoIngressSubscription()
}
override fun subscriptionBuilder(headScheduler: Scheduler): (Multistream) -> EgressSubscription {
return { _ -> EmptyEgressSubscription }
}
override fun callSelector(caches: Caches): CallSelector? {
return null
}
}
abstract class AbstractPollChainSpecific : AbstractChainSpecific() {
override fun getLatestBlock(api: JsonRpcReader, upstreamId: String): Mono<BlockContainer> {
return api.read(latestBlockRequest()).map {
parseBlock(it.getResult(), upstreamId)
}
}
abstract fun latestBlockRequest(): JsonRpcRequest
abstract fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer
}

View File

@@ -3,6 +3,7 @@ package io.emeraldpay.dshackle.upstream.generic
import io.emeraldpay.dshackle.BlockchainType.BITCOIN import io.emeraldpay.dshackle.BlockchainType.BITCOIN
import io.emeraldpay.dshackle.BlockchainType.ETHEREUM import io.emeraldpay.dshackle.BlockchainType.ETHEREUM
import io.emeraldpay.dshackle.BlockchainType.POLKADOT import io.emeraldpay.dshackle.BlockchainType.POLKADOT
import io.emeraldpay.dshackle.BlockchainType.SOLANA
import io.emeraldpay.dshackle.BlockchainType.STARKNET import io.emeraldpay.dshackle.BlockchainType.STARKNET
import io.emeraldpay.dshackle.BlockchainType.UNKNOWN import io.emeraldpay.dshackle.BlockchainType.UNKNOWN
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
@@ -25,6 +26,7 @@ import io.emeraldpay.dshackle.upstream.ethereum.EthereumChainSpecific
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.dshackle.upstream.polkadot.PolkadotChainSpecific import io.emeraldpay.dshackle.upstream.polkadot.PolkadotChainSpecific
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.solana.SolanaChainSpecific
import io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific import io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific
import org.apache.commons.collections4.Factory import org.apache.commons.collections4.Factory
import org.springframework.cloud.sleuth.Tracer import org.springframework.cloud.sleuth.Tracer
@@ -36,11 +38,9 @@ typealias LocalReaderBuilder = (CachingReader, CallMethods, Head) -> Mono<JsonRp
typealias CachingReaderBuilder = (Multistream, Caches, Factory<CallMethods>) -> CachingReader typealias CachingReaderBuilder = (Multistream, Caches, Factory<CallMethods>) -> CachingReader
interface ChainSpecific { interface ChainSpecific {
fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer
fun parseHeader(data: ByteArray, upstreamId: String): BlockContainer fun parseHeader(data: ByteArray, upstreamId: String): BlockContainer
fun latestBlockRequest(): JsonRpcRequest fun getLatestBlock(api: JsonRpcReader, upstreamId: String): Mono<BlockContainer>
fun listenNewHeadsRequest(): JsonRpcRequest fun listenNewHeadsRequest(): JsonRpcRequest
@@ -56,8 +56,6 @@ interface ChainSpecific {
fun labelDetector(chain: Chain, reader: JsonRpcReader): LabelsDetector? fun labelDetector(chain: Chain, reader: JsonRpcReader): LabelsDetector?
fun subscriptionTopics(upstream: GenericUpstream): List<String>
fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription
fun callSelector(caches: Caches): CallSelector? fun callSelector(caches: Caches): CallSelector?
@@ -71,6 +69,7 @@ object ChainSpecificRegistry {
ETHEREUM -> EthereumChainSpecific ETHEREUM -> EthereumChainSpecific
STARKNET -> StarknetChainSpecific STARKNET -> StarknetChainSpecific
POLKADOT -> PolkadotChainSpecific POLKADOT -> PolkadotChainSpecific
SOLANA -> SolanaChainSpecific
BITCOIN -> throw IllegalArgumentException("bitcoin should use custom streams implementation") BITCOIN -> throw IllegalArgumentException("bitcoin should use custom streams implementation")
UNKNOWN -> throw IllegalArgumentException("unknown chain") UNKNOWN -> throw IllegalArgumentException("unknown chain")
} }

View File

@@ -0,0 +1,26 @@
package io.emeraldpay.dshackle.upstream.generic
import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Selector.Matcher
import reactor.core.publisher.Flux
import reactor.core.scheduler.Scheduler
class GenericEgressSubscription(
val multistream: Multistream,
val scheduler: Scheduler,
val methods: List<String>,
) : EgressSubscription {
override fun getAvailableTopics(): List<String> {
return methods
}
override fun subscribe(topic: String, params: Any?, matcher: Matcher): Flux<ByteArray> {
val up = multistream.getUpstreams()
.filter { it.isAvailable() }
.shuffled()
.first { matcher.matches(it) } as GenericUpstream
return up.getIngressSubscription().get<ByteArray>(topic, params)?.connect(matcher) ?: Flux.empty()
}
}

View File

@@ -34,12 +34,9 @@ open class GenericHead(
) : Head, AbstractHead(forkChoice, headScheduler, blockValidator, 60_000, upstreamId) { ) : Head, AbstractHead(forkChoice, headScheduler, blockValidator, 60_000, upstreamId) {
fun getLatestBlock(api: JsonRpcReader): Mono<BlockContainer> { fun getLatestBlock(api: JsonRpcReader): Mono<BlockContainer> {
return api.read(chainSpecific.latestBlockRequest()) return chainSpecific.getLatestBlock(api, upstreamId)
.subscribeOn(headScheduler) .subscribeOn(headScheduler)
.timeout(Defaults.timeout, Mono.error(Exception("Block data not received"))) .timeout(Defaults.timeout, Mono.error(Exception("Block data not received")))
.map {
chainSpecific.parseBlock(it.getResult(), upstreamId)
}
.onErrorResume { err -> .onErrorResume { err ->
log.error("Failed to fetch latest block: ${err.message} $upstreamId", err) log.error("Failed to fetch latest block: ${err.message} $upstreamId", err)
Mono.empty() Mono.empty()

View File

@@ -1,4 +1,4 @@
package io.emeraldpay.dshackle.upstream.polkadot package io.emeraldpay.dshackle.upstream.generic
import io.emeraldpay.dshackle.upstream.IngressSubscription import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.SubscriptionConnect import io.emeraldpay.dshackle.upstream.SubscriptionConnect
@@ -8,28 +8,39 @@ import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
import reactor.core.publisher.Mono import reactor.core.publisher.Mono
import java.time.Duration import java.time.Duration
import java.util.concurrent.ConcurrentHashMap
class PolkadotIngressSubscription(val conn: WsSubscriptions) : IngressSubscription { class GenericIngressSubscription(val conn: WsSubscriptions) : IngressSubscription {
override fun getAvailableTopics(): List<String> { override fun getAvailableTopics(): List<String> {
return emptyList() // not used now return emptyList() // not used now
} }
private val holders = ConcurrentHashMap<Pair<String, Any?>, SubscriptionConnect<out Any>>()
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
override fun <T> get(topic: String): SubscriptionConnect<T> { override fun <T> get(topic: String, params: Any?): SubscriptionConnect<T> {
return PolkaConnect(conn, topic) as SubscriptionConnect<T> return holders.computeIfAbsent(topic to params, { key -> GenericSubscriptionConnect(conn, key.first, key.second) }) as SubscriptionConnect<T>
} }
} }
class PolkaConnect( class GenericSubscriptionConnect(
val conn: WsSubscriptions, val conn: WsSubscriptions,
val topic: String, val topic: String,
val params: Any?,
) : GenericPersistentConnect() { ) : GenericPersistentConnect() {
@Suppress("UNCHECKED_CAST") @Suppress("UNCHECKED_CAST")
override fun createConnection(): Flux<Any> { override fun createConnection(): Flux<Any> {
return conn.subscribe(JsonRpcRequest(topic, listOf())) return conn.subscribe(JsonRpcRequest(topic, getParams(params)))
.data .data
.timeout(Duration.ofSeconds(60), Mono.empty()) .timeout(Duration.ofSeconds(60), Mono.empty())
.onErrorResume { Mono.empty() } as Flux<Any> .onErrorResume { Mono.empty() } as Flux<Any>
} }
private fun getParams(params: Any?): List<Any?> {
if (params == null) {
return listOf()
}
return params as List<Any?>
}
} }

View File

@@ -20,7 +20,9 @@ import io.emeraldpay.api.proto.BlockchainOuterClass
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.cache.Caches import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.UpstreamsConfig import io.emeraldpay.dshackle.config.UpstreamsConfig
import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.reader.JsonRpcReader import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.CachingReader import io.emeraldpay.dshackle.upstream.CachingReader
import io.emeraldpay.dshackle.upstream.DistanceExtractor import io.emeraldpay.dshackle.upstream.DistanceExtractor
import io.emeraldpay.dshackle.upstream.DynamicMergedHead import io.emeraldpay.dshackle.upstream.DynamicMergedHead
@@ -35,8 +37,11 @@ import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.Selector.Matcher import io.emeraldpay.dshackle.upstream.Selector.Matcher
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.calls.CallSelector import io.emeraldpay.dshackle.upstream.calls.CallSelector
import io.emeraldpay.dshackle.upstream.ethereum.EnrichedMergedHead
import io.emeraldpay.dshackle.upstream.ethereum.EthereumCachingReader
import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.PriorityForkChoice
import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream import io.emeraldpay.dshackle.upstream.grpc.GrpcUpstream
import io.emeraldpay.etherjar.domain.BlockHash
import org.springframework.util.ConcurrentReferenceHashMap import org.springframework.util.ConcurrentReferenceHashMap
import org.springframework.util.ConcurrentReferenceHashMap.ReferenceType.WEAK import org.springframework.util.ConcurrentReferenceHashMap.ReferenceType.WEAK
import reactor.core.publisher.Flux import reactor.core.publisher.Flux
@@ -146,9 +151,26 @@ open class GenericMultistream(
return head return head
} }
override fun getEnrichedHead(mather: Matcher): Head { override fun getEnrichedHead(mather: Selector.Matcher): Head =
return getHead() filteredHeads.computeIfAbsent(mather.describeInternal().intern()) { _ ->
} upstreams.filter { mather.matches(it) }
.apply {
log.debug("Found $size upstreams matching [${mather.describeInternal()}]")
}.let {
val selected = it.map { source -> source.getHead() }
EnrichedMergedHead(
selected,
getHead(),
headScheduler,
object :
Reader<BlockHash, BlockContainer> {
override fun read(key: BlockHash): Mono<BlockContainer> {
return (cachingReader as EthereumCachingReader).blocksByHashAsCont().read(key).map { res -> res.data }
}
},
)
}
}
override fun getLabels(): Collection<UpstreamsConfig.Labels> { override fun getLabels(): Collection<UpstreamsConfig.Labels> {
return upstreams.flatMap { it.getLabels() } return upstreams.flatMap { it.getLabels() }

View File

@@ -40,7 +40,6 @@ open class GenericUpstream(
connectorFactory: ConnectorFactory, connectorFactory: ConnectorFactory,
validatorBuilder: UpstreamValidatorBuilder, validatorBuilder: UpstreamValidatorBuilder,
labelsDetectorBuilder: LabelsDetectorBuilder, labelsDetectorBuilder: LabelsDetectorBuilder,
private val subscriptionTopics: (GenericUpstream) -> List<String>,
) : DefaultUpstream(id, hash, null, UpstreamAvailability.OK, options, role, targets, node, chainConfig), Lifecycle { ) : DefaultUpstream(id, hash, null, UpstreamAvailability.OK, options, role, targets, node, chainConfig), Lifecycle {
private val validator: UpstreamValidator? = validatorBuilder(chain, this, getOptions(), chainConfig) private val validator: UpstreamValidator? = validatorBuilder(chain, this, getOptions(), chainConfig)
@@ -64,10 +63,6 @@ open class GenericUpstream(
return node?.let { listOf(it.labels) } ?: emptyList() return node?.let { listOf(it.labels) } ?: emptyList()
} }
override fun getSubscriptionTopics(): List<String> {
return subscriptionTopics(this)
}
// outdated, looks like applicable only for bitcoin and our ws_head trick // outdated, looks like applicable only for bitcoin and our ws_head trick
override fun getCapabilities(): Set<Capability> { override fun getCapabilities(): Set<Capability> {
return if (hasLiveSubscriptionHead.get()) { return if (hasLiveSubscriptionHead.get()) {

View File

@@ -10,11 +10,12 @@ import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.IngressSubscription import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.Lifecycle import io.emeraldpay.dshackle.upstream.Lifecycle
import io.emeraldpay.dshackle.upstream.MergedHead import io.emeraldpay.dshackle.upstream.MergedHead
import io.emeraldpay.dshackle.upstream.NoIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.GenericWsHead import io.emeraldpay.dshackle.upstream.ethereum.GenericWsHead
import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessValidator import io.emeraldpay.dshackle.upstream.ethereum.HeadLivenessValidator
import io.emeraldpay.dshackle.upstream.ethereum.NoEthereumIngressSubscription
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPool import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPool
import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPoolFactory import io.emeraldpay.dshackle.upstream.ethereum.WsConnectionPoolFactory
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptionsImpl import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptionsImpl
import io.emeraldpay.dshackle.upstream.forkchoice.AlwaysForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.AlwaysForkChoice
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
@@ -41,10 +42,12 @@ class GenericRpcConnector(
headScheduler: Scheduler, headScheduler: Scheduler,
headLivenessScheduler: Scheduler, headLivenessScheduler: Scheduler,
expectedBlockTime: Duration, expectedBlockTime: Duration,
chainSpecific: ChainSpecific, private val chainSpecific: ChainSpecific,
) : GenericConnector, CachesEnabled { ) : GenericConnector, CachesEnabled {
private val id = upstream.getId() private val id = upstream.getId()
private val pool: WsConnectionPool? private val pool: WsConnectionPool?
private val wsSubs: WsSubscriptions?
private val ingressSubscription: IngressSubscription?
private val head: Head private val head: Head
private val liveness: HeadLivenessValidator private val liveness: HeadLivenessValidator
@@ -58,6 +61,8 @@ class GenericRpcConnector(
init { init {
pool = wsFactory?.create(upstream) pool = wsFactory?.create(upstream)
wsSubs = pool?.let { WsSubscriptionsImpl(it) }
ingressSubscription = wsSubs?.let { chainSpecific.makeIngressSubscription(it) }
head = when (connectorType) { head = when (connectorType) {
RPC_ONLY -> { RPC_ONLY -> {
@@ -83,7 +88,7 @@ class GenericRpcConnector(
AlwaysForkChoice(), AlwaysForkChoice(),
blockValidator, blockValidator,
getIngressReader(), getIngressReader(),
WsSubscriptionsImpl(pool!!), wsSubs!!,
wsConnectionResubscribeScheduler, wsConnectionResubscribeScheduler,
headScheduler, headScheduler,
upstream, upstream,
@@ -108,7 +113,7 @@ class GenericRpcConnector(
AlwaysForkChoice(), AlwaysForkChoice(),
blockValidator, blockValidator,
getIngressReader(), getIngressReader(),
WsSubscriptionsImpl(pool!!), wsSubs!!,
wsConnectionResubscribeScheduler, wsConnectionResubscribeScheduler,
headScheduler, headScheduler,
upstream, upstream,
@@ -152,7 +157,7 @@ class GenericRpcConnector(
} }
override fun getIngressSubscription(): IngressSubscription { override fun getIngressSubscription(): IngressSubscription {
return NoEthereumIngressSubscription.DEFAULT return ingressSubscription ?: NoIngressSubscription()
} }
override fun getHead(): Head { override fun getHead(): Head {

View File

@@ -89,10 +89,6 @@ class BitcoinGrpcUpstream(
block block
} }
override fun getSubscriptionTopics(): List<String> {
return listOf()
}
private val reloadBlock: Function<BlockContainer, Publisher<BlockContainer>> = Function { existingBlock -> private val reloadBlock: Function<BlockContainer, Publisher<BlockContainer>> = Function { existingBlock ->
// head comes without transaction data // head comes without transaction data
// need to download transactions for the block // need to download transactions for the block

View File

@@ -102,9 +102,6 @@ open class GenericGrpcUpstream(
private val defaultReader: JsonRpcReader = client.getReader() private val defaultReader: JsonRpcReader = client.getReader()
// private val ethereumSubscriptions = EthereumDshackleIngressSubscription(chain, remote)
private var subscriptionTopics = listOf<String>()
override fun start() { override fun start() {
} }
@@ -115,10 +112,6 @@ open class GenericGrpcUpstream(
override fun stop() { override fun stop() {
} }
override fun getSubscriptionTopics(): List<String> {
return subscriptionTopics
}
override fun getBuildInfo(): BuildInfo { override fun getBuildInfo(): BuildInfo {
return buildInfo return buildInfo
} }
@@ -130,11 +123,8 @@ open class GenericGrpcUpstream(
val upstreamStatusChanged = (upstreamStatus.update(conf) || (newCapabilities != capabilities)).also { val upstreamStatusChanged = (upstreamStatus.update(conf) || (newCapabilities != capabilities)).also {
capabilities = newCapabilities capabilities = newCapabilities
} }
val subsChanged = (conf.supportedSubscriptionsList != subscriptionTopics).also {
subscriptionTopics = conf.supportedSubscriptionsList
}
conf.status?.let { status -> onStatus(status) } conf.status?.let { status -> onStatus(status) }
return buildInfoChanged || upstreamStatusChanged || subsChanged return buildInfoChanged || upstreamStatusChanged
} }
override fun getQuorumByLabel(): QuorumForLabels { override fun getQuorumByLabel(): QuorumForLabels {

View File

@@ -4,7 +4,6 @@ import com.fasterxml.jackson.annotation.JsonIgnoreProperties
import com.fasterxml.jackson.annotation.JsonProperty import com.fasterxml.jackson.annotation.JsonProperty
import io.emeraldpay.dshackle.Chain import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.ChainsConfig.ChainConfig import io.emeraldpay.dshackle.config.ChainsConfig.ChainConfig
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
@@ -20,11 +19,12 @@ import io.emeraldpay.dshackle.upstream.NoopCachingReader
import io.emeraldpay.dshackle.upstream.Upstream import io.emeraldpay.dshackle.upstream.Upstream
import io.emeraldpay.dshackle.upstream.UpstreamValidator import io.emeraldpay.dshackle.upstream.UpstreamValidator
import io.emeraldpay.dshackle.upstream.calls.CallMethods import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.CallSelector import io.emeraldpay.dshackle.upstream.calls.DefaultPolkadotMethods
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.dshackle.upstream.generic.AbstractPollChainSpecific
import io.emeraldpay.dshackle.upstream.generic.CachingReaderBuilder import io.emeraldpay.dshackle.upstream.generic.CachingReaderBuilder
import io.emeraldpay.dshackle.upstream.generic.ChainSpecific import io.emeraldpay.dshackle.upstream.generic.GenericEgressSubscription
import io.emeraldpay.dshackle.upstream.generic.GenericUpstream import io.emeraldpay.dshackle.upstream.generic.GenericIngressSubscription
import io.emeraldpay.dshackle.upstream.generic.LocalReader import io.emeraldpay.dshackle.upstream.generic.LocalReader
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import org.springframework.cloud.sleuth.Tracer import org.springframework.cloud.sleuth.Tracer
@@ -33,7 +33,7 @@ import reactor.core.scheduler.Scheduler
import java.math.BigInteger import java.math.BigInteger
import java.time.Instant import java.time.Instant
object PolkadotChainSpecific : ChainSpecific { object PolkadotChainSpecific : AbstractPollChainSpecific() {
override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer { override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer {
val response = Global.objectMapper.readValue(data, PolkadotBlockResponse::class.java) val response = Global.objectMapper.readValue(data, PolkadotBlockResponse::class.java)
@@ -79,7 +79,7 @@ object PolkadotChainSpecific : ChainSpecific {
} }
override fun subscriptionBuilder(headScheduler: Scheduler): (Multistream) -> EgressSubscription { override fun subscriptionBuilder(headScheduler: Scheduler): (Multistream) -> EgressSubscription {
return { ms -> PolkadotEgressSubscription(ms, headScheduler) } return { ms -> GenericEgressSubscription(ms, headScheduler, DefaultPolkadotMethods.subs.map { it.first }) }
} }
override fun makeCachingReaderBuilder(tracer: Tracer): CachingReaderBuilder { override fun makeCachingReaderBuilder(tracer: Tracer): CachingReaderBuilder {
@@ -99,16 +99,8 @@ object PolkadotChainSpecific : ChainSpecific {
return null return null
} }
override fun subscriptionTopics(upstream: GenericUpstream): List<String> {
return emptyList()
}
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription { override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
return PolkadotIngressSubscription(ws) return GenericIngressSubscription(ws)
}
override fun callSelector(caches: Caches): CallSelector? {
return null
} }
} }

View File

@@ -1,23 +0,0 @@
package io.emeraldpay.dshackle.upstream.polkadot
import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.Selector.Matcher
import io.emeraldpay.dshackle.upstream.calls.DefaultPolkadotMethods
import io.emeraldpay.dshackle.upstream.generic.GenericUpstream
import reactor.core.publisher.Flux
import reactor.core.scheduler.Scheduler
class PolkadotEgressSubscription(
val upstream: Multistream,
val scheduler: Scheduler,
) : EgressSubscription {
override fun getAvailableTopics(): List<String> {
return DefaultPolkadotMethods.subs.map { it.first }
}
override fun subscribe(topic: String, params: Any?, matcher: Matcher): Flux<ByteArray> {
val up = upstream.getUpstreams().shuffled().first { matcher.matches(it) } as GenericUpstream
return up.getIngressSubscription().get<ByteArray>(topic)?.connect(matcher) ?: Flux.empty()
}
}

View File

@@ -97,7 +97,10 @@ class JsonRpcResponse(
if (str.startsWith("\"") && str.endsWith("\"")) { if (str.startsWith("\"") && str.endsWith("\"")) {
return str.substring(1, str.length - 1) return str.substring(1, str.length - 1)
} }
throw IllegalStateException("Not as JS string") if (str.all { it.isDigit() }) {
return str
}
throw IllegalStateException("Not as JS string - [$str]")
} }
fun requireResult(): Mono<ByteArray> { fun requireResult(): Mono<ByteArray> {

View File

@@ -0,0 +1,127 @@
package io.emeraldpay.dshackle.upstream.solana
import com.fasterxml.jackson.annotation.JsonIgnoreProperties
import com.fasterxml.jackson.annotation.JsonProperty
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.DefaultSolanaMethods
import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.LabelsDetector
import io.emeraldpay.dshackle.upstream.Multistream
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.dshackle.upstream.generic.AbstractChainSpecific
import io.emeraldpay.dshackle.upstream.generic.GenericEgressSubscription
import io.emeraldpay.dshackle.upstream.generic.GenericIngressSubscription
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import reactor.core.publisher.Mono
import reactor.core.scheduler.Scheduler
import java.math.BigInteger
import java.time.Instant
object SolanaChainSpecific : AbstractChainSpecific() {
override fun getLatestBlock(api: JsonRpcReader, upstreamId: String): Mono<BlockContainer> {
return api.read(JsonRpcRequest("getLatestBlockhash", listOf())).flatMap {
val response = Global.objectMapper.readValue(it.getResult(), SolanaLatest::class.java)
api.read(
JsonRpcRequest(
"getBlock",
listOf(
response.context.slot,
mapOf(
"showRewards" to false,
"transactionDetails" to "none",
"maxSupportedTransactionVersion" to 0,
),
),
),
).map {
val raw = it.getResult()
val block = Global.objectMapper.readValue(it.getResult(), SolanaBlock::class.java)
makeBlock(raw, block, upstreamId)
}
}
}
override fun parseHeader(data: ByteArray, upstreamId: String): BlockContainer {
val res = Global.objectMapper.readValue(data, SolanaWrapper::class.java)
return makeBlock(data, res.value.block, upstreamId)
}
private fun makeBlock(raw: ByteArray, block: SolanaBlock, upstreamId: String): BlockContainer {
return BlockContainer(
height = block.height,
hash = BlockId.fromBase64(block.hash),
difficulty = BigInteger.ZERO,
timestamp = Instant.ofEpochMilli(block.timestamp),
full = false,
json = raw,
parsed = block,
transactions = emptyList(),
upstreamId = upstreamId,
parentHash = BlockId.fromBase64(block.parent),
)
}
override fun listenNewHeadsRequest(): JsonRpcRequest {
return JsonRpcRequest(
"blockSubscribe",
listOf(
"all",
mapOf(
"showRewards" to false,
"transactionDetails" to "none",
),
),
)
}
override fun unsubscribeNewHeadsRequest(subId: String): JsonRpcRequest {
return JsonRpcRequest("blockUnsubscribe", listOf(subId))
}
override fun labelDetector(chain: Chain, reader: JsonRpcReader): LabelsDetector? {
return null
}
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
return GenericIngressSubscription(ws)
}
override fun subscriptionBuilder(headScheduler: Scheduler): (Multistream) -> EgressSubscription {
return { ms -> GenericEgressSubscription(ms, headScheduler, DefaultSolanaMethods.subs.map { it.first }) }
}
}
@JsonIgnoreProperties(ignoreUnknown = true)
data class SolanaLatest(
@JsonProperty("context") var context: SolanaContext,
)
@JsonIgnoreProperties(ignoreUnknown = true)
data class SolanaWrapper(
@JsonProperty("context") var context: SolanaContext,
@JsonProperty("value") var value: SolanaResult,
)
@JsonIgnoreProperties(ignoreUnknown = true)
data class SolanaContext(
@JsonProperty("slot") var slot: Long,
)
@JsonIgnoreProperties(ignoreUnknown = true)
data class SolanaResult(
@JsonProperty("block") var block: SolanaBlock,
)
@JsonIgnoreProperties(ignoreUnknown = true)
data class SolanaBlock(
@JsonProperty("blockHeight") var height: Long,
@JsonProperty("blockTime") var timestamp: Long,
@JsonProperty("blockhash") var hash: String,
@JsonProperty("previousBlockhash") var parent: String,
)

View File

@@ -2,40 +2,15 @@ package io.emeraldpay.dshackle.upstream.starknet
import com.fasterxml.jackson.annotation.JsonIgnoreProperties import com.fasterxml.jackson.annotation.JsonIgnoreProperties
import com.fasterxml.jackson.annotation.JsonProperty import com.fasterxml.jackson.annotation.JsonProperty
import io.emeraldpay.dshackle.Chain
import io.emeraldpay.dshackle.Global import io.emeraldpay.dshackle.Global
import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.config.ChainsConfig.ChainConfig
import io.emeraldpay.dshackle.data.BlockContainer import io.emeraldpay.dshackle.data.BlockContainer
import io.emeraldpay.dshackle.data.BlockId import io.emeraldpay.dshackle.data.BlockId
import io.emeraldpay.dshackle.foundation.ChainOptions.Options import io.emeraldpay.dshackle.upstream.generic.AbstractPollChainSpecific
import io.emeraldpay.dshackle.reader.JsonRpcReader
import io.emeraldpay.dshackle.upstream.CachingReader
import io.emeraldpay.dshackle.upstream.EgressSubscription
import io.emeraldpay.dshackle.upstream.EmptyEgressSubscription
import io.emeraldpay.dshackle.upstream.Head
import io.emeraldpay.dshackle.upstream.IngressSubscription
import io.emeraldpay.dshackle.upstream.LabelsDetector
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.UpstreamValidator
import io.emeraldpay.dshackle.upstream.calls.CallMethods
import io.emeraldpay.dshackle.upstream.calls.CallSelector
import io.emeraldpay.dshackle.upstream.ethereum.WsSubscriptions
import io.emeraldpay.dshackle.upstream.generic.CachingReaderBuilder
import io.emeraldpay.dshackle.upstream.generic.ChainSpecific
import io.emeraldpay.dshackle.upstream.generic.GenericUpstream
import io.emeraldpay.dshackle.upstream.generic.LocalReader
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import org.springframework.cloud.sleuth.Tracer
import reactor.core.publisher.Mono
import reactor.core.scheduler.Scheduler
import java.math.BigInteger import java.math.BigInteger
import java.time.Instant import java.time.Instant
object StarknetChainSpecific : ChainSpecific { object StarknetChainSpecific : AbstractPollChainSpecific() {
override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer { override fun parseBlock(data: ByteArray, upstreamId: String): BlockContainer {
val block = Global.objectMapper.readValue(data, StarknetBlock::class.java) val block = Global.objectMapper.readValue(data, StarknetBlock::class.java)
@@ -57,9 +32,6 @@ object StarknetChainSpecific : ChainSpecific {
throw NotImplementedError() throw NotImplementedError()
} }
override fun latestBlockRequest(): JsonRpcRequest =
JsonRpcRequest("starknet_getBlockWithTxHashes", listOf("latest"))
override fun listenNewHeadsRequest(): JsonRpcRequest { override fun listenNewHeadsRequest(): JsonRpcRequest {
throw NotImplementedError() throw NotImplementedError()
} }
@@ -68,46 +40,8 @@ object StarknetChainSpecific : ChainSpecific {
throw NotImplementedError() throw NotImplementedError()
} }
override fun localReaderBuilder( override fun latestBlockRequest(): JsonRpcRequest =
cachingReader: CachingReader, JsonRpcRequest("starknet_getBlockWithTxHashes", listOf("latest"))
methods: CallMethods,
head: Head,
): Mono<JsonRpcReader> {
return Mono.just(LocalReader(methods))
}
override fun subscriptionBuilder(headScheduler: Scheduler): (Multistream) -> EgressSubscription {
return { _ -> EmptyEgressSubscription }
}
override fun makeCachingReaderBuilder(tracer: Tracer): CachingReaderBuilder {
return { _, _, _ -> NoopCachingReader }
}
override fun validator(
chain: Chain,
upstream: Upstream,
options: Options,
config: ChainConfig,
): UpstreamValidator? {
return null
}
override fun labelDetector(chain: Chain, reader: JsonRpcReader): LabelsDetector? {
return null
}
override fun subscriptionTopics(upstream: GenericUpstream): List<String> {
return emptyList()
}
override fun makeIngressSubscription(ws: WsSubscriptions): IngressSubscription {
return NoIngressSubscription()
}
override fun callSelector(caches: Caches): CallSelector? {
return null
}
} }
@JsonIgnoreProperties(ignoreUnknown = true) @JsonIgnoreProperties(ignoreUnknown = true)

View File

@@ -41,11 +41,11 @@ class NativeSubscribeSpec extends Specification {
def subscribe = Mock(EthereumEgressSubscription) { def subscribe = Mock(EthereumEgressSubscription) {
1 * it.subscribe("newHeads", null, _ as Selector.AnyLabelMatcher) >> Flux.just("{}") 1 * it.subscribe("newHeads", null, _ as Selector.AnyLabelMatcher) >> Flux.just("{}")
1 * it.getAvailableTopics() >> ["newHeads"]
} }
def up = Mock(GenericMultistream) { def up = Mock(GenericMultistream) {
1 * it.tryProxySubscribe(_ as Selector.AnyLabelMatcher, call) >> null 1 * it.tryProxySubscribe(_ as Selector.AnyLabelMatcher, call) >> null
1 * it.getEgressSubscription() >> subscribe 2 * it.getEgressSubscription() >> subscribe
1 * it.getSubscriptionTopics() >> ["newHeads"]
} }
def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM__MAINNET, up), signer) def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM__MAINNET, up), signer)
@@ -81,11 +81,11 @@ class NativeSubscribeSpec extends Specification {
println("ok: $ok") println("ok: $ok")
ok ok
}, _ as Selector.AnyLabelMatcher) >> Flux.just("{}") }, _ as Selector.AnyLabelMatcher) >> Flux.just("{}")
1 * it.getAvailableTopics() >> ["logs"]
} }
def up = Mock(GenericMultistream) { def up = Mock(GenericMultistream) {
1 * it.tryProxySubscribe(_ as Selector.AnyLabelMatcher, call) >> null 1 * it.tryProxySubscribe(_ as Selector.AnyLabelMatcher, call) >> null
1 * it.getEgressSubscription() >> subscribe 2 * it.getEgressSubscription() >> subscribe
1 * it.getSubscriptionTopics() >> ["logs"]
} }
def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM__MAINNET, up), signer) def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM__MAINNET, up), signer)
@@ -106,10 +106,14 @@ class NativeSubscribeSpec extends Specification {
.setChainValue(Chain.ETHEREUM__MAINNET.id) .setChainValue(Chain.ETHEREUM__MAINNET.id)
.setMethod("newHeads") .setMethod("newHeads")
.build() .build()
def subscribe = Mock(EthereumEgressSubscription) {
1 * it.getAvailableTopics() >> ["newHeads"]
}
def up = Mock(GenericMultistream) { def up = Mock(GenericMultistream) {
1 * it.tryProxySubscribe(_ as Selector.AnyLabelMatcher, call) >> Flux.just("{}") 1 * it.tryProxySubscribe(_ as Selector.AnyLabelMatcher, call) >> Flux.just("{}")
0 * it.getEgressSubscription() 1 * it.getEgressSubscription() >> subscribe
1 * it.getSubscriptionTopics() >> ["newHeads"]
} }
def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM__MAINNET, up), signer) def nativeSubscribe = new NativeSubscribe(new MultistreamHolderMock(Chain.ETHEREUM__MAINNET, up), signer)

View File

@@ -74,7 +74,6 @@ class GenericUpstreamMock extends GenericUpstream {
new ConnectorFactoryMock(api, new EthereumHeadMock()), new ConnectorFactoryMock(api, new EthereumHeadMock()),
io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&validator, io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&validator,
io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&labelDetector, io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&labelDetector,
io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&subscriptionTopics,
) )
this.ethereumHeadMock = this.getHead() as EthereumHeadMock this.ethereumHeadMock = this.getHead() as EthereumHeadMock
setLag(0) setLag(0)

View File

@@ -79,7 +79,6 @@ class FilteredApisSpec extends Specification {
connectorFactory, connectorFactory,
cs.&validator, cs.&validator,
cs.&labelDetector, cs.&labelDetector,
cs.&subscriptionTopics,
) )
} }
def matcher = new Selector.LabelMatcher("test", ["foo"]) def matcher = new Selector.LabelMatcher("test", ["foo"])

View File

@@ -0,0 +1,35 @@
package io.emeraldpay.dshackle.upstream.solana
import io.emeraldpay.dshackle.data.BlockId
import org.assertj.core.api.Assertions
import org.junit.jupiter.api.Test
val example = """{
"context": {
"slot": 112301554
},
"value": {
"slot": 112301554,
"block": {
"previousBlockhash": "GJp125YAN4ufCSUvZJVdCyWQJ7RPWMmwxoyUQySydZA",
"blockhash": "6ojMHjctdqfB55JDpEpqfHnP96fiaHEcvzEQ2NNcxzHP",
"parentSlot": 112301553,
"blockTime": 1639926816,
"blockHeight": 101210751
},
"err": null
}
}
""".trimIndent()
class SolanaChainSpecificTest {
@Test
fun parseBlock() {
val result = SolanaChainSpecific.parseHeader(example.toByteArray(), "1")
Assertions.assertThat(result.height).isEqualTo(101210751)
Assertions.assertThat(result.hash).isEqualTo(BlockId.fromBase64("6ojMHjctdqfB55JDpEpqfHnP96fiaHEcvzEQ2NNcxzHP"))
Assertions.assertThat(result.upstreamId).isEqualTo("1")
Assertions.assertThat(result.parentHash).isEqualTo(BlockId.fromBase64("GJp125YAN4ufCSUvZJVdCyWQJ7RPWMmwxoyUQySydZA"))
}
}