Lower bounds fixes (#777)
This commit is contained in:
@@ -20,6 +20,8 @@ import com.fasterxml.jackson.annotation.JsonSubTypes
|
||||
import com.fasterxml.jackson.annotation.JsonTypeInfo
|
||||
import io.emeraldpay.dshackle.foundation.ChainOptions
|
||||
import io.emeraldpay.dshackle.upstream.generic.connectors.GenericConnectorFactory.ConnectorMode
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.ManualLowerBoundType
|
||||
import java.net.URI
|
||||
import java.util.Arrays
|
||||
import java.util.Locale
|
||||
@@ -42,6 +44,7 @@ data class UpstreamsConfig(
|
||||
var methodGroups: MethodGroups? = null,
|
||||
var role: UpstreamRole = UpstreamRole.PRIMARY,
|
||||
var customHeaders: Map<String, String> = emptyMap(),
|
||||
var additionalSettings: AdditionalSettings? = null,
|
||||
) {
|
||||
|
||||
@Suppress("UNCHECKED_CAST")
|
||||
@@ -188,6 +191,15 @@ data class UpstreamsConfig(
|
||||
}
|
||||
}
|
||||
|
||||
data class AdditionalSettings(
|
||||
val manualLowerBounds: Map<LowerBoundType, ManualBoundSetting>,
|
||||
)
|
||||
|
||||
data class ManualBoundSetting(
|
||||
val type: ManualLowerBoundType,
|
||||
val value: Long,
|
||||
)
|
||||
|
||||
data class Methods(
|
||||
val enabled: Set<Method>,
|
||||
val disabled: Set<Method>,
|
||||
|
||||
@@ -17,12 +17,16 @@
|
||||
package io.emeraldpay.dshackle.config
|
||||
|
||||
import io.emeraldpay.dshackle.FileResolver
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig.ManualBoundSetting
|
||||
import io.emeraldpay.dshackle.foundation.ChainOptions
|
||||
import io.emeraldpay.dshackle.foundation.ChainOptionsReader
|
||||
import io.emeraldpay.dshackle.foundation.YamlConfigReader
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.ManualLowerBoundType
|
||||
import org.apache.commons.lang3.StringUtils
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.yaml.snakeyaml.nodes.MappingNode
|
||||
import org.yaml.snakeyaml.nodes.Tag
|
||||
import java.io.InputStream
|
||||
import java.net.URI
|
||||
import java.util.Locale
|
||||
@@ -287,6 +291,7 @@ class UpstreamsConfigReader(
|
||||
upstream.options = optionsReader.read(upNode)
|
||||
upstream.methods = tryReadMethods(upNode)
|
||||
upstream.methodGroups = tryReadMethodGroups(upNode)
|
||||
upstream.additionalSettings = readAdditionalSettings(upNode)
|
||||
getValueAsBool(upNode, "enabled")?.let {
|
||||
upstream.isEnabled = it
|
||||
}
|
||||
@@ -311,6 +316,38 @@ class UpstreamsConfigReader(
|
||||
}
|
||||
}
|
||||
|
||||
private fun readAdditionalSettings(upNode: MappingNode): UpstreamsConfig.AdditionalSettings? {
|
||||
return getMapping(upNode, "additional-settings")
|
||||
?.let { settings ->
|
||||
val manualLowerBounds = getMapping(settings, "manual-lower-bounds")
|
||||
?.let { bounds ->
|
||||
bounds.value.asSequence()
|
||||
.filter {
|
||||
StringUtils.isNotBlank(it.keyNode.valueAsString()) && it.valueNode.tag == Tag.MAP
|
||||
}
|
||||
.map {
|
||||
val setting = readManualBoundSetting(it.valueNode as MappingNode)
|
||||
|
||||
if (setting != null) {
|
||||
LowerBoundType.byName(it.keyNode.valueAsString()!!) to setting
|
||||
} else {
|
||||
null
|
||||
}
|
||||
}
|
||||
.filterNotNull()
|
||||
.associate { it.first to it.second }
|
||||
} ?: emptyMap()
|
||||
UpstreamsConfig.AdditionalSettings(manualLowerBounds)
|
||||
}
|
||||
}
|
||||
|
||||
private fun readManualBoundSetting(node: MappingNode): ManualBoundSetting? {
|
||||
val type = getValueAsString(node, "type")
|
||||
?.let { ManualLowerBoundType.byName(it) } ?: return null
|
||||
val value = getValueAsLong(node, "value") ?: return null
|
||||
return ManualBoundSetting(type, value)
|
||||
}
|
||||
|
||||
private fun readUpstreamGrpc(
|
||||
upNode: MappingNode,
|
||||
) {
|
||||
|
||||
@@ -172,6 +172,10 @@ abstract class DefaultUpstream(
|
||||
return 0
|
||||
}
|
||||
|
||||
override fun getAdditionalSettings(): UpstreamsConfig.AdditionalSettings? {
|
||||
return null
|
||||
}
|
||||
|
||||
protected fun sendUpstreamStateEvent(eventType: UpstreamChangeEvent.ChangeType) {
|
||||
stateEventStream.emitNext(
|
||||
UpstreamChangeEvent(chain, this, eventType),
|
||||
|
||||
@@ -261,6 +261,10 @@ abstract class Multistream(
|
||||
return 0
|
||||
}
|
||||
|
||||
override fun getAdditionalSettings(): UpstreamsConfig.AdditionalSettings? {
|
||||
return null
|
||||
}
|
||||
|
||||
override fun getStatus(): UpstreamAvailability {
|
||||
return state.getStatus()
|
||||
}
|
||||
|
||||
@@ -56,6 +56,7 @@ interface Upstream : Lifecycle {
|
||||
fun getUpstreamSettingsData(): UpstreamSettingsData?
|
||||
fun updateLowerBound(lowerBound: Long, type: LowerBoundType)
|
||||
fun predictLowerBound(type: LowerBoundType, timeOffsetSeconds: Long): Long
|
||||
fun getAdditionalSettings(): UpstreamsConfig.AdditionalSettings?
|
||||
|
||||
fun getChain(): Chain
|
||||
|
||||
|
||||
@@ -21,7 +21,9 @@ object EthereumStateLowerBoundErrorHandler : EthereumLowerBoundErrorHandler() {
|
||||
private val applicableMethods = firstTagIndexMethods + secondTagIndexMethods
|
||||
|
||||
override fun canHandle(request: ChainRequest, errorMessage: String?): Boolean {
|
||||
return stateErrors.any { errorMessage?.contains(it) ?: false } && applicableMethods.contains(request.method)
|
||||
return !(errorMessage?.contains("execution reverted") ?: false) &&
|
||||
stateErrors.any { errorMessage?.contains(it) ?: false } &&
|
||||
applicableMethods.contains(request.method)
|
||||
}
|
||||
|
||||
override fun tagIndex(method: String): Int {
|
||||
|
||||
@@ -29,6 +29,7 @@ class EthereumLowerBoundBlockDetector(
|
||||
"Unexpected error", // hyperliquid
|
||||
"invalid block height", // hyperliquid
|
||||
"pruned history unavailable", // xlayer
|
||||
"no transactions snapshot file for",
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -64,6 +64,7 @@ open class GenericUpstream(
|
||||
lowerBoundServiceBuilder: LowerBoundServiceBuilder,
|
||||
finalizationDetectorBuilder: FinalizationDetectorBuilder,
|
||||
versionRules: Supplier<CompatibleVersionsRules?>,
|
||||
private val additionalSettings: UpstreamsConfig.AdditionalSettings?,
|
||||
) : DefaultUpstream(id, hash, null, UpstreamAvailability.OK, options, role, targets, node, chainConfig, chain),
|
||||
Lifecycle {
|
||||
constructor(
|
||||
@@ -96,6 +97,7 @@ open class GenericUpstream(
|
||||
lowerBoundServiceBuilder,
|
||||
finalizationDetectorBuilder,
|
||||
versionRules,
|
||||
config.additionalSettings,
|
||||
) {
|
||||
rpcMethodsDetector = upstreamRpcMethodsDetectorBuilder(this, config)
|
||||
detectRpcMethods(config, buildMethods)
|
||||
@@ -443,6 +445,10 @@ open class GenericUpstream(
|
||||
return lowerBoundService.predictLowerBound(type, timeOffsetSeconds)
|
||||
}
|
||||
|
||||
override fun getAdditionalSettings(): UpstreamsConfig.AdditionalSettings? {
|
||||
return additionalSettings
|
||||
}
|
||||
|
||||
fun isValid(): Boolean = isUpstreamValid.get()
|
||||
|
||||
companion object {
|
||||
|
||||
@@ -20,7 +20,23 @@ data class LowerBoundData(
|
||||
}
|
||||
|
||||
enum class LowerBoundType {
|
||||
UNKNOWN, STATE, SLOT, BLOCK, TX, LOGS, TRACE, PROOF, BLOB, EPOCH, RECEIPTS
|
||||
UNKNOWN, STATE, SLOT, BLOCK, TX, LOGS, TRACE, PROOF, BLOB, EPOCH, RECEIPTS;
|
||||
|
||||
companion object {
|
||||
fun byName(name: String): LowerBoundType {
|
||||
return entries.firstOrNull { it.name.equals(name, ignoreCase = true) } ?: UNKNOWN
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
enum class ManualLowerBoundType {
|
||||
HEAD, FIXED;
|
||||
|
||||
companion object {
|
||||
fun byName(name: String): ManualLowerBoundType? {
|
||||
return entries.firstOrNull { it.name.equals(name, ignoreCase = true) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun BlockchainOuterClass.LowerBoundType.fromProtoType(): LowerBoundType {
|
||||
|
||||
@@ -18,7 +18,7 @@ abstract class LowerBoundDetector(
|
||||
protected val lowerBounds = LowerBounds(chain)
|
||||
private val lowerBoundSink = Sinks.many().multicast().directBestEffort<LowerBoundData>()
|
||||
|
||||
fun detectLowerBound(): Flux<LowerBoundData> {
|
||||
fun detectLowerBound(manualBoundsService: ManualLowerBoundService): Flux<LowerBoundData> {
|
||||
val notProcessing = AtomicBoolean(true)
|
||||
|
||||
return Flux.merge(
|
||||
@@ -29,11 +29,21 @@ abstract class LowerBoundDetector(
|
||||
)
|
||||
.filter { notProcessing.get() }
|
||||
.flatMap {
|
||||
if (types() == manualBoundsService.manualBoundTypes()) {
|
||||
Flux.fromIterable(types())
|
||||
.mapNotNull { manualBoundsService.manualLowerBound(it) }
|
||||
} else {
|
||||
notProcessing.set(false)
|
||||
Flux.merge(
|
||||
internalDetectLowerBound()
|
||||
.filter { !manualBoundsService.hasManualBound(it.type) }
|
||||
.onErrorResume { Mono.just(LowerBoundData.default()) }
|
||||
.switchIfEmpty(Flux.just(LowerBoundData.default()))
|
||||
.doFinally { notProcessing.set(true) }
|
||||
.doFinally { notProcessing.set(true) },
|
||||
Flux.fromIterable(types())
|
||||
.mapNotNull { manualBoundsService.manualLowerBound(it) },
|
||||
)
|
||||
}
|
||||
},
|
||||
)
|
||||
.filter {
|
||||
|
||||
@@ -16,10 +16,14 @@ abstract class LowerBoundService(
|
||||
|
||||
private val lowerBounds = ConcurrentHashMap<LowerBoundType, LowerBoundData>()
|
||||
private val detectors: List<LowerBoundDetector> by lazy { detectors() }
|
||||
private val manualBoundsService = BaseManualLowerBoundService(
|
||||
upstream,
|
||||
upstream.getAdditionalSettings()?.manualLowerBounds ?: emptyMap(),
|
||||
)
|
||||
|
||||
fun detectLowerBounds(): Flux<LowerBoundData> {
|
||||
return Flux.merge(
|
||||
detectors.map { it.detectLowerBound() },
|
||||
detectors.map { it.detectLowerBound(manualBoundsService) },
|
||||
)
|
||||
.doOnNext {
|
||||
log.info("Lower bound of type ${it.type} is ${it.lowerBound} for upstream ${upstream.getId()} of chain $chain")
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
package io.emeraldpay.dshackle.upstream.lowerbound
|
||||
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
|
||||
interface ManualLowerBoundService {
|
||||
fun manualLowerBound(type: LowerBoundType): LowerBoundData?
|
||||
fun hasManualBound(type: LowerBoundType): Boolean
|
||||
fun manualBoundTypes(): Set<LowerBoundType>
|
||||
}
|
||||
|
||||
class NoopManualLowerBoundService : ManualLowerBoundService {
|
||||
override fun manualLowerBound(type: LowerBoundType): LowerBoundData? {
|
||||
return null
|
||||
}
|
||||
|
||||
override fun hasManualBound(type: LowerBoundType): Boolean {
|
||||
return false
|
||||
}
|
||||
|
||||
override fun manualBoundTypes(): Set<LowerBoundType> {
|
||||
return emptySet()
|
||||
}
|
||||
}
|
||||
|
||||
class BaseManualLowerBoundService(
|
||||
private val upstream: Upstream,
|
||||
private val manualLowerBoundSettings: Map<LowerBoundType, UpstreamsConfig.ManualBoundSetting>,
|
||||
) : ManualLowerBoundService {
|
||||
private val providers = manualLowerBoundSettings.mapValues { provider(it.value) }
|
||||
|
||||
override fun manualLowerBound(type: LowerBoundType): LowerBoundData? {
|
||||
val provider = providers[type] ?: return null
|
||||
if (!provider.canHandle()) {
|
||||
return null
|
||||
}
|
||||
return LowerBoundData(provider.getManualLowerBound(), type)
|
||||
}
|
||||
|
||||
override fun hasManualBound(type: LowerBoundType): Boolean {
|
||||
return manualLowerBoundSettings.contains(type)
|
||||
}
|
||||
|
||||
override fun manualBoundTypes(): Set<LowerBoundType> {
|
||||
return manualLowerBoundSettings.keys
|
||||
}
|
||||
|
||||
private fun provider(setting: UpstreamsConfig.ManualBoundSetting): ManualLowerBoundProvider {
|
||||
return when (setting.type) {
|
||||
ManualLowerBoundType.HEAD -> HeadManualLowerBoundProvider(upstream, setting)
|
||||
ManualLowerBoundType.FIXED -> FixedManualLowerBoundProvider(setting)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
interface ManualLowerBoundProvider {
|
||||
fun getManualLowerBound(): Long
|
||||
fun canHandle(): Boolean
|
||||
}
|
||||
|
||||
class HeadManualLowerBoundProvider(
|
||||
private val upstream: Upstream,
|
||||
private val setting: UpstreamsConfig.ManualBoundSetting,
|
||||
) : ManualLowerBoundProvider {
|
||||
override fun getManualLowerBound(): Long {
|
||||
return upstream.getHead().getCurrentHeight()!! + setting.value
|
||||
}
|
||||
|
||||
override fun canHandle(): Boolean {
|
||||
return setting.type == ManualLowerBoundType.HEAD &&
|
||||
setting.value < 0 &&
|
||||
upstream.getHead().getCurrentHeight() != null &&
|
||||
upstream.getHead().getCurrentHeight()!! + setting.value >= 0
|
||||
}
|
||||
}
|
||||
|
||||
class FixedManualLowerBoundProvider(
|
||||
private val setting: UpstreamsConfig.ManualBoundSetting,
|
||||
) : ManualLowerBoundProvider {
|
||||
override fun getManualLowerBound(): Long {
|
||||
return setting.value
|
||||
}
|
||||
|
||||
override fun canHandle(): Boolean {
|
||||
return setting.type == ManualLowerBoundType.FIXED && setting.value > 0
|
||||
}
|
||||
}
|
||||
@@ -79,6 +79,7 @@ class GenericUpstreamMock extends GenericUpstream {
|
||||
io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&lowerBoundService,
|
||||
io.emeraldpay.dshackle.upstream.starknet.StarknetChainSpecific.INSTANCE.&finalizationDetectorBuilder,
|
||||
[get: { null }] as java.util.function.Supplier,
|
||||
null,
|
||||
)
|
||||
this.ethereumHeadMock = this.getHead() as EthereumHeadMock
|
||||
setLag(0)
|
||||
|
||||
@@ -80,6 +80,7 @@ class FilteredApisSpec extends Specification {
|
||||
cs.&lowerBoundService,
|
||||
cs.&finalizationDetectorBuilder,
|
||||
[get: { null }] as java.util.function.Supplier,
|
||||
null,
|
||||
)
|
||||
}
|
||||
def matcher = new Selector.LabelMatcher("test", ["foo"])
|
||||
|
||||
@@ -2,6 +2,8 @@ package io.emeraldpay.dshackle.config
|
||||
|
||||
import io.emeraldpay.dshackle.FileResolver
|
||||
import io.emeraldpay.dshackle.foundation.ChainOptionsReader
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.ManualLowerBoundType
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
import org.junit.jupiter.api.Assertions.assertNotNull
|
||||
import org.junit.jupiter.api.Assertions.assertTrue
|
||||
@@ -10,6 +12,62 @@ import java.io.File
|
||||
|
||||
class UpstreamsConfigReaderTest {
|
||||
|
||||
@Test
|
||||
fun `should parse additional settings with lower bounds`() {
|
||||
val yaml = """
|
||||
version: v1
|
||||
upstreams:
|
||||
- id: test-upstream
|
||||
chain: ethereum
|
||||
additional-settings:
|
||||
manual-lower-bounds:
|
||||
state:
|
||||
type: head
|
||||
value: -5000000
|
||||
block:
|
||||
type: fixed
|
||||
value: 1
|
||||
tx:
|
||||
type: fixed
|
||||
value: 123
|
||||
receipts:
|
||||
type: fixed
|
||||
value: 15661
|
||||
slot:
|
||||
type: head
|
||||
value: -34566
|
||||
connection:
|
||||
ethereum:
|
||||
rpc:
|
||||
url: "http://localhost:8545"
|
||||
""".trimIndent()
|
||||
|
||||
val reader = UpstreamsConfigReader(
|
||||
FileResolver(File(".")),
|
||||
ChainOptionsReader(),
|
||||
)
|
||||
|
||||
val config = reader.readInternal(yaml.byteInputStream())
|
||||
|
||||
assertNotNull(config)
|
||||
assertEquals(1, config.upstreams.size)
|
||||
|
||||
val upstream = config.upstreams[0]
|
||||
assertEquals("test-upstream", upstream.id)
|
||||
assertEquals(
|
||||
UpstreamsConfig.AdditionalSettings(
|
||||
mapOf(
|
||||
LowerBoundType.SLOT to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, -34566),
|
||||
LowerBoundType.RECEIPTS to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 15661),
|
||||
LowerBoundType.TX to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 123),
|
||||
LowerBoundType.BLOCK to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 1),
|
||||
LowerBoundType.STATE to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, -5000000),
|
||||
),
|
||||
),
|
||||
upstream.additionalSettings,
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `should parse customHeaders from YAML`() {
|
||||
val yaml = """
|
||||
|
||||
@@ -42,6 +42,20 @@ class EthereumStateLowerBoundErrorHandlerTest {
|
||||
verify(upstream, never()).updateLowerBound(anyLong(), any())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `no update lower bound if execution reverted`() {
|
||||
val upstream = mock<Upstream>()
|
||||
val request = ChainRequest(
|
||||
"eth_call",
|
||||
ListParams("0x343", "0xCB5A0A8"),
|
||||
)
|
||||
val handler = EthereumStateLowerBoundErrorHandler
|
||||
|
||||
handler.handle(upstream, request, "execution reverted: Fallback not supported")
|
||||
|
||||
verify(upstream, never()).updateLowerBound(anyLong(), any())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `no update lower bound if there is a non-state method`() {
|
||||
val upstream = mock<Upstream>()
|
||||
|
||||
@@ -3,6 +3,7 @@ package io.emeraldpay.dshackle.upstream.kadena
|
||||
import io.emeraldpay.dshackle.Chain
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundData
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.NoopManualLowerBoundService
|
||||
import org.junit.jupiter.api.Test
|
||||
import reactor.test.StepVerifier
|
||||
import java.time.Duration
|
||||
@@ -13,7 +14,7 @@ class KadenaLowerBoundStateDetectorTest {
|
||||
fun `kadena lower block is 1`() {
|
||||
val detector = KadenaLowerBoundStateDetector(Chain.KADENA__MAINNET)
|
||||
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound() }
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound(NoopManualLowerBoundService()) }
|
||||
.expectSubscription()
|
||||
.expectNoEvent(Duration.ofSeconds(15))
|
||||
.expectNext(LowerBoundData(1, LowerBoundType.STATE))
|
||||
|
||||
@@ -0,0 +1,176 @@
|
||||
package io.emeraldpay.dshackle.upstream.lowerbound
|
||||
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import io.emeraldpay.dshackle.upstream.Head
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import org.assertj.core.api.Assertions.assertThat
|
||||
import org.junit.jupiter.api.Test
|
||||
import org.junit.jupiter.params.ParameterizedTest
|
||||
import org.junit.jupiter.params.provider.Arguments
|
||||
import org.junit.jupiter.params.provider.MethodSource
|
||||
import org.mockito.kotlin.doReturn
|
||||
import org.mockito.kotlin.mock
|
||||
|
||||
class BaseManualLowerBoundServiceTest {
|
||||
|
||||
@Test
|
||||
fun `no manual settings then return nothing`() {
|
||||
val upstream = mock<Upstream>()
|
||||
val service = BaseManualLowerBoundService(upstream, emptyMap())
|
||||
|
||||
assertThat(service.manualBoundTypes()).isEmpty()
|
||||
LowerBoundType.entries
|
||||
.forEach {
|
||||
assertThat(service.manualLowerBound(it)).isNull()
|
||||
assertThat(service.hasManualBound(it)).isFalse
|
||||
}
|
||||
}
|
||||
|
||||
@ParameterizedTest
|
||||
@MethodSource("headProviderData")
|
||||
fun `test head provider - canHandle`(
|
||||
upstream: Upstream,
|
||||
setting: UpstreamsConfig.ManualBoundSetting,
|
||||
expected: Boolean,
|
||||
) {
|
||||
val provider = HeadManualLowerBoundProvider(upstream, setting)
|
||||
|
||||
assertThat(provider.canHandle()).isEqualTo(expected)
|
||||
}
|
||||
|
||||
@ParameterizedTest
|
||||
@MethodSource("fixedProviderData")
|
||||
fun `test fixed provider - canHandle`(
|
||||
setting: UpstreamsConfig.ManualBoundSetting,
|
||||
expected: Boolean,
|
||||
) {
|
||||
val provider = FixedManualLowerBoundProvider(setting)
|
||||
|
||||
assertThat(provider.canHandle()).isEqualTo(expected)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `get value from head provider`() {
|
||||
val upstream = mock<Upstream> {
|
||||
val head = mock<Head> {
|
||||
on { getCurrentHeight() } doReturn 100
|
||||
}
|
||||
on { getHead() } doReturn head
|
||||
}
|
||||
val provider = HeadManualLowerBoundProvider(upstream, UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, -45))
|
||||
|
||||
assertThat(provider.getManualLowerBound()).isEqualTo(55)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `get value from fixed provider`() {
|
||||
val provider = FixedManualLowerBoundProvider(UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 200))
|
||||
|
||||
assertThat(provider.getManualLowerBound()).isEqualTo(200)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `no supported manual type then null otherwise get a value`() {
|
||||
val upstream = mock<Upstream>()
|
||||
val settings = mapOf(
|
||||
LowerBoundType.STATE to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 100),
|
||||
)
|
||||
val service = BaseManualLowerBoundService(upstream, settings)
|
||||
|
||||
assertThat(service.manualLowerBound(LowerBoundType.STATE))
|
||||
.usingRecursiveComparison()
|
||||
.ignoringFields("timestamp")
|
||||
.isEqualTo(LowerBoundData(100, LowerBoundType.STATE))
|
||||
assertThat(service.manualBoundTypes()).containsExactly(LowerBoundType.STATE)
|
||||
assertThat(service.hasManualBound(LowerBoundType.STATE)).isTrue()
|
||||
LowerBoundType.entries
|
||||
.filter { it != LowerBoundType.STATE }
|
||||
.forEach {
|
||||
assertThat(service.manualLowerBound(it)).isNull()
|
||||
assertThat(service.hasManualBound(it)).isFalse
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `can not be handled then null`() {
|
||||
val upstream = mock<Upstream>()
|
||||
val settings = mapOf(
|
||||
LowerBoundType.STATE to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, -100),
|
||||
)
|
||||
val service = BaseManualLowerBoundService(upstream, settings)
|
||||
|
||||
assertThat(service.manualLowerBound(LowerBoundType.STATE)).isNull()
|
||||
assertThat(service.manualBoundTypes()).containsExactly(LowerBoundType.STATE)
|
||||
assertThat(service.hasManualBound(LowerBoundType.STATE)).isTrue()
|
||||
LowerBoundType.entries
|
||||
.filter { it != LowerBoundType.STATE }
|
||||
.forEach {
|
||||
assertThat(service.manualLowerBound(it)).isNull()
|
||||
assertThat(service.hasManualBound(it)).isFalse
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
@JvmStatic
|
||||
fun fixedProviderData(): List<Arguments> =
|
||||
listOf(
|
||||
Arguments.of(
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, 1),
|
||||
false,
|
||||
),
|
||||
Arguments.of(
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, -121),
|
||||
false,
|
||||
),
|
||||
Arguments.of(
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 55),
|
||||
true,
|
||||
),
|
||||
)
|
||||
|
||||
@JvmStatic
|
||||
fun headProviderData(): List<Arguments> =
|
||||
listOf(
|
||||
Arguments.of(
|
||||
mock<Upstream>(),
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 1),
|
||||
false,
|
||||
),
|
||||
Arguments.of(
|
||||
mock<Upstream>(),
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, 1241),
|
||||
false,
|
||||
),
|
||||
Arguments.of(
|
||||
mock<Upstream> {
|
||||
val head = mock<Head> {
|
||||
on { getCurrentHeight() } doReturn null
|
||||
}
|
||||
on { getHead() } doReturn head
|
||||
},
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, -100),
|
||||
false,
|
||||
),
|
||||
Arguments.of(
|
||||
mock<Upstream> {
|
||||
val head = mock<Head> {
|
||||
on { getCurrentHeight() } doReturn 20
|
||||
}
|
||||
on { getHead() } doReturn head
|
||||
},
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, -100),
|
||||
false,
|
||||
),
|
||||
Arguments.of(
|
||||
mock<Upstream> {
|
||||
val head = mock<Head> {
|
||||
on { getCurrentHeight() } doReturn 150
|
||||
}
|
||||
on { getHead() } doReturn head
|
||||
},
|
||||
UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.HEAD, -100),
|
||||
true,
|
||||
),
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,95 @@
|
||||
package io.emeraldpay.dshackle.upstream.lowerbound
|
||||
|
||||
import io.emeraldpay.dshackle.Chain
|
||||
import org.junit.jupiter.api.Test
|
||||
import org.mockito.kotlin.doReturn
|
||||
import org.mockito.kotlin.mock
|
||||
import reactor.core.publisher.Flux
|
||||
import reactor.test.StepVerifier
|
||||
import java.time.Duration
|
||||
|
||||
class LowerBoundDetectorTest {
|
||||
|
||||
@Test
|
||||
fun `receive only manual bounds`() {
|
||||
val detector = TestLowerBoundDetector()
|
||||
val manualBoundService = mock<ManualLowerBoundService> {
|
||||
on { manualBoundTypes() } doReturn setOf(LowerBoundType.STATE, LowerBoundType.BLOCK)
|
||||
on { hasManualBound(LowerBoundType.STATE) } doReturn true
|
||||
on { hasManualBound(LowerBoundType.BLOCK) } doReturn true
|
||||
on { manualLowerBound(LowerBoundType.STATE) } doReturn LowerBoundData(50L, 1000, LowerBoundType.STATE)
|
||||
on { manualLowerBound(LowerBoundType.BLOCK) } doReturn LowerBoundData(80L, 1000, LowerBoundType.BLOCK)
|
||||
}
|
||||
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound(manualBoundService) }
|
||||
.expectSubscription()
|
||||
.expectNoEvent(Duration.ofSeconds(15))
|
||||
.expectNextSequence(
|
||||
setOf(
|
||||
LowerBoundData(50L, 1000, LowerBoundType.STATE),
|
||||
LowerBoundData(80L, 1000, LowerBoundType.BLOCK),
|
||||
),
|
||||
)
|
||||
.thenCancel()
|
||||
.verify(Duration.ofSeconds(3))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `receive one manual and one calculated bounds`() {
|
||||
val detector = TestLowerBoundDetector()
|
||||
val manualBoundService = mock<ManualLowerBoundService> {
|
||||
on { manualBoundTypes() } doReturn setOf(LowerBoundType.STATE)
|
||||
on { hasManualBound(LowerBoundType.STATE) } doReturn true
|
||||
on { hasManualBound(LowerBoundType.BLOCK) } doReturn false
|
||||
on { manualLowerBound(LowerBoundType.STATE) } doReturn LowerBoundData(50L, 1000, LowerBoundType.STATE)
|
||||
on { manualLowerBound(LowerBoundType.BLOCK) } doReturn null
|
||||
}
|
||||
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound(manualBoundService) }
|
||||
.expectSubscription()
|
||||
.expectNoEvent(Duration.ofSeconds(15))
|
||||
.expectNextSequence(
|
||||
setOf(
|
||||
LowerBoundData(5000L, 1000, LowerBoundType.BLOCK),
|
||||
LowerBoundData(50L, 1000, LowerBoundType.STATE),
|
||||
),
|
||||
)
|
||||
.thenCancel()
|
||||
.verify(Duration.ofSeconds(3))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `receive only calculated bounds`() {
|
||||
val detector = TestLowerBoundDetector()
|
||||
val manualBoundService = BaseManualLowerBoundService(mock(), emptyMap())
|
||||
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound(manualBoundService) }
|
||||
.expectSubscription()
|
||||
.expectNoEvent(Duration.ofSeconds(15))
|
||||
.expectNextSequence(
|
||||
setOf(
|
||||
LowerBoundData(1000L, 1000, LowerBoundType.STATE),
|
||||
LowerBoundData(5000L, 1000, LowerBoundType.BLOCK),
|
||||
),
|
||||
)
|
||||
.thenCancel()
|
||||
.verify(Duration.ofSeconds(3))
|
||||
}
|
||||
}
|
||||
|
||||
private class TestLowerBoundDetector : LowerBoundDetector(Chain.CORE__MAINNET) {
|
||||
override fun period(): Long {
|
||||
return 1
|
||||
}
|
||||
|
||||
override fun internalDetectLowerBound(): Flux<LowerBoundData> {
|
||||
return Flux.just(
|
||||
LowerBoundData(1000L, 1000, LowerBoundType.STATE),
|
||||
LowerBoundData(5000L, 1000, LowerBoundType.BLOCK),
|
||||
)
|
||||
}
|
||||
|
||||
override fun types(): Set<LowerBoundType> {
|
||||
return setOf(LowerBoundType.STATE, LowerBoundType.BLOCK)
|
||||
}
|
||||
}
|
||||
@@ -2,12 +2,14 @@ package io.emeraldpay.dshackle.upstream.solana
|
||||
|
||||
import io.emeraldpay.dshackle.Chain
|
||||
import io.emeraldpay.dshackle.Global
|
||||
import io.emeraldpay.dshackle.config.UpstreamsConfig
|
||||
import io.emeraldpay.dshackle.reader.ChainReader
|
||||
import io.emeraldpay.dshackle.upstream.ChainRequest
|
||||
import io.emeraldpay.dshackle.upstream.ChainResponse
|
||||
import io.emeraldpay.dshackle.upstream.Upstream
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundData
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.ManualLowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.rpcclient.ListParams
|
||||
import org.assertj.core.api.Assertions.assertThat
|
||||
import org.junit.jupiter.api.Test
|
||||
@@ -76,4 +78,68 @@ class SolanaLowerBoundServiceTest {
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `get solana lower block and slot with manual bounds`() {
|
||||
val reader = mock<ChainReader> {
|
||||
on { read(ChainRequest("getFirstAvailableBlock", ListParams())) } doReturn
|
||||
Mono.just(ChainResponse("25000000".toByteArray(), null))
|
||||
on {
|
||||
read(
|
||||
ChainRequest(
|
||||
"getBlock",
|
||||
ListParams(
|
||||
25000000L,
|
||||
mapOf(
|
||||
"showRewards" to false,
|
||||
"transactionDetails" to "none",
|
||||
"maxSupportedTransactionVersion" to 0,
|
||||
),
|
||||
),
|
||||
),
|
||||
)
|
||||
} doReturn Mono.just(
|
||||
ChainResponse(
|
||||
Global.objectMapper.writeValueAsBytes(
|
||||
mapOf(
|
||||
"blockHeight" to 21000000,
|
||||
"blockTime" to 111,
|
||||
"blockhash" to "22",
|
||||
"previousBlockhash" to "33",
|
||||
),
|
||||
),
|
||||
null,
|
||||
),
|
||||
)
|
||||
}
|
||||
val settings = UpstreamsConfig.AdditionalSettings(
|
||||
mapOf(
|
||||
LowerBoundType.SLOT to UpstreamsConfig.ManualBoundSetting(ManualLowerBoundType.FIXED, 54L),
|
||||
),
|
||||
)
|
||||
val upstream = mock<Upstream> {
|
||||
on { getIngressReader() } doReturn reader
|
||||
on { getChain() } doReturn Chain.UNSPECIFIED
|
||||
on { getAdditionalSettings() } doReturn settings
|
||||
}
|
||||
|
||||
val detector = SolanaLowerBoundService(Chain.UNSPECIFIED, upstream)
|
||||
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBounds() }
|
||||
.expectSubscription()
|
||||
.expectNoEvent(Duration.ofSeconds(15))
|
||||
.expectNextMatches { it.lowerBound == 21000000L && it.type == LowerBoundType.STATE }
|
||||
.expectNextMatches { it.lowerBound == 54L && it.type == LowerBoundType.SLOT }
|
||||
.thenCancel()
|
||||
.verify(Duration.ofSeconds(3))
|
||||
|
||||
assertThat(detector.getLowerBounds().toList())
|
||||
.usingRecursiveFieldByFieldElementComparatorIgnoringFields("timestamp")
|
||||
.hasSameElementsAs(
|
||||
listOf(
|
||||
LowerBoundData(21000000L, LowerBoundType.STATE),
|
||||
LowerBoundData(54L, LowerBoundType.SLOT),
|
||||
),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package io.emeraldpay.dshackle.upstream.starknet
|
||||
import io.emeraldpay.dshackle.Chain
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundData
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType
|
||||
import io.emeraldpay.dshackle.upstream.lowerbound.NoopManualLowerBoundService
|
||||
import org.junit.jupiter.api.Test
|
||||
import reactor.test.StepVerifier
|
||||
import java.time.Duration
|
||||
@@ -13,7 +14,7 @@ class StarknetLowerBoundStateDetectorTest {
|
||||
fun `starknet lower block is 1`() {
|
||||
val detector = StarknetLowerBoundStateDetector(Chain.STARKNET__MAINNET)
|
||||
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound() }
|
||||
StepVerifier.withVirtualTime { detector.detectLowerBound(NoopManualLowerBoundService()) }
|
||||
.expectSubscription()
|
||||
.expectNoEvent(Duration.ofSeconds(15))
|
||||
.expectNext(LowerBoundData(1, LowerBoundType.STATE))
|
||||
|
||||
Reference in New Issue
Block a user