diff --git a/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/block/BlockManager.java b/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/block/BlockManager.java index 477cc6d2f6d..d07637fb288 100644 --- a/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/block/BlockManager.java +++ b/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/block/BlockManager.java @@ -13,10 +13,14 @@ package tech.pegasys.teku.statetransition.block; +import static tech.pegasys.teku.statetransition.validation.InternalValidationResult.ACCEPT; + import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.concurrent.RejectedExecutionException; import java.util.function.Supplier; import org.apache.logging.log4j.LogManager; @@ -64,7 +68,7 @@ public class BlockManager extends Service private final TimeProvider timeProvider; private final EventLogger eventLogger; - private final FutureItems futureBlocks; + private final FutureBlockTracker futureBlockTracker; // in the invalidBlockRoots map we are going to store blocks whose import result is invalid // and will not require any further retry. Descendants of these blocks will be considered invalid // as well. @@ -95,7 +99,7 @@ public BlockManager( this.blockEventsListener = blockEventsListener; this.executionPayloadEventsListenerSupplier = executionPayloadEventsListenerSupplier; this.pendingBlockPool = pendingBlockPool; - this.futureBlocks = futureBlocks; + this.futureBlockTracker = new FutureBlockTracker(futureBlocks); this.invalidBlockRoots = invalidBlockRoots; this.blockValidator = blockValidator; this.timeProvider = timeProvider; @@ -124,7 +128,7 @@ public SafeFuture importBlock( blockValidator.initiateBroadcastValidation(block, broadcastValidationLevel); final SafeFuture importResult = - doImportBlock(block, Optional.empty(), blockBroadcastValidator, origin); + doImportBlock(block, Optional.empty(), blockBroadcastValidator, origin, false); // we want to intercept any early import exceptions happening before the consensus validation is // completed @@ -172,7 +176,8 @@ public SafeFuture validateAndImportBlock( block, blockImportPerformance, BlockBroadcastValidator.NOOP, - Optional.of(RemoteOrigin.GOSSIP)) + Optional.of(RemoteOrigin.GOSSIP), + result.isSaveForFuture()) .finish(err -> LOG.error("Failed to process received block.", err)); // block failed gossip validation, let's drop it from the pool, so it won't be served @@ -187,8 +192,7 @@ public SafeFuture validateAndImportBlock( @Override public void onSlot(final UInt64 slot) { - futureBlocks.onSlot(slot); - futureBlocks.prune(slot).forEach(this::importBlockIgnoringResult); + futureBlockTracker.prune(slot).forEach(this::processFutureBlock); } public void subscribeFailedPayloadExecution(final FailedPayloadExecutionSubscriber subscriber) { @@ -248,20 +252,51 @@ public void onExecutionPayloadImported( private void importBlockIgnoringResult(final SignedBeaconBlock block) { // we don't care about origin here because flow calls this function for retries only - doImportBlock(block, Optional.empty(), BlockBroadcastValidator.NOOP, Optional.empty()) + doImportBlock(block, Optional.empty(), BlockBroadcastValidator.NOOP, Optional.empty(), false) .finishStackTrace(); } + private void processFutureBlock(final QueuedFutureBlock futureBlock) { + // Future-block retries need to handle the exact sequence below. + // + // 1. Gossip validation may accept a near-future block because of clock tolerance. + // 2. Import can still return BLOCK_IS_FROM_FUTURE, so we queue it for later. + // 3. When the slot arrives, retry gossip may return IGNORE_ALREADY_SEEN. + // 4. In that case, retry the import without gossip validation so the block is not lost. + // 5. Retry gossip may also return IGNORE_EQUIVOCATION_DETECTED. + // 6. In that case, drop the block listeners and stop. + if (futureBlock.needsGossipValidation()) { + validateAndImportBlock(futureBlock.block(), Optional.empty()) + .thenAccept( + result -> { + if (result.isIgnoreAlreadySeen()) { + importBlockIgnoringResult(futureBlock.block()); + } else if (result.isIgnoreEquivocationDetected()) { + blockEventsListener.removeAllForBlock(futureBlock.block().getSlotAndBlockRoot()); + } + }) + .finishError(LOG); + } else { + importBlockIgnoringResult(futureBlock.block()); + } + } + private SafeFuture doImportBlock( final SignedBeaconBlock block, final Optional blockImportPerformance, final BlockBroadcastValidator blockBroadcastValidator, - final Optional origin) { + final Optional origin, + final boolean needsGossipValidationOnRetry) { return handleInvalidBlock(block) .or(() -> handleKnownBlock(block)) .orElseGet( () -> - handleBlockImport(block, blockImportPerformance, blockBroadcastValidator, origin) + handleBlockImport( + block, + blockImportPerformance, + blockBroadcastValidator, + origin, + needsGossipValidationOnRetry) .thenPeek( result -> lateBlockImportCheck(blockImportPerformance, block, result))); } @@ -288,7 +323,7 @@ private Optional> handleInvalidBlock( } private Optional> handleKnownBlock(final SignedBeaconBlock block) { - if (pendingBlockPool.contains(block) || futureBlocks.contains(block)) { + if (pendingBlockPool.contains(block) || futureBlockTracker.contains(block)) { // Pending and future blocks can't have been executed yet so must be marked optimistic return Optional.of(SafeFuture.completedFuture(BlockImportResult.knownBlock(block, true))); } @@ -303,7 +338,8 @@ private SafeFuture handleBlockImport( final SignedBeaconBlock block, final Optional blockImportPerformance, final BlockBroadcastValidator blockBroadcastValidator, - final Optional origin) { + final Optional origin, + final boolean needsGossipValidationOnRetry) { blockEventsListener.onNewBlock(block, origin); preImportBlockSubscribers.deliver(l -> l.onNewBlock(block, origin)); @@ -331,7 +367,8 @@ private SafeFuture handleBlockImport( case UNKNOWN_PARENT_EXECUTION_PAYLOAD -> { addBlockPendingParentExecutionPayload(block); } - case BLOCK_IS_FROM_FUTURE -> futureBlocks.add(block); + case BLOCK_IS_FROM_FUTURE -> + futureBlockTracker.add(block, needsGossipValidationOnRetry); case FAILED_EXECUTION_PAYLOAD_EXECUTION_SYNCING -> { LOG.warn( "Unable to import block {} with execution payload {}: Execution Client is still syncing", @@ -515,4 +552,35 @@ public interface PreImportBlockListener { public interface RequiredParentExecutionPayloadSubscriber { void onRequiredParentExecutionPayload(ParentExecutionPayloadDependency dependency); } + + private static final class FutureBlockTracker { + private final FutureItems futureBlocks; + private final Set gossipRetryRoots = new HashSet<>(); + + private FutureBlockTracker(final FutureItems futureBlocks) { + this.futureBlocks = futureBlocks; + } + + synchronized void add( + final SignedBeaconBlock block, final boolean needsGossipValidationOnRetry) { + // FutureItems drops blocks that are beyond the future-slot tolerance, so only remember the + // retry mode when the block actually made it into the queue. + if (futureBlocks.add(block) && needsGossipValidationOnRetry) { + gossipRetryRoots.add(block.getRoot()); + } + } + + synchronized List prune(final UInt64 slot) { + futureBlocks.onSlot(slot); + return futureBlocks.prune(slot).stream() + .map(block -> new QueuedFutureBlock(block, gossipRetryRoots.remove(block.getRoot()))) + .toList(); + } + + synchronized boolean contains(final SignedBeaconBlock block) { + return futureBlocks.contains(block); + } + } + + private record QueuedFutureBlock(SignedBeaconBlock block, boolean needsGossipValidation) {} } diff --git a/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/util/FutureItems.java b/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/util/FutureItems.java index 5cf402703b5..4b9cde2f235 100644 --- a/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/util/FutureItems.java +++ b/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/util/FutureItems.java @@ -74,20 +74,22 @@ public void onSlot(final UInt64 slot) { } /** - * Add a item to the future items set + * Add an item to the future items set. * - * @param item The item to add + * @param item the item to add + * @return true if the item was accepted for future processing, even if it was already queued */ - public void add(final T item) { + public boolean add(final T item) { final UInt64 slot = slotFunction.apply(item); if (slot.isGreaterThan(currentSlot.plus(futureSlotTolerance))) { // Item is too far in the future - return; + return false; } LOG.trace("Save future item at slot {} for later import: {}", slot, item); queuedFutureItems.computeIfAbsent(slot, key -> createNewSet()).add(item); futureItemsCounter.set(size(), type); + return true; } /** diff --git a/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/validation/InternalValidationResult.java b/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/validation/InternalValidationResult.java index ef39d88648e..82e5837ac7d 100644 --- a/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/validation/InternalValidationResult.java +++ b/ethereum/statetransition/src/main/java/tech/pegasys/teku/statetransition/validation/InternalValidationResult.java @@ -141,6 +141,13 @@ public boolean isIgnoreAlreadySeen() { .orElse(false); } + public boolean isIgnoreEquivocationDetected() { + return isIgnore() + && this.validationResultSubCode + .map(subCode -> subCode.equals(ValidationResultSubCode.IGNORE_EQUIVOCATION_DETECTED)) + .orElse(false); + } + public boolean isReject() { return this.validationResultCode.equals(ValidationResultCode.REJECT); } diff --git a/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/block/BlockManagerTest.java b/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/block/BlockManagerTest.java index d8ae2ac94d3..b83b6063ede 100644 --- a/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/block/BlockManagerTest.java +++ b/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/block/BlockManagerTest.java @@ -45,6 +45,8 @@ import static tech.pegasys.teku.statetransition.block.BlockImportPerformance.TRANSACTION_COMMITTED_EVENT_LABEL; import static tech.pegasys.teku.statetransition.block.BlockImportPerformance.TRANSACTION_PREPARED_EVENT_LABEL; import static tech.pegasys.teku.statetransition.validation.BlockBroadcastValidator.BroadcastValidationResult.SUCCESS; +import static tech.pegasys.teku.statetransition.validation.ValidationResultCode.ValidationResultSubCode.IGNORE_ALREADY_SEEN; +import static tech.pegasys.teku.statetransition.validation.ValidationResultCode.ValidationResultSubCode.IGNORE_EQUIVOCATION_DETECTED; import java.util.ArrayList; import java.util.Collections; @@ -65,6 +67,7 @@ import tech.pegasys.teku.infrastructure.async.ExceptionThrowingFutureSupplier; import tech.pegasys.teku.infrastructure.async.SafeFuture; import tech.pegasys.teku.infrastructure.async.SafeFutureAssert; +import tech.pegasys.teku.infrastructure.async.Waiter; import tech.pegasys.teku.infrastructure.async.eventthread.InlineEventThread; import tech.pegasys.teku.infrastructure.collections.LimitedMap; import tech.pegasys.teku.infrastructure.logging.EventLogger; @@ -547,6 +550,66 @@ public void onProposedBlock_futureBlock() { verifyNoInteractions(blobSidecarManager); } + @Test + public void onProposedBlock_futureBlock_shouldRerunGossipValidationAndImportOnRetry() { + incrementSlot(); + final UInt64 nextSlot = currentSlot.plus(UInt64.ONE); + final SignedBeaconBlock futureBlock = + localChain.chainBuilder().generateBlockAtSlot(nextSlot).getBlock(); + + when(blockValidator.validateGossip(eq(futureBlock))) + .thenReturn(SafeFuture.completedFuture(InternalValidationResult.SAVE_FOR_FUTURE)) + .thenReturn( + SafeFuture.completedFuture( + InternalValidationResult.ignore( + IGNORE_ALREADY_SEEN, "retry validation found already seen"))); + + assertThatSafeFuture(blockManager.validateAndImportBlock(futureBlock, Optional.empty())) + .isCompletedWithValue(InternalValidationResult.SAVE_FOR_FUTURE); + Waiter.waitFor(() -> assertThat(futureBlocks.size()).isEqualTo(1)); + assertThat(futureBlocks.contains(futureBlock)).isTrue(); + assertThat(invalidBlockRoots).isEmpty(); + + incrementSlot(); + + Waiter.waitFor(() -> verify(blockValidator, times(2)).validateGossip(eq(futureBlock))); + Waiter.waitFor( + () -> + verify(receivedBlockEventsChannelPublisher, times(1)) + .onBlockImported(futureBlock, false)); + Waiter.waitFor(() -> assertThat(invalidBlockRoots).isEmpty()); + Waiter.waitFor(() -> assertThat(futureBlocks.size()).isEqualTo(0)); + verify(blockEventsListenerRouter).onBlockImported(futureBlock); + } + + @Test + public void onProposedBlock_futureBlock_shouldCleanupEquivocationOnRetry() { + incrementSlot(); + final UInt64 nextSlot = currentSlot.plus(UInt64.ONE); + final SignedBeaconBlock futureBlock = + localChain.chainBuilder().generateBlockAtSlot(nextSlot).getBlock(); + + when(blockValidator.validateGossip(eq(futureBlock))) + .thenReturn(SafeFuture.completedFuture(InternalValidationResult.SAVE_FOR_FUTURE)) + .thenReturn( + SafeFuture.completedFuture( + InternalValidationResult.ignore( + IGNORE_EQUIVOCATION_DETECTED, "retry validation found equivocation"))); + + assertThatSafeFuture(blockManager.validateAndImportBlock(futureBlock, Optional.empty())) + .isCompletedWithValue(InternalValidationResult.SAVE_FOR_FUTURE); + Waiter.waitFor(() -> assertThat(futureBlocks.size()).isEqualTo(1)); + assertThat(futureBlocks.contains(futureBlock)).isTrue(); + assertThat(invalidBlockRoots).isEmpty(); + + incrementSlot(); + + Waiter.waitFor(() -> verify(blockValidator, times(2)).validateGossip(eq(futureBlock))); + Waiter.waitFor(() -> assertThat(invalidBlockRoots).isEmpty()); + Waiter.waitFor(() -> assertThat(futureBlocks.size()).isEqualTo(0)); + verify(blockEventsListenerRouter).removeAllForBlock(futureBlock.getSlotAndBlockRoot()); + } + @Test public void onBlockImported_withPendingBlocks() { final int blockCount = 3; diff --git a/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/util/FutureItemsTest.java b/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/util/FutureItemsTest.java index 34bd65a749c..690fe51a914 100644 --- a/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/util/FutureItemsTest.java +++ b/ethereum/statetransition/src/test/java/tech/pegasys/teku/statetransition/util/FutureItemsTest.java @@ -39,7 +39,18 @@ public void add_success() { final UInt64 itemSlot = currentSlot.plus(FutureItems.DEFAULT_FUTURE_SLOT_TOLERANCE); final Item item = new Item(itemSlot); - futureItems.add(item); + assertThat(futureItems.add(item)).isTrue(); + assertThat(futureItems.size()).isEqualTo(1); + assertThat(futureItems.contains(item)).isTrue(); + } + + @Test + public void add_duplicate_isStillAccepted() { + final UInt64 itemSlot = currentSlot.plus(FutureItems.DEFAULT_FUTURE_SLOT_TOLERANCE); + final Item item = new Item(itemSlot); + + assertThat(futureItems.add(item)).isTrue(); + assertThat(futureItems.add(item)).isTrue(); assertThat(futureItems.size()).isEqualTo(1); assertThat(futureItems.contains(item)).isTrue(); } @@ -50,8 +61,8 @@ public void add_ignored() { final Item itemA = new Item(itemSlot); final Item itemB = new Item(itemSlot.plus(10)); - futureItems.add(itemA); - futureItems.add(itemB); + assertThat(futureItems.add(itemA)).isFalse(); + assertThat(futureItems.add(itemB)).isFalse(); assertThat(futureItems.size()).isEqualTo(0); assertThat(futureItems.contains(itemA)).isFalse(); assertThat(futureItems.contains(itemB)).isFalse(); @@ -62,7 +73,7 @@ public void prune_nothingToPrune() { final UInt64 itemSlot = currentSlot.plus(FutureItems.DEFAULT_FUTURE_SLOT_TOLERANCE); final Item item = new Item(itemSlot); - futureItems.add(item); + assertThat(futureItems.add(item)).isTrue(); final UInt64 priorSlot = item.getSlot().minus(UInt64.ONE); final List pruned = futureItems.prune(priorSlot); @@ -77,7 +88,7 @@ public void prune_itemAtSlot() { final UInt64 itemSlot = currentSlot.plus(FutureItems.DEFAULT_FUTURE_SLOT_TOLERANCE); final Item item = new Item(itemSlot); - futureItems.add(item); + assertThat(futureItems.add(item)).isTrue(); final List pruned = futureItems.prune(item.getSlot()); assertThat(pruned).containsExactly(item); @@ -90,7 +101,7 @@ public void metrics_shouldIncreaseAndDecrease() { final UInt64 itemSlot = currentSlot.plus(FutureItems.DEFAULT_FUTURE_SLOT_TOLERANCE); final Item item = new Item(itemSlot); - futureItems.add(item); + assertThat(futureItems.add(item)).isTrue(); verify(gauge).set(1L, "items"); final UInt64 pruneSlot = item.getSlot().plus(UInt64.ONE); @@ -104,7 +115,7 @@ public void prune_itemPriorToSlot() { final UInt64 itemSlot = currentSlot.plus(FutureItems.DEFAULT_FUTURE_SLOT_TOLERANCE); final Item item = new Item(itemSlot); - futureItems.add(item); + assertThat(futureItems.add(item)).isTrue(); final List pruned = futureItems.prune(item.getSlot().plus(UInt64.ONE)); assertThat(pruned).containsExactly(item);