Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 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 @@ -173,4 +173,7 @@ public interface ISyncConfig : IConfig

[ConfigItem(Description = "_Technical._ Max tx in forward sync buffer.", DefaultValue = "200000", HiddenFromDocs = true)]
int MaxTxInForwardSyncBuffer { get; set; }

[ConfigItem(Description = "_Technical._ Max tx waiting for processing before stopping.", DefaultValue = "400000", HiddenFromDocs = true)]
int MaxTxInProcessingQueue { get; set; }
}
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,8 @@ public string? PivotHash

public ulong FastHeadersMemoryBudget { get; set; } = (ulong)128.MB();
public bool EnableSnapSyncStorageRangeSplit { get; set; } = false;
public int MaxTxInForwardSyncBuffer { get; set; } = 20000;
public int MaxTxInForwardSyncBuffer { get; set; } = 200000;
public int MaxTxInProcessingQueue { get; set; } = 400000;

public override string ToString()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
using Nethermind.Blockchain;
using Nethermind.Blockchain.Receipts;
using Nethermind.Blockchain.Synchronization;
using Nethermind.Consensus.Processing;
using Nethermind.Consensus.Validators;
using Nethermind.Core;
using Nethermind.Core.Specs;
Expand All @@ -18,40 +19,46 @@

namespace Nethermind.Merge.Plugin.Synchronization
{
public class MergeBlockDownloader : BlockDownloader
public class MergeBlockDownloader(
IBeaconPivot beaconPivot,
IBlockTree blockTree,
IBlockValidator blockValidator,
ISyncReport syncReport,
IReceiptStorage receiptStorage,
ISpecProvider specProvider,
IBetterPeerStrategy betterPeerStrategy,
IFullStateFinder fullStateFinder,
IForwardHeaderProvider forwardHeaderProvider,
ISyncPeerPool syncPeerPool,
IReceiptsRecovery receiptsRecovery,
IBlockProcessingQueue blockProcessingQueue,
ISyncConfig syncConfig,
ILogManager logManager)
: BlockDownloader(
blockTree,
blockValidator,
syncReport,
receiptStorage,
specProvider,
betterPeerStrategy,
fullStateFinder,
forwardHeaderProvider,
syncPeerPool,
receiptsRecovery,
blockProcessingQueue,
syncConfig,
logManager)
{
private readonly IBeaconPivot _beaconPivot;
private readonly IBlockTree _blockTree;
private readonly ILogger _logger;

public MergeBlockDownloader(
IBeaconPivot beaconPivot,
IBlockTree? blockTree,
IBlockValidator? blockValidator,
ISyncReport? syncReport,
IReceiptStorage? receiptStorage,
ISpecProvider specProvider,
IBetterPeerStrategy betterPeerStrategy,
IFullStateFinder fullStateFinder,
IForwardHeaderProvider forwardHeaderProvider,
ISyncPeerPool syncPeerPool,
ISyncConfig syncConfig,
ILogManager logManager)
: base(blockTree, blockValidator, syncReport, receiptStorage,
specProvider, betterPeerStrategy, fullStateFinder, forwardHeaderProvider, syncPeerPool, syncConfig, logManager)
{
_blockTree = blockTree ?? throw new ArgumentNullException(nameof(blockTree));
_beaconPivot = beaconPivot;
_logger = logManager.GetClassLogger();
}
private readonly IBlockTree _blockTree = blockTree;
private readonly ILogger _logger = logManager.GetClassLogger();

protected override BlockTreeSuggestOptions GetSuggestOption(bool shouldProcess, Block currentBlock)
{
BlockTreeSuggestOptions suggestOptions =
shouldProcess ? BlockTreeSuggestOptions.ShouldProcess : BlockTreeSuggestOptions.None;

bool isKnownBeaconBlock = _blockTree.IsKnownBeaconBlock(currentBlock.Number, currentBlock.GetOrCalculateHash());
if (_logger.IsTrace) _logger.Trace($"Current block {currentBlock}, BeaconPivot: {_beaconPivot.PivotNumber}, IsKnownBeaconBlock: {isKnownBeaconBlock}");
if (_logger.IsTrace) _logger.Trace($"Current block {currentBlock}, BeaconPivot: {beaconPivot.PivotNumber}, IsKnownBeaconBlock: {isKnownBeaconBlock}");

if (isKnownBeaconBlock)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
using Nethermind.Blockchain;
using Nethermind.Blockchain.Receipts;
using Nethermind.Blockchain.Synchronization;
using Nethermind.Consensus.Processing;
using Nethermind.Consensus.Validators;
using Nethermind.Core;
using Nethermind.Core.Collections;
Expand Down Expand Up @@ -47,10 +48,12 @@ public class BlockDownloader : IForwardSyncController
private readonly IFullStateFinder _fullStateFinder;
private readonly IForwardHeaderProvider _forwardHeaderProvider;
private readonly ISyncPeerPool _syncPeerPool;
private readonly IBlockProcessingQueue _processingQueue;
private readonly ILogger _logger;

// Estimated maximum tx in buffer used to estimate memory limit. Each tx is on average about 1KB.
private readonly int _maxTxInBuffer;
private readonly int _maxTxInInProcessingQueue;
private const int MinEstimateTxPerBlock = 10;

// Header lookup need to be limited, because `IForwardHeaderProvider.GetBlockHeaders` can be slow.
Expand All @@ -72,31 +75,34 @@ public class BlockDownloader : IForwardSyncController
private SemaphoreSlim _requestLock = new(1);

public BlockDownloader(
IBlockTree? blockTree,
IBlockValidator? blockValidator,
ISyncReport? syncReport,
IReceiptStorage? receiptStorage,
ISpecProvider? specProvider,
IBlockTree blockTree,
IBlockValidator blockValidator,
ISyncReport syncReport,
IReceiptStorage receiptStorage,
ISpecProvider specProvider,
IBetterPeerStrategy betterPeerStrategy,
IFullStateFinder fullStateFinder,
IForwardHeaderProvider forwardHeaderProvider,
ISyncPeerPool syncPeerPool,
IReceiptsRecovery receiptsRecovery,
IBlockProcessingQueue processingQueue,
ISyncConfig syncConfig,
ILogManager? logManager)
ILogManager logManager)
{
_blockTree = blockTree ?? throw new ArgumentNullException(nameof(blockTree));
_blockValidator = blockValidator ?? throw new ArgumentNullException(nameof(blockValidator));
_syncReport = syncReport ?? throw new ArgumentNullException(nameof(syncReport));
_receiptStorage = receiptStorage ?? throw new ArgumentNullException(nameof(receiptStorage));
_specProvider = specProvider ?? throw new ArgumentNullException(nameof(specProvider));
_betterPeerStrategy = betterPeerStrategy ?? throw new ArgumentNullException(nameof(betterPeerStrategy));
_fullStateFinder = fullStateFinder ?? throw new ArgumentNullException(nameof(fullStateFinder));
_blockTree = blockTree;
_blockValidator = blockValidator;
_syncReport = syncReport;
_receiptStorage = receiptStorage;
_specProvider = specProvider;
_betterPeerStrategy = betterPeerStrategy;
_fullStateFinder = fullStateFinder;
_forwardHeaderProvider = forwardHeaderProvider;
_syncPeerPool = syncPeerPool;
_logger = logManager?.GetClassLogger() ?? throw new ArgumentNullException(nameof(logManager));
_logger = logManager.GetClassLogger();
_maxTxInBuffer = syncConfig.MaxTxInForwardSyncBuffer;

_receiptsRecovery = new ReceiptsRecovery(new EthereumEcdsa(_specProvider.ChainId), _specProvider);
_maxTxInInProcessingQueue = syncConfig.MaxTxInProcessingQueue;
_receiptsRecovery = receiptsRecovery;
_processingQueue = processingQueue;
_blockTree.NewHeadBlock += BlockTreeOnNewHeadBlock;
}

Expand Down Expand Up @@ -145,6 +151,12 @@ private void BlockTreeOnNewHeadBlock(object? sender, BlockEventArgs e)

while (true)
{
if (_processingQueue.Count > _maxTxInInProcessingQueue / _estimateTxPerBlock)
{
if (_logger.IsTrace) _logger.Trace("Processing queue full");
return null;
}

using IOwnedReadOnlyList<BlockHeader?>? headers = await _forwardHeaderProvider.GetBlockHeaders(fastSyncLag, HeaderLookupSize + 1, cancellation);
if (cancellation.IsCancellationRequested) return null; // check before every heavy operation
if (headers is null || headers.Count <= 1) return null;
Expand Down