Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,11 @@
package tech.pegasys.teku.statetransition.block;

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 @@ -46,6 +48,7 @@
import tech.pegasys.teku.statetransition.validation.BlockBroadcastValidator;
import tech.pegasys.teku.statetransition.validation.BlockValidator;
import tech.pegasys.teku.statetransition.validation.InternalValidationResult;
import tech.pegasys.teku.statetransition.validation.ValidationResultCode.ValidationResultSubCode;
import tech.pegasys.teku.storage.client.RecentChainData;

public class BlockManager extends Service
Expand All @@ -64,7 +67,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 +98,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 +127,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 +175,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 +191,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 +251,48 @@ 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) {
// This future-block retry path handles a block that gossip accepts but import still defers.
//
// 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, gossip revalidation may return IGNORE_ALREADY_SEEN because the
// block was already marked as seen.
// 4. Retry the import without gossip validation so the block is not lost.
if (futureBlock.needsGossipValidation()) {
validateAndImportBlock(futureBlock.block(), Optional.empty())
.thenAccept(
result -> {
if (result.isIgnoreAlreadySeen()) {
importBlockIgnoringResult(futureBlock.block());
}
})
.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 +319,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 +334,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 +363,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 +548,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 @@ -45,6 +45,7 @@
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_EQUIVOCATION_DETECTED;

import java.util.ArrayList;
import java.util.Collections;
Expand All @@ -65,6 +66,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 +549,60 @@ public void onProposedBlock_futureBlock() {
verifyNoInteractions(blobSidecarManager);
}

@Test
public void onProposedBlock_futureBlock_shouldRerunGossipValidationOnRetry() {
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.reject("retry validation failed")));

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 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