Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -64,7 +68,7 @@ public class BlockManager extends Service
private final TimeProvider timeProvider;
private final EventLogger eventLogger;

private final FutureItems<SignedBeaconBlock> 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.
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -124,7 +128,7 @@ public SafeFuture<BlockImportAndBroadcastValidationResults> importBlock(
blockValidator.initiateBroadcastValidation(block, broadcastValidationLevel);

final SafeFuture<BlockImportResult> 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
Expand Down Expand Up @@ -172,7 +176,8 @@ public SafeFuture<InternalValidationResult> 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
Expand All @@ -187,8 +192,7 @@ public SafeFuture<InternalValidationResult> 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) {
Expand Down Expand Up @@ -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);
Comment thread
cursor[bot] marked this conversation as resolved.
} else {
importBlockIgnoringResult(futureBlock.block());
}
}

private SafeFuture<BlockImportResult> doImportBlock(
final SignedBeaconBlock block,
final Optional<BlockImportPerformance> blockImportPerformance,
final BlockBroadcastValidator blockBroadcastValidator,
final Optional<RemoteOrigin> origin) {
final Optional<RemoteOrigin> 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)));
}
Expand All @@ -288,7 +323,7 @@ private Optional<SafeFuture<BlockImportResult>> handleInvalidBlock(
}

private Optional<SafeFuture<BlockImportResult>> 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)));
}
Expand All @@ -303,7 +338,8 @@ private SafeFuture<BlockImportResult> handleBlockImport(
final SignedBeaconBlock block,
final Optional<BlockImportPerformance> blockImportPerformance,
final BlockBroadcastValidator blockBroadcastValidator,
final Optional<RemoteOrigin> origin) {
final Optional<RemoteOrigin> origin,
final boolean needsGossipValidationOnRetry) {
blockEventsListener.onNewBlock(block, origin);
preImportBlockSubscribers.deliver(l -> l.onNewBlock(block, origin));

Expand Down Expand Up @@ -331,7 +367,8 @@ private SafeFuture<BlockImportResult> 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",
Expand Down Expand Up @@ -515,4 +552,35 @@ public interface PreImportBlockListener {
public interface RequiredParentExecutionPayloadSubscriber {
void onRequiredParentExecutionPayload(ParentExecutionPayloadDependency dependency);
}

private static final class FutureBlockTracker {
private final FutureItems<SignedBeaconBlock> futureBlocks;
private final Set<Bytes32> gossipRetryRoots = new HashSet<>();

private FutureBlockTracker(final FutureItems<SignedBeaconBlock> 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<QueuedFutureBlock> 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) {}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
Expand All @@ -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();
Expand All @@ -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<Item> pruned = futureItems.prune(priorSlot);
Expand All @@ -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<Item> pruned = futureItems.prune(item.getSlot());
assertThat(pruned).containsExactly(item);
Expand All @@ -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);
Expand All @@ -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<Item> pruned = futureItems.prune(item.getSlot().plus(UInt64.ONE));
assertThat(pruned).containsExactly(item);
Expand Down
Loading