From a34f2186ccb33a5a75e826939e7850af8d6f91ed Mon Sep 17 00:00:00 2001 From: EugeneDrpc Date: Fri, 8 Aug 2025 18:03:10 +0200 Subject: [PATCH] Log index validator (#699) * log index validator - wip * fix: Handle edge cases in LogIndexValidator and fix GitHub Actions gix critical edge case: when first transaction has no logs, second transaction starting at logIndex=0 is VALID behavior null handling fix 'repo_token' parameter * - Use Mono.empty() instead of Mono.just(null) - Use defaultIfEmpty for handling empty results * wip : add tests for LogIndexValidator * test fixes * fix: Fix ktlint issues in LogIndexValidatorTest - Remove wildcard import - Fix trailing spaces - Add missing trailing commas - Fix import order - Add newline at end of file * add disable option * add logs for error case * save prevState on doOnNext * also check for first tx log index order * add not continius indexes test * make code thread safe * use upstream.getHead * search CHECK_TX_COUNT instead of 0 and 1 indexes * refactor double rpc call --- .github/workflows/test.yaml | 2 +- .../dshackle/foundation/ChainOptions.kt | 6 +- .../dshackle/foundation/ChainOptionsReader.kt | 3 + .../ethereum/EthereumChainSpecific.kt | 7 + .../upstream/ethereum/LogIndexValidator.kt | 382 ++++++++++++++++++ .../config/UpstreamsConfigReaderSpec.groovy | 2 +- .../ethereum/LogIndexValidatorTest.kt | 363 +++++++++++++++++ 7 files changed, 762 insertions(+), 3 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidator.kt create mode 100644 src/test/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidatorTest.kt diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index f2a3139a..3311e335 100644 --- a/.github/workflows/test.yaml +++ b/.github/workflows/test.yaml @@ -38,7 +38,7 @@ jobs: - name: Setup gradle uses: gradle/gradle-build-action@v2 with: - repo_token: ${{ secrets.GITHUB_TOKEN }} + github-token: ${{ secrets.GITHUB_TOKEN }} - name: Check run: make test env: diff --git a/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptions.kt b/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptions.kt index 645d04a0..16470c70 100644 --- a/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptions.kt +++ b/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptions.kt @@ -19,6 +19,7 @@ class ChainOptions { val disableLivenessSubscriptionValidation: Boolean, val disableBoundValidation: Boolean = false, val valdateErigonBug: Boolean, + val disableLogIndexValidation: Boolean = false, ) data class DefaultOptions( @@ -41,7 +42,8 @@ class ChainOptions { var callLimitSize: Int? = null, var disableLivenessSubscriptionValidation: Boolean? = null, var disableBoundValidation: Boolean? = null, - var validateErigonBug: Boolean? = null + var validateErigonBug: Boolean? = null, + var disableLogIndexValidation: Boolean? = null ) { companion object { @JvmStatic @@ -73,6 +75,7 @@ class ChainOptions { copy.disableLivenessSubscriptionValidation = overwrites.disableLivenessSubscriptionValidation ?: this.disableLivenessSubscriptionValidation copy.disableBoundValidation = overwrites.disableBoundValidation ?: this.disableBoundValidation copy.validateErigonBug = overwrites.validateErigonBug ?: this.validateErigonBug + copy.disableLogIndexValidation = overwrites.disableLogIndexValidation ?: this.disableLogIndexValidation return copy } @@ -93,6 +96,7 @@ class ChainOptions { this.disableLivenessSubscriptionValidation ?: false, this.disableBoundValidation ?: false, this.validateErigonBug ?: true, + this.disableLogIndexValidation ?: false, ) } } diff --git a/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptionsReader.kt b/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptionsReader.kt index 1fac61a6..b2b951b8 100644 --- a/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptionsReader.kt +++ b/foundation/src/main/kotlin/io/emeraldpay/dshackle/foundation/ChainOptionsReader.kt @@ -58,6 +58,9 @@ class ChainOptionsReader : YamlConfigReader() { getValueAsBool(values, "disable-bound-validation")?.let { options.disableBoundValidation = it } + getValueAsBool(values, "disable-log-index-validation")?.let { + options.disableLogIndexValidation = it + } return options } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumChainSpecific.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumChainSpecific.kt index ca906214..c19c56be 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumChainSpecific.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumChainSpecific.kt @@ -160,6 +160,13 @@ object EthereumChainSpecific : AbstractPollChainSpecific() { if (options.valdateErigonBug) { validators.add(ErigonBuggedValidator(upstream)) } + + // Add LogIndexValidator to detect incorrect logIndex numbering + // Can be disabled via options.disableLogIndexValidation + if (!options.disableLogIndexValidation) { + validators.add(LogIndexValidator(upstream)) + } + val limitValidator = EthCallLimitValidator(upstream, options, config) if (limitValidator.isEnabled()) { validators.add(limitValidator) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidator.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidator.kt new file mode 100644 index 00000000..b038c592 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidator.kt @@ -0,0 +1,382 @@ +/** + * Copyright (c) 2024 DRPC + * + * 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 + +import com.fasterxml.jackson.databind.JsonNode +import io.emeraldpay.dshackle.Global +import io.emeraldpay.dshackle.upstream.ChainRequest +import io.emeraldpay.dshackle.upstream.ChainResponse +import io.emeraldpay.dshackle.upstream.SingleValidator +import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult +import io.emeraldpay.dshackle.upstream.rpcclient.ListParams +import org.slf4j.Logger +import org.slf4j.LoggerFactory +import reactor.core.publisher.Mono +import reactor.kotlin.extra.retry.retryRandomBackoff +import java.time.Duration +import java.util.concurrent.atomic.AtomicInteger +import java.util.concurrent.atomic.AtomicReference +/** + * Validator to detect incorrect logIndex numbering in Ethereum nodes. + * * Some Erigon nodes (versions 3.0.8-3.0.14) return local logIndex within transaction + * instead of global logIndex within block. This validator detects such misconfiguration. + * * Detection principle: + * - In correct implementation: logIndex is global across the entire block + * - In buggy implementation: logIndex resets to 0 for each transaction + * - Test: Check if the first log in the second transaction has logIndex = 0 + * (which would indicate local numbering) + */ +class LogIndexValidator( + private val upstream: Upstream, +) : SingleValidator { + + private data class TransactionValidationData( + val firstTxHash: String, + val secondTxHash: String, + val firstReceipt: JsonNode, + val secondReceipt: JsonNode, + ) + + companion object { + @JvmStatic + val log: Logger = LoggerFactory.getLogger(LogIndexValidator::class.java) + + // Check every N validations to save resources + private const val CHECK_FREQUENCY = 10 + private const val CHECK_TX_COUNT = 6 + + // Maximum attempts to find a suitable block for validation + private const val MAX_BLOCK_SEARCH_ATTEMPTS = 5 + } + + private val callCount = AtomicInteger(0) + private val lastResult = AtomicReference(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + + override fun validate(onError: ValidateUpstreamSettingsResult): Mono { + // Perform validation periodically to save resources + val currentCount = callCount.incrementAndGet() + if ((currentCount - 1) % CHECK_FREQUENCY != 0) { + return Mono.just(lastResult.get()) + } + + log.debug("Starting logIndex validation for upstream ${upstream.getId()}, check #${currentCount / CHECK_FREQUENCY}") + + return findBlockWithLogsAndValidate() + .doOnNext { result -> + // Update lastResult only when we actually performed validation + lastResult.set(result) + if (result == ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) { + log.error("LogIndex validation failed for upstream ${upstream.getId()}, marking as fatal") + } + } + .timeout(Duration.ofSeconds(15)) + .retryRandomBackoff(2, Duration.ofMillis(200), Duration.ofMillis(1000)) { ctx -> + log.debug( + "Retry logIndex validation for ${upstream.getId()}, iteration ${ctx.iteration()}, " + + "error: ${ctx.exception().message}", + ) + } + .onErrorResume { err -> + log.warn("Error during logIndex validation for ${upstream.getId()}: ${err.message}") + // In case of error, return last known state to avoid false positives + // We only mark as invalid if we explicitly detect the bug + Mono.just(lastResult.get()) + } + } + + /** + * Find a suitable block with at least 2 transactions that have logs and validate it + */ + private fun findBlockWithLogsAndValidate(): Mono { + // Try to use upstream's head height to avoid extra RPC call + return try { + val head = upstream.getHead() + val currentHeight = head.getCurrentHeight() + if (currentHeight != null && currentHeight > 0) { + searchForSuitableBlock(currentHeight, 0) + } else { + // No head height available, skip validation + log.debug("No head height available for ${upstream.getId()}, skipping validation") + Mono.just(lastResult.get()) + } + } catch (e: Exception) { + // If getHead() is not available (e.g., in tests), fallback to RPC for compatibility + getLatestBlockNumber() + .flatMap { latestBlockNum -> + searchForSuitableBlock(latestBlockNum, 0) + } + .onErrorResume { + log.debug("Cannot get block number for ${upstream.getId()}: ${it.message}, skipping validation") + Mono.just(lastResult.get()) + } + } + } + + /** + * Get the latest block number - used as fallback when head is not available + */ + private fun getLatestBlockNumber(): Mono { + return upstream.getIngressReader() + .read(ChainRequest("eth_blockNumber", ListParams())) + .flatMap(ChainResponse::requireStringResult) + .map { it.removePrefix("0x").toLong(16) } + } + + /** + * Search backwards from the latest block to find one suitable for validation + */ + private fun searchForSuitableBlock( + startBlockNum: Long, + attempt: Int, + ): Mono { + if (attempt >= MAX_BLOCK_SEARCH_ATTEMPTS || startBlockNum <= 0) { + log.debug("Could not find suitable block for logIndex validation in ${upstream.getId()} after $attempt attempts, returning last result: ${lastResult.get()}") + return Mono.just(lastResult.get()) + } + + val blockNum = startBlockNum - attempt + val blockHex = "0x${blockNum.toString(16)}" + + return getBlock(blockHex) + .flatMap { block -> + val transactions = block.get("transactions") + if (transactions != null && transactions.isArray && transactions.size() >= 2) { + // We have at least 2 transactions, now check if they have logs + validateBlockTransactions(transactions) + } else { + // Not enough transactions, try previous block + searchForSuitableBlock(startBlockNum, attempt + 1) + } + } + .onErrorResume { + // Error getting block, try previous one + searchForSuitableBlock(startBlockNum, attempt + 1) + } + } + + /** + * Validate transactions in a block to check for logIndex bug + */ + private fun validateBlockTransactions(transactions: JsonNode): Mono { + // We need at least 2 transactions to validate + if (transactions.size() < 2) { + log.debug("Block has less than 2 transactions, cannot validate, returning last result: ${lastResult.get()}") + return Mono.just(lastResult.get()) + } + + // Try to find two transactions with logs (receipts are already fetched) + return findTwoTransactionsWithLogs(transactions) + .map { validationData -> + val firstLogs = validationData.firstReceipt.get("logs") + val secondLogs = validationData.secondReceipt.get("logs") + + // Both transactions are guaranteed to have logs (found by search function) + validateLogIndices( + firstLogs, + secondLogs, + validationData.firstTxHash, + validationData.secondTxHash, + ) + } + .defaultIfEmpty(lastResult.get()) + } + + /** + * Find first two transactions that both have logs (not necessarily consecutive) + * Returns validation data with receipts to avoid double RPC calls + */ + private fun findTwoTransactionsWithLogs(transactions: JsonNode): Mono { + // Search for first two transactions that have logs (not necessarily consecutive) + val maxSearchCount = minOf(CHECK_TX_COUNT, transactions.size()) // Limit search to avoid excessive RPC calls + + return searchForTransactionsWithLogs(transactions, 0, maxSearchCount, mutableListOf()) + } + + /** + * Recursively search for transactions with logs, checking receipts one by one + * Stops as soon as we find two transactions with logs + */ + private fun searchForTransactionsWithLogs( + transactions: JsonNode, + currentIndex: Int, + maxSearchCount: Int, + foundTxData: MutableList>, + ): Mono { + // Stop if we found enough transactions or reached the limit + if (foundTxData.size >= 2) { + val first = foundTxData[0] + val second = foundTxData[1] + return Mono.just( + TransactionValidationData( + firstTxHash = first.first, + secondTxHash = second.first, + firstReceipt = first.second, + secondReceipt = second.second, + ), + ) + } + + if (currentIndex >= maxSearchCount || currentIndex >= transactions.size()) { + return Mono.empty() // Not enough transactions found + } + + val txHash = transactions[currentIndex].get("hash")?.asText() + if (txHash == null) { + log.warn("Transaction at index $currentIndex has no hash, aborting validation due to corrupted block data") + return Mono.empty() + } + + return getTransactionReceipt(txHash) + .flatMap { receipt -> + val logs = receipt.get("logs") + if (logs != null && logs.isArray && logs.size() > 0) { + // This transaction has logs, add it to our list + foundTxData.add(Pair(txHash, receipt)) + log.debug("Found transaction with logs: $txHash (${logs.size()} logs)") + } + + // Continue searching - either we have enough or need to find more + searchForTransactionsWithLogs(transactions, currentIndex + 1, maxSearchCount, foundTxData) + } + .onErrorResume { error -> + log.warn("Error getting receipt for $txHash: ${error.message}, aborting validation to avoid false results") + Mono.empty() + } + } + + /** + * Validate that logIndex is global across the block, not local to transaction + */ + private fun validateLogIndices( + firstTxLogs: JsonNode, + secondTxLogs: JsonNode, + firstTxHash: String, + secondTxHash: String, + ): ValidateUpstreamSettingsResult { + // CRITICAL: We need BOTH transactions to have logs to detect the bug + // If first tx has no logs and second tx starts at 0, that's CORRECT behavior + + // Parse all log indices safely + val firstTxFirstLogIndex = if (firstTxLogs.size() > 0) { + parseLogIndex(firstTxLogs[0].get("logIndex")?.asText()) + } else { + null + } + + val firstTxLastLogIndex = if (firstTxLogs.size() > 0) { + parseLogIndex(firstTxLogs[firstTxLogs.size() - 1].get("logIndex")?.asText()) + } else { + null + } + + val secondTxFirstLogIndex = if (secondTxLogs.size() > 0) { + parseLogIndex(secondTxLogs[0].get("logIndex")?.asText()) + } else { + null + } + + // Can't validate without proper data + if (firstTxFirstLogIndex == null || firstTxLastLogIndex == null || secondTxFirstLogIndex == null) { + log.debug( + "Cannot validate logIndex for {}: missing log data, returning last result: {}", + upstream.getId(), + lastResult.get(), + ) + return lastResult.get() + } + + // Validation logic: + // 1. First transaction's first log MUST start at 0 (error if not) + // 2. Second transaction MUST continue from where first ended (error if resets to 0) + + if (firstTxFirstLogIndex != 0L) { + // This is also a critical issue - first log in block should always be 0 + log.error( + "Node ${upstream.getId()} has incorrect logIndex start in first transaction $firstTxHash: " + + "first log has logIndex=$firstTxFirstLogIndex instead of 0. " + + "This indicates a serious issue with logIndex numbering.", + ) + return ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR + } + + // The key check: if first tx has logs AND second tx starts at 0, it's a BUG + val expectedSecondTxStart = firstTxLastLogIndex + 1 + + if (secondTxFirstLogIndex == 0L && firstTxLogs.size() > 0) { + // This is the bug! Second transaction should not start at 0 when first has logs + log.error( + "Node ${upstream.getId()} uses LOCAL logIndex instead of GLOBAL. " + + "First tx ($firstTxHash) has ${firstTxLogs.size()} logs (indices 0-$firstTxLastLogIndex), " + + "but second tx ($secondTxHash) starts at logIndex=0 instead of $expectedSecondTxStart. " + + "This indicates Erigon bug with incorrect logIndex numbering.", + ) + return ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR + } + + // Additional validation: check if indices are continuous + if (secondTxFirstLogIndex != expectedSecondTxStart) { + log.error( + "Node ${upstream.getId()} has non-continuous logIndex between transactions $firstTxHash and $secondTxHash: " + + "second transaction starts at $secondTxFirstLogIndex, expected $expectedSecondTxStart. " + + "This indicates missing or incorrectly numbered logs.", + ) + return ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR + } + + // If we got here, everything is correct + return ValidateUpstreamSettingsResult.UPSTREAM_VALID + } + + /** + * Parse logIndex from hex string to Long + * Returns null if parsing fails or input is invalid + */ + private fun parseLogIndex(logIndexHex: String?): Long? { + if (logIndexHex == null || logIndexHex.isBlank()) return null + return try { + val cleaned = logIndexHex.trim().removePrefix("0x").removePrefix("0X") + if (cleaned.isEmpty()) { + 0L // "0x" or "0x0" -> 0 + } else { + cleaned.toLong(16) + } + } catch (e: NumberFormatException) { + log.warn("Failed to parse logIndex: $logIndexHex", e) + null + } + } + + /** + * Get block by number with full transaction details + */ + private fun getBlock(blockNumber: String): Mono { + return upstream.getIngressReader() + .read(ChainRequest("eth_getBlockByNumber", ListParams(blockNumber, true))) + .flatMap(ChainResponse::requireResult) + .map { Global.objectMapper.readTree(it) } + } + + /** + * Get transaction receipt by hash + */ + private fun getTransactionReceipt(txHash: String): Mono { + return upstream.getIngressReader() + .read(ChainRequest("eth_getTransactionReceipt", ListParams(txHash))) + .flatMap(ChainResponse::requireResult) + .map { Global.objectMapper.readTree(it) } + } +} diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy index 92d85b56..2f8b83ed 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy @@ -636,7 +636,7 @@ class UpstreamsConfigReaderSpec extends Specification { def options = partialOptions.buildOptions() then: options == new ChainOptions.Options( - false, false, 30, Duration.ofSeconds(60), null, true, 1, true, true, true, true, 1_000_000, false, false, true + false, false, 30, Duration.ofSeconds(60), null, true, 1, true, true, true, true, 1_000_000, false, false, true, false ) } } diff --git a/src/test/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidatorTest.kt b/src/test/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidatorTest.kt new file mode 100644 index 00000000..a28fdadf --- /dev/null +++ b/src/test/kotlin/io/emeraldpay/dshackle/upstream/ethereum/LogIndexValidatorTest.kt @@ -0,0 +1,363 @@ +/** + * Copyright (c) 2024 DRPC + * + * 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 + +import io.emeraldpay.dshackle.reader.ChainReader +import io.emeraldpay.dshackle.upstream.ChainCallError +import io.emeraldpay.dshackle.upstream.ChainRequest +import io.emeraldpay.dshackle.upstream.ChainResponse +import io.emeraldpay.dshackle.upstream.Upstream +import io.emeraldpay.dshackle.upstream.ValidateUpstreamSettingsResult +import io.emeraldpay.dshackle.upstream.rpcclient.ListParams +import org.junit.jupiter.api.BeforeEach +import org.junit.jupiter.api.Test +import org.mockito.kotlin.argThat +import org.mockito.kotlin.doReturn +import org.mockito.kotlin.mock +import reactor.core.publisher.Mono +import reactor.test.StepVerifier +import java.time.Duration + +class LogIndexValidatorTest { + + private lateinit var reader: ChainReader + private lateinit var upstream: Upstream + private lateinit var validator: LogIndexValidator + + @BeforeEach + fun setup() { + reader = mock {} + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + } + + @Test + fun `detects local logIndex numbering bug - basic case`() { + setupMockForBugDetection( + firstTxLogs = listOf("0x0", "0x1", "0x2"), + secondTxLogs = listOf("0x0", "0x1"), // BUG: should be 0x3, 0x4 + ) + + // First call (callCount=0) triggers validation immediately + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test + fun `validates correct global logIndex numbering`() { + setupMockForBugDetection( + firstTxLogs = listOf("0x0", "0x1", "0x2"), + secondTxLogs = listOf("0x3", "0x4"), // CORRECT: continues globally + ) + + // First call triggers validation immediately + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test fun `handles single log transactions correctly`() { + setupMockForBugDetection( + firstTxLogs = listOf("0x0"), + secondTxLogs = listOf("0x0"), // BUG: should be 0x1 + ) + + // First call triggers validation immediately + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test + fun `handles transactions without logs gracefully`() { + reader = mock { + on { read(ChainRequest("eth_blockNumber", ListParams())) } doReturn + Mono.just(ChainResponse("\"0x1000\"".toByteArray(), null)) + + on { read(ChainRequest("eth_getBlockByNumber", ListParams("0x1000", true))) } doReturn + Mono.just(ChainResponse(createBlock("0xaaa", "0xbbb").toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams("0xaaa"))) } doReturn + Mono.just(ChainResponse("""{"logs": []}""".toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams("0xbbb"))) } doReturn + Mono.just(ChainResponse("""{"logs": []}""".toByteArray(), null)) + } + + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test + fun `handles block with insufficient transactions`() { + reader = mock { + on { read(ChainRequest("eth_blockNumber", ListParams())) } doReturn + Mono.just(ChainResponse("\"0x1000\"".toByteArray(), null)) + + // Block with only 1 transaction + on { read(ChainRequest("eth_getBlockByNumber", ListParams("0x1000", true))) } doReturn + Mono.just(ChainResponse("""{"transactions": [{"hash": "0xaaa"}]}""".toByteArray(), null)) + + // Should try previous block + on { read(ChainRequest("eth_getBlockByNumber", ListParams("0xfff", true))) } doReturn + Mono.just(ChainResponse(createBlock("0xccc", "0xddd").toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams("0xccc"))) } doReturn + Mono.just(ChainResponse(createReceipt(listOf("0x0")).toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams("0xddd"))) } doReturn + Mono.just(ChainResponse(createReceipt(listOf("0x1")).toByteArray(), null)) + } + + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test + fun `performs validation only once every 10 calls`() { + // First call (callCount=0) will perform validation (0 % 10 == 0) + // So we need to set up mocks for the first call + setupMockForBugDetection( + firstTxLogs = listOf("0x0", "0x1"), + secondTxLogs = listOf("0x2", "0x3"), // Correct numbering + ) + + // First call - performs validation (callCount=0) + val firstResult = validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR).block() + assert(firstResult == ValidateUpstreamSettingsResult.UPSTREAM_VALID) { "First call should validate and return VALID" } + + // Calls 2-10 should skip validation and return VALID immediately + repeat(9) { _ -> + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + .expectComplete() + .verify(Duration.ofMillis(100)) // Should be fast since no actual validation + } + + // 11th call (callCount=10) should perform validation again + val eleventhResult = validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR).block() + assert(eleventhResult == ValidateUpstreamSettingsResult.UPSTREAM_VALID) { "11th call should validate again" } + } + + @Test + fun `handles RPC errors gracefully`() { + reader = mock { + on { read(ChainRequest("eth_blockNumber", ListParams())) } doReturn + Mono.just(ChainResponse(null, ChainCallError(123, "Node error"))) + } + + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_VALID) // Should not fail upstream + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test + fun `detects non-continuous logIndex as error`() { + setupMockForBugDetection( + firstTxLogs = listOf("0x0", "0x1", "0x2"), + secondTxLogs = listOf("0x5", "0x6"), // Gap - should be 0x3, 0x4 + ) + + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) // Should fail due to gap + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + // Helper methods + private fun setupMockForBugDetection(firstTxLogs: List, secondTxLogs: List) { + reader = mock { + on { read(ChainRequest("eth_blockNumber", ListParams())) } doReturn + Mono.just(ChainResponse("\"0x1000\"".toByteArray(), null)) + + on { read(ChainRequest("eth_getBlockByNumber", ListParams("0x1000", true))) } doReturn + Mono.just(ChainResponse(createBlock("0xaaa", "0xbbb").toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams("0xaaa"))) } doReturn + Mono.just(ChainResponse(createReceipt(firstTxLogs).toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams("0xbbb"))) } doReturn + Mono.just(ChainResponse(createReceipt(secondTxLogs).toByteArray(), null)) + } + + // Update upstream mock with new reader + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + } + + private fun setupReceiptsForTxs( + firstHash: String, + secondHash: String, + firstLogs: List, + secondLogs: List, + ) { + reader = mock { + on { read(ChainRequest("eth_blockNumber", ListParams())) } doReturn + Mono.just(ChainResponse("\"0x1000\"".toByteArray(), null)) + + on { read(ChainRequest("eth_getBlockByNumber", ListParams("0x1000", true))) } doReturn + Mono.just(ChainResponse(createBlock(firstHash, secondHash).toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams(firstHash))) } doReturn + Mono.just(ChainResponse(createReceipt(firstLogs).toByteArray(), null)) + + on { read(ChainRequest("eth_getTransactionReceipt", ListParams(secondHash))) } doReturn + Mono.just(ChainResponse(createReceipt(secondLogs).toByteArray(), null)) + } + + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + } + + private fun createBlock(tx1Hash: String, tx2Hash: String) = """ + { + "transactions": [ + {"hash": "$tx1Hash"}, + {"hash": "$tx2Hash"} + ] + } + """.trimIndent() + + private fun createReceipt(logIndexes: List) = """ + { + "logs": [ + ${logIndexes.joinToString(",") { """{"logIndex": "$it"}""" }} + ] + } + """.trimIndent() + + @Test + fun `detects incorrect logIndex start in first transaction`() { + setupMockForBugDetection( + firstTxLogs = listOf("0x5", "0x6", "0x7"), // BUG: should start from 0x0 + secondTxLogs = listOf("0x8", "0x9"), + ) + + // Should detect that first transaction doesn't start at 0 + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } + + @Test + fun `remembers error state between validations`() { + // First call detects bug + setupMockForBugDetection( + firstTxLogs = listOf("0x0", "0x1", "0x2"), + secondTxLogs = listOf("0x0", "0x1"), // BUG: should be 0x3, 0x4 + ) + + // First call (callCount=0) triggers validation and detects bug + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) + .expectComplete() + .verify(Duration.ofSeconds(3)) + + // Calls 2-10 should skip validation but return the error state + repeat(9) { _ -> + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_FATAL_SETTINGS_ERROR) // Should remember error + .expectComplete() + .verify(Duration.ofMillis(100)) + } + } + + @Test + fun `returns last result when cannot find suitable block`() { + // Setup: only blocks with 1 transaction + reader = mock { + on { read(ChainRequest("eth_blockNumber", ListParams())) } doReturn + Mono.just(ChainResponse("\"0x1000\"".toByteArray(), null)) + + // All blocks return only 1 transaction + on { read(argThat { method == "eth_getBlockByNumber" }) } doReturn + Mono.just(ChainResponse("""{"transactions": [{"hash": "0xaaa"}]}""".toByteArray(), null)) + } + + upstream = mock { + on { getIngressReader() } doReturn reader + on { getId() } doReturn "test-upstream" + } + validator = LogIndexValidator(upstream) + + // First validation should return UPSTREAM_VALID (initial state) + StepVerifier.create( + validator.validate(ValidateUpstreamSettingsResult.UPSTREAM_SETTINGS_ERROR), + ) + .expectNext(ValidateUpstreamSettingsResult.UPSTREAM_VALID) + .expectComplete() + .verify(Duration.ofSeconds(3)) + } +}