diff --git a/cl/beacon/handler/block_production.go b/cl/beacon/handler/block_production.go index 3d156ef8dbc..8593177aac2 100644 --- a/cl/beacon/handler/block_production.go +++ b/cl/beacon/handler/block_production.go @@ -45,6 +45,7 @@ import ( peerdasutils "github.com/erigontech/erigon/cl/das/utils" "github.com/erigontech/erigon/cl/gossip" "github.com/erigontech/erigon/cl/persistence/beacon_indicies" + "github.com/erigontech/erigon/cl/persistence/blob_storage" "github.com/erigontech/erigon/cl/phase1/core/state" "github.com/erigontech/erigon/cl/phase1/forkchoice" "github.com/erigontech/erigon/cl/phase1/network/subnets" @@ -2079,7 +2080,12 @@ func (a *ApiHandler) storeBlockAndBlobs( if err != nil { return err } - // TODO: write column sidecars if needed + commitments := block.GetBlobKzgCommitments() + if block.Version() >= clparams.FuluVersion && commitments != nil && commitments.Len() > 0 { + if err := storeProducedDataColumns(ctx, a.columnStorage, a.beaconChainCfg, blockRoot, columnSidecars); err != nil { + return err + } + } if block.Version() < clparams.FuluVersion { if err := a.blobStoage.WriteBlobSidecars(ctx, blockRoot, sidecars); err != nil { @@ -2141,6 +2147,38 @@ func (a *ApiHandler) storeBlockAndBlobs( return nil } +func storeProducedDataColumns( + ctx context.Context, + storage blob_storage.DataColumnStorage, + cfg *clparams.BeaconChainConfig, + blockRoot common.Hash, + columns []*cltypes.DataColumnSidecar, +) error { + expected := int(cfg.NumberOfColumns) + if len(columns) != expected { + return fmt.Errorf("expected %d data columns, got %d", expected, len(columns)) + } + if storage == nil { + return errors.New("data column storage is not configured") + } + seen := make([]bool, expected) + for _, column := range columns { + if column == nil || column.Index >= cfg.NumberOfColumns { + return errors.New("produced data column has an invalid index") + } + if seen[column.Index] { + return fmt.Errorf("produced data column %d is duplicated", column.Index) + } + seen[column.Index] = true + } + for _, column := range columns { + if err := storage.WriteColumnSidecars(ctx, blockRoot, int64(column.Index), column); err != nil { + return fmt.Errorf("store data column %d: %w", column.Index, err) + } + } + return nil +} + func (a *ApiHandler) selectedHeadState(auxiliaryRoot common.Hash) (common.Hash, uint64, *state.CachingBeaconState, error) { auxiliaryState, err := a.forkchoiceStore.GetStateAtBlockRoot(auxiliaryRoot, false) if err != nil { diff --git a/cl/beacon/handler/block_production_test.go b/cl/beacon/handler/block_production_test.go index fddce8a9234..6bbaef996ea 100644 --- a/cl/beacon/handler/block_production_test.go +++ b/cl/beacon/handler/block_production_test.go @@ -36,6 +36,7 @@ import ( "github.com/erigontech/erigon/cl/clparams" "github.com/erigontech/erigon/cl/cltypes" "github.com/erigontech/erigon/cl/cltypes/solid" + blobstoragemock "github.com/erigontech/erigon/cl/persistence/blob_storage/mock_services" "github.com/erigontech/erigon/cl/phase1/core/state" "github.com/erigontech/erigon/cl/phase1/execution_client" "github.com/erigontech/erigon/common" @@ -52,6 +53,29 @@ import ( "github.com/erigontech/erigon/node/gointerfaces/typesproto" ) +func TestStoreProducedDataColumns(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + cfg.NumberOfColumns = 2 + blockRoot := common.Hash{1} + columns := []*cltypes.DataColumnSidecar{{Index: 0}, {Index: 1}} + + ctrl := gomock.NewController(t) + storage := blobstoragemock.NewMockDataColumnStorage(ctrl) + for _, column := range columns { + storage.EXPECT().WriteColumnSidecars(gomock.Any(), blockRoot, int64(column.Index), column).Return(nil) + } + require.NoError(t, storeProducedDataColumns(t.Context(), storage, &cfg, blockRoot, columns)) +} + +func TestStoreProducedDataColumnsRequiresEveryColumn(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + cfg.NumberOfColumns = 2 + storage := blobstoragemock.NewMockDataColumnStorage(gomock.NewController(t)) + + err := storeProducedDataColumns(t.Context(), storage, &cfg, common.Hash{}, []*cltypes.DataColumnSidecar{{Index: 0}}) + require.Error(t, err) +} + func TestBlockBuilderWindowPreGloas(t *testing.T) { cfg := &clparams.BeaconChainConfig{ SecondsPerSlot: 12, diff --git a/cl/cltypes/beacon_block_interface.go b/cl/cltypes/beacon_block_interface.go index fa7f63880b3..3a42c2656de 100644 --- a/cl/cltypes/beacon_block_interface.go +++ b/cl/cltypes/beacon_block_interface.go @@ -6,10 +6,7 @@ import ( "github.com/erigontech/erigon/common" ) -// ColumnSyncableSignedBlock is implemented by both SignedBeaconBlock and SignedBlindedBeaconBlock -// for PeerDAS column synchronization operations. -// [New in Gloas:EIP7732] This interface allows PeerDAS to work with both block types -// without needing to call Blinded() which fails for GLOAS blocks. +// ColumnSyncableSignedBlock provides the block metadata needed to request PeerDAS columns. type ColumnSyncableSignedBlock interface { Version() clparams.StateVersion GetSlot() uint64 diff --git a/cl/das/availability.go b/cl/das/availability.go new file mode 100644 index 00000000000..d09a058f5f9 --- /dev/null +++ b/cl/das/availability.go @@ -0,0 +1,32 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package das + +import "github.com/erigontech/erigon/cl/clparams" + +// IsDataAvailabilityRequired reports whether a Fulu block is inside the protocol's column request window. +func IsDataAvailabilityRequired(cfg *clparams.BeaconChainConfig, currentSlot, blockSlot uint64, version clparams.StateVersion) bool { + if version != clparams.FuluVersion { + return false + } + if cfg.SlotsPerEpoch == 0 { + return true + } + currentEpoch := currentSlot / cfg.SlotsPerEpoch + blockEpoch := blockSlot / cfg.SlotsPerEpoch + return blockEpoch >= currentEpoch || currentEpoch-blockEpoch <= cfg.MinEpochsForDataColumnSidecarsRequests +} diff --git a/cl/das/availability_test.go b/cl/das/availability_test.go new file mode 100644 index 00000000000..66b6bf28cee --- /dev/null +++ b/cl/das/availability_test.go @@ -0,0 +1,51 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package das + +import ( + "testing" + + "github.com/stretchr/testify/require" + + "github.com/erigontech/erigon/cl/clparams" +) + +func TestIsDataAvailabilityRequired(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + lastSlotInRetentionWindow := (cfg.MinEpochsForDataColumnSidecarsRequests+1)*cfg.SlotsPerEpoch - 1 + firstSlotAfterRetentionWindow := lastSlotInRetentionWindow + 1 + tests := []struct { + name string + version clparams.StateVersion + currentSlot uint64 + blockSlot uint64 + want bool + }{ + {name: "current Fulu block", version: clparams.FuluVersion, currentSlot: 100, blockSlot: 100, want: true}, + {name: "future Fulu block", version: clparams.FuluVersion, currentSlot: 100, blockSlot: 101, want: true}, + {name: "last slot in retention window", version: clparams.FuluVersion, currentSlot: lastSlotInRetentionWindow, blockSlot: 0, want: true}, + {name: "first slot after retention window", version: clparams.FuluVersion, currentSlot: firstSlotAfterRetentionWindow, blockSlot: 0}, + {name: "Electra", version: clparams.ElectraVersion, currentSlot: 100, blockSlot: 100}, + {name: "Gloas envelope owns availability", version: clparams.GloasVersion, currentSlot: 100, blockSlot: 100}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + require.Equal(t, tt.want, IsDataAvailabilityRequired(&cfg, tt.currentSlot, tt.blockSlot, tt.version)) + }) + } +} diff --git a/cl/das/peer_das.go b/cl/das/peer_das.go index 6e625e43f2e..a4d6698e7c8 100644 --- a/cl/das/peer_das.go +++ b/cl/das/peer_das.go @@ -3,6 +3,7 @@ package das import ( "context" "errors" + "maps" "math" "sync" "time" @@ -66,6 +67,39 @@ type PeerDas interface { var numOfBlobRecoveryWorkers = 8 +const ( + maxDeferredColumnSyncBlocks = 64 + deferredColumnSyncTTL = 30 * time.Minute +) + +type deferredColumnSyncJob struct { + block cltypes.ColumnSyncableSignedBlock + addedAt time.Time +} + +type deferredColumnSyncBlock struct { + slot uint64 + root common.Hash + version clparams.StateVersion + commitments *solid.ListSSZ[*cltypes.KZGCommitment] +} + +func (b *deferredColumnSyncBlock) Version() clparams.StateVersion { + return b.version +} + +func (b *deferredColumnSyncBlock) GetSlot() uint64 { + return b.slot +} + +func (b *deferredColumnSyncBlock) BlockHashSSZ() ([32]byte, error) { + return b.root, nil +} + +func (b *deferredColumnSyncBlock) GetBlobKzgCommitments() *solid.ListSSZ[*cltypes.KZGCommitment] { + return b.commitments +} + type peerdas struct { state *peerdasstate.PeerDasState nodeID enode.ID @@ -79,9 +113,10 @@ type peerdas struct { gossipManager gossipmgr.Gossip recoverBlobsQueue chan recoverBlobsRequest - recoveringMutex sync.Mutex - isRecovering map[common.Hash]bool - blocksToCheckSync sync.Map // blockRoot -> ColumnSyncableSignedBlock (SignedBeaconBlock or SignedBlindedBeaconBlock) + recoveringMutex sync.Mutex + isRecovering map[common.Hash]bool + blocksToCheckSync map[common.Hash]deferredColumnSyncJob + blocksToCheckSyncMu sync.Mutex // [New in Gloas:EIP7732] For fetching blocks to get kzg_commitments forkChoice BlockGetter @@ -122,7 +157,7 @@ func NewPeerDas( recoveringMutex: sync.Mutex{}, isRecovering: make(map[common.Hash]bool), - blocksToCheckSync: sync.Map{}, + blocksToCheckSync: make(map[common.Hash]deferredColumnSyncJob), blockReader: blockReader, indiciesDB: indiciesDB, @@ -647,8 +682,11 @@ var allColumns = func() map[cltypes.CustodyIndex]bool { return columns }() -// DownloadMissingColumns downloads the missing columns for the given blocks but not recover the blobs +// DownloadOnlyCustodyColumns downloads custody columns without reconstructing blobs. func (d *peerdas) DownloadOnlyCustodyColumns(ctx context.Context, blocks []cltypes.ColumnSyncableSignedBlock) error { + if d.rpc == nil { + return errors.New("peer DAS RPC is not configured") + } custodyColumns, err := d.state.GetMyCustodyColumns() if err != nil { return err @@ -1085,11 +1123,13 @@ func (d *downloadRequest) requestData() *solid.ListSSZ[*cltypes.DataColumnsByRoo } func (d *peerdas) SyncColumnDataLater(block *cltypes.SignedBeaconBlock) error { - if block.Version() < clparams.FuluVersion { + blockVersion := block.Version() + if d.beaconConfig != nil && d.beaconConfig.SlotsPerEpoch != 0 { + blockVersion = d.beaconConfig.GetCurrentStateVersion(block.Block.Slot / d.beaconConfig.SlotsPerEpoch) + } + if blockVersion < clparams.FuluVersion { return nil } - // [Modified in Gloas:EIP7732] Use GetBlobKzgCommitments() which is version-aware - // For GLOAS, commitments are in SignedExecutionPayloadBid.Message kzgCommitments := block.GetBlobKzgCommitments() if kzgCommitments == nil || kzgCommitments.Len() == 0 { return nil @@ -1098,12 +1138,58 @@ func (d *peerdas) SyncColumnDataLater(block *cltypes.SignedBeaconBlock) error { if err != nil { return err } - // [Modified in Gloas:EIP7732] Store SignedBeaconBlock directly via ColumnSyncableSignedBlock interface - // instead of calling Blinded() which fails for GLOAS blocks - d.blocksToCheckSync.Store(common.Hash(blockRoot), block) + queuedBlock := &deferredColumnSyncBlock{ + slot: block.Block.Slot, + root: common.Hash(blockRoot), + version: blockVersion, + commitments: kzgCommitments.ShallowCopy(), + } + d.storeDeferredColumnSyncJob(common.Hash(blockRoot), deferredColumnSyncJob{ + block: queuedBlock, + addedAt: time.Now(), + }) return nil } +func (d *peerdas) storeDeferredColumnSyncJob(root common.Hash, job deferredColumnSyncJob) { + d.blocksToCheckSyncMu.Lock() + defer d.blocksToCheckSyncMu.Unlock() + if d.blocksToCheckSync == nil { + d.blocksToCheckSync = make(map[common.Hash]deferredColumnSyncJob) + } + if _, exists := d.blocksToCheckSync[root]; exists { + return + } + if len(d.blocksToCheckSync) >= maxDeferredColumnSyncBlocks { + var oldestRoot common.Hash + var oldestTime time.Time + for candidateRoot, candidate := range d.blocksToCheckSync { + if oldestTime.IsZero() || candidate.addedAt.Before(oldestTime) { + oldestRoot = candidateRoot + oldestTime = candidate.addedAt + } + } + if !oldestTime.IsZero() { + delete(d.blocksToCheckSync, oldestRoot) + } + } + d.blocksToCheckSync[root] = job +} + +func (d *peerdas) deleteDeferredColumnSyncJob(root common.Hash) { + d.blocksToCheckSyncMu.Lock() + defer d.blocksToCheckSyncMu.Unlock() + delete(d.blocksToCheckSync, root) +} + +func (d *peerdas) deferredColumnSyncJobs() map[common.Hash]deferredColumnSyncJob { + d.blocksToCheckSyncMu.Lock() + defer d.blocksToCheckSyncMu.Unlock() + jobs := make(map[common.Hash]deferredColumnSyncJob, len(d.blocksToCheckSync)) + maps.Copy(jobs, d.blocksToCheckSync) + return jobs +} + func (d *peerdas) syncColumnDataWorker(ctx context.Context) { ticker := time.NewTicker(time.Minute) defer ticker.Stop() @@ -1112,60 +1198,84 @@ func (d *peerdas) syncColumnDataWorker(ctx context.Context) { case <-ctx.Done(): return case <-ticker.C: - // check peers count - if d.rpc != nil { - if peersCount, err := d.rpc.Peers(); err != nil { - log.Warn("failed to get peers count", "err", err) - continue - } else if peersCount == 0 { - log.Info("[syncColumnDataWorker] no peers available, skipping sync") - continue - } - } + d.syncDeferredColumnData(ctx) + } + } +} - // [Modified in Gloas:EIP7732] Use ColumnSyncableSignedBlock interface - blocks := []cltypes.ColumnSyncableSignedBlock{} - roots := []common.Hash{} - d.blocksToCheckSync.Range(func(key, value any) bool { - root := key.(common.Hash) - block := value.(cltypes.ColumnSyncableSignedBlock) - curSlot := d.ethClock.GetCurrentSlot() - if curSlot-block.GetSlot() < 5 { // wait slow data from peers - // skip blocks that are too close to the current slot - return true - } - available, err := d.IsDataAvailable(block.GetSlot(), root) - switch { - case err != nil: - log.Warn("failed to check if data is available", "err", err) - case available: - log.Trace("[syncColumnDataWorker] column data is already available, removing from sync queue", "slot", block.GetSlot(), "blockRoot", root) - d.blocksToCheckSync.Delete(root) - default: - blocks = append(blocks, block) - roots = append(roots, root) - } - return true - }) - if len(blocks) == 0 { - continue - } - log.Debug("[syncColumnDataWorker] syncing column data", "blocks_count", len(blocks)) - if d.IsArchivedMode() { - if err := d.DownloadColumnsAndRecoverBlobs(ctx, blocks); err != nil { - log.Warn("failed to download columns and recover blobs", "err", err) - continue - } - } else { - if err := d.DownloadOnlyCustodyColumns(ctx, blocks); err != nil { - log.Warn("failed to download only custody columns", "err", err) - continue - } - } - for i, root := range roots { - d.blocksToCheckSync.Delete(root) - log.Debug("[syncColumnDataWorker] column data is synced, removing from sync queue", "slot", blocks[i].GetSlot(), "blockRoot", root) - } +func (d *peerdas) syncDeferredColumnData(ctx context.Context) { + jobs := d.deferredColumnSyncJobs() + now := time.Now() + for root, job := range jobs { + if now.Sub(job.addedAt) > deferredColumnSyncTTL { + d.deleteDeferredColumnSyncJob(root) + delete(jobs, root) + } + } + if d.rpc == nil { + log.Debug("[syncColumnDataWorker] PeerDAS RPC is not configured") + return + } + if peersCount, err := d.rpc.Peers(); err != nil { + log.Warn("failed to get peers count", "err", err) + return + } else if peersCount == 0 { + log.Info("[syncColumnDataWorker] no peers available, skipping sync") + return + } + + blocks := make([]cltypes.ColumnSyncableSignedBlock, 0, len(jobs)) + roots := make([]common.Hash, 0, len(jobs)) + curSlot := d.ethClock.GetCurrentSlot() + for root, job := range jobs { + block := job.block + if block.GetSlot() > curSlot || curSlot-block.GetSlot() < 5 { + continue + } + available, err := d.IsDataAvailable(block.GetSlot(), root) + switch { + case err != nil: + log.Warn("failed to check if data is available", "err", err) + case available: + log.Trace("[syncColumnDataWorker] column data is already available, removing from sync queue", "slot", block.GetSlot(), "blockRoot", root) + d.deleteDeferredColumnSyncJob(root) + default: + blocks = append(blocks, block) + roots = append(roots, root) + } + } + if len(blocks) == 0 { + return + } + log.Debug("[syncColumnDataWorker] syncing column data", "blocks_count", len(blocks)) + downloadTimeout := time.Duration(d.beaconConfig.SecondsPerSlot) * time.Second + if downloadTimeout <= 0 { + downloadTimeout = time.Second + } + downloadCtx, cancel := context.WithTimeout(ctx, downloadTimeout) + if d.IsArchivedMode() { + if err := d.DownloadColumnsAndRecoverBlobs(downloadCtx, blocks); err != nil { + cancel() + log.Warn("failed to download columns and recover blobs", "err", err) + return + } + } else { + if err := d.DownloadOnlyCustodyColumns(downloadCtx, blocks); err != nil { + cancel() + log.Warn("failed to download only custody columns", "err", err) + return + } + } + cancel() + for i, root := range roots { + available, err := d.IsDataAvailable(blocks[i].GetSlot(), root) + if err != nil { + log.Warn("failed to recheck data column availability", "slot", blocks[i].GetSlot(), "blockRoot", root, "err", err) + continue + } + if available { + d.deleteDeferredColumnSyncJob(root) + log.Debug("[syncColumnDataWorker] column data is synced, removing from sync queue", "slot", blocks[i].GetSlot(), "blockRoot", root) } } } diff --git a/cl/das/peer_das_test.go b/cl/das/peer_das_test.go index e56b6042e4c..486c99e2427 100644 --- a/cl/das/peer_das_test.go +++ b/cl/das/peer_das_test.go @@ -17,9 +17,11 @@ package das import ( + "context" "errors" "fmt" "testing" + "time" "github.com/stretchr/testify/require" @@ -29,6 +31,99 @@ import ( "github.com/erigontech/erigon/common" ) +func TestSyncColumnDataLaterStoresCompactFuluBlock(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + d := &peerdas{} + block := cltypes.NewSignedBeaconBlock(&cfg, clparams.FuluVersion) + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + require.NoError(t, d.SyncColumnDataLater(block)) + + value, ok := d.blocksToCheckSync[common.Hash(blockRoot)] + require.True(t, ok) + queued, compact := value.block.(*deferredColumnSyncBlock) + require.True(t, compact) + require.Equal(t, block.Block.Slot, queued.slot) + require.Equal(t, common.Hash(blockRoot), queued.root) + require.Equal(t, block.Block.Body.BlobKzgCommitments.Len(), queued.commitments.Len()) +} + +func TestSyncColumnDataLaterBoundsQueue(t *testing.T) { + const queueLimit = maxDeferredColumnSyncBlocks + cfg := clparams.MainnetBeaconConfig + d := &peerdas{} + for slot := uint64(1); slot <= queueLimit+1; slot++ { + block := cltypes.NewSignedBeaconBlock(&cfg, clparams.FuluVersion) + block.Block.Slot = slot + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + require.NoError(t, d.SyncColumnDataLater(block)) + } + + require.LessOrEqual(t, len(d.blocksToCheckSync), queueLimit) +} + +func TestSyncColumnDataLaterUsesConfiguredFork(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + cfg.AltairForkEpoch = 0 + cfg.BellatrixForkEpoch = 0 + cfg.CapellaForkEpoch = 0 + cfg.DenebForkEpoch = 0 + cfg.ElectraForkEpoch = 0 + cfg.FuluForkEpoch = 2 + cfg.InitializeForkSchedule() + + tests := []struct { + name string + blockSlot uint64 + decodedVersion clparams.StateVersion + wantQueued bool + }{ + {name: "Fulu slot decoded as Electra", blockSlot: 2 * cfg.SlotsPerEpoch, decodedVersion: clparams.ElectraVersion, wantQueued: true}, + {name: "Electra slot decoded as Fulu", blockSlot: cfg.SlotsPerEpoch, decodedVersion: clparams.FuluVersion}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + d := &peerdas{beaconConfig: &cfg} + block := cltypes.NewSignedBeaconBlock(&cfg, tt.decodedVersion) + block.Block.Slot = tt.blockSlot + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + require.NoError(t, d.SyncColumnDataLater(block)) + job, queued := d.blocksToCheckSync[common.Hash(blockRoot)] + require.Equal(t, tt.wantQueued, queued) + if tt.wantQueued { + require.Equal(t, clparams.FuluVersion, job.block.Version()) + } + }) + } +} + +func TestDownloadOnlyCustodyColumnsWithoutRPC(t *testing.T) { + d := &peerdas{} + require.Error(t, d.DownloadOnlyCustodyColumns(context.Background(), nil)) +} + +func TestSyncDeferredColumnDataPrunesExpiredJobsWithoutRPC(t *testing.T) { + now := time.Now() + expiredRoot := common.Hash{1} + activeRoot := common.Hash{2} + d := &peerdas{blocksToCheckSync: map[common.Hash]deferredColumnSyncJob{ + expiredRoot: {block: &deferredColumnSyncBlock{}, addedAt: now.Add(-deferredColumnSyncTTL - time.Second)}, + activeRoot: {block: &deferredColumnSyncBlock{}, addedAt: now}, + }} + + d.syncDeferredColumnData(t.Context()) + + jobs := d.deferredColumnSyncJobs() + require.NotContains(t, jobs, expiredRoot) + require.Contains(t, jobs, activeRoot) +} + // initTestBeaconConfig installs cfg as the global config if no test has done so // yet. InitGlobalStaticConfig panics on a second call, so tests in this package // must agree on every global-only field; they may differ only in fork epochs, diff --git a/cl/phase1/forkchoice/data_availability.go b/cl/phase1/forkchoice/data_availability.go new file mode 100644 index 00000000000..37a45b1ce81 --- /dev/null +++ b/cl/phase1/forkchoice/data_availability.go @@ -0,0 +1,42 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package forkchoice + +import ( + "fmt" + + "github.com/erigontech/erigon/cl/cltypes" + "github.com/erigontech/erigon/common" + "github.com/erigontech/erigon/common/log/v3" +) + +func (f *ForkChoiceStore) requireDataColumnAvailability(block *cltypes.SignedBeaconBlock, blockRoot common.Hash) error { + if f.peerDas == nil { + return fmt.Errorf("%w: peer DAS is not configured", ErrEIP7594ColumnDataNotAvailable) + } + available, err := f.peerDas.IsDataAvailable(block.Block.Slot, blockRoot) + if err != nil { + return fmt.Errorf("failed to check data column availability: %w", err) + } + if available { + return nil + } + if err := f.peerDas.SyncColumnDataLater(block); err != nil { + log.Warn("failed to schedule data column sync", "slot", block.Block.Slot, "blockRoot", blockRoot, "err", err) + } + return ErrEIP7594ColumnDataNotAvailable +} diff --git a/cl/phase1/forkchoice/on_block.go b/cl/phase1/forkchoice/on_block.go index d32b9904b08..d0f99671749 100644 --- a/cl/phase1/forkchoice/on_block.go +++ b/cl/phase1/forkchoice/on_block.go @@ -28,6 +28,7 @@ import ( "github.com/erigontech/erigon/cl/clparams" "github.com/erigontech/erigon/cl/cltypes" "github.com/erigontech/erigon/cl/cltypes/solid" + "github.com/erigontech/erigon/cl/das" "github.com/erigontech/erigon/cl/monitor" "github.com/erigontech/erigon/cl/persistence/beacon_indicies" "github.com/erigontech/erigon/cl/phase1/core/state" @@ -55,23 +56,11 @@ var ( ) func verifyKzgCommitmentsAgainstTransactions(cfg *clparams.BeaconChainConfig, block *cltypes.BeaconBlock) error { - expectedBlobHashes := []common.Hash{} + expectedBlobHashes := kzgCommitmentsToVersionedHashes(block.Body.BlobKzgCommitments) transactions, err := types.DecodeTransactions(block.Body.ExecutionPayload.Transactions.UnderlyngReference()) if err != nil { return fmt.Errorf("unable to decode transactions: %w", err) } - block.Body.BlobKzgCommitments.Range(func(index int, value *cltypes.KZGCommitment, length int) bool { - var kzg common.Hash - kzg, err = utils.KzgCommitmentToVersionedHash(common.Bytes48(*value)) - if err != nil { - return false - } - expectedBlobHashes = append(expectedBlobHashes, kzg) - return true - }) - if err != nil { - return err - } maxBlobsPerBlock := cfg.MaxBlobsPerBlockByVersion(block.Version()) if block.Version() >= clparams.FuluVersion { @@ -80,6 +69,19 @@ func verifyKzgCommitmentsAgainstTransactions(cfg *clparams.BeaconChainConfig, bl return misc.ValidateBlobs(block.Body.ExecutionPayload.BlobGasUsed, cfg.MaxBlobGasPerBlock, maxBlobsPerBlock, expectedBlobHashes, &transactions) } +func kzgCommitmentsToVersionedHashes(commitments *solid.ListSSZ[*cltypes.KZGCommitment]) []common.Hash { + if commitments == nil { + return nil + } + hashes := make([]common.Hash, 0, commitments.Len()) + commitments.Range(func(_ int, commitment *cltypes.KZGCommitment, _ int) bool { + versionedHash, _ := utils.KzgCommitmentToVersionedHash(common.Bytes48(*commitment)) + hashes = append(hashes, versionedHash) + return true + }) + return hashes +} + func collectOnBlockLatencyToUnixTime(ethClock eth_clock.EthereumClock, slot, currentSlotOnEntry uint64) { if slot != currentSlotOnEntry { return @@ -88,7 +90,21 @@ func collectOnBlockLatencyToUnixTime(ethClock eth_clock.EthereumClock, slot, cur monitor.ObserveBlockImportingLatency(initialSlotTime) } -func (f *ForkChoiceStore) OnBlock(ctx context.Context, block *cltypes.SignedBeaconBlock, newPayload, fullValidation, checkDataAvaiability bool) error { +func hasCompleteBlobDataForAllCommitments(blobs [][]byte, proofs [][][]byte, expectedCount int) bool { + if len(blobs) != expectedCount || len(proofs) != expectedCount { + return false + } + for i, blob := range blobs { + if len(blob) != cltypes.BYTES_PER_BLOB || + len(proofs[i]) != 1 || + len(proofs[i][0]) != cltypes.BYTES_KZG_PROOF { + return false + } + } + return true +} + +func (f *ForkChoiceStore) OnBlock(ctx context.Context, block *cltypes.SignedBeaconBlock, newPayload, fullValidation, checkDataAvailability bool) error { f.mu.Lock() unlocked := false defer f.drainQueuedWork() @@ -136,6 +152,7 @@ func (f *ForkChoiceStore) OnBlock(ctx context.Context, block *cltypes.SignedBeac // Validate parent payload status path early (before expensive operations) blockEpoch := f.computeEpochAtSlot(block.Block.Slot) blockVersion := f.beaconCfg.GetCurrentStateVersion(blockEpoch) + checkDataAvailability = checkDataAvailability || das.IsDataAvailabilityRequired(f.beaconCfg, f.Slot(), block.Block.Slot, blockVersion) isGloas := blockVersion >= clparams.GloasVersion headBeforeBlock := common.Hash{} if isGloas && f.Slot() == block.Block.Slot { @@ -159,47 +176,27 @@ func (f *ForkChoiceStore) OnBlock(ctx context.Context, block *cltypes.SignedBeac startEngine := time.Now() isVerifiedExecutionPayload := f.verifiedExecutionPayload.Contains(blockRoot) if blockVersion < clparams.GloasVersion { - // Find the versioned hashes from blob commitments + blobCommitmentCount := block.Block.Body.BlobKzgCommitments.Len() var versionedHashes []common.Hash if newPayload && f.engine != nil && block.Version() >= clparams.DenebVersion { - versionedHashes = []common.Hash{} - solid.RangeErr[*cltypes.KZGCommitment](block.Block.Body.BlobKzgCommitments, func(i1 int, k *cltypes.KZGCommitment, i2 int) error { - versionedHash, err := utils.KzgCommitmentToVersionedHash(common.Bytes48(*k)) - if err != nil { - return err - } - versionedHashes = append(versionedHashes, versionedHash) - return nil - }) + versionedHashes = kzgCommitmentsToVersionedHashes(block.Block.Body.BlobKzgCommitments) } - // Check if EL has blobs - elHasBlobs := false - if f.engine != nil && f.peerDas != nil && checkDataAvaiability && block.Block.Body.BlobKzgCommitments.Len() > 0 && !f.peerDas.IsArchivedMode() { - blobsWithProof, proofs, err := f.engine.GetBlobs(ctx, versionedHashes, block.Version()) + elHasBlobData := false + if newPayload && f.engine != nil && f.peerDas != nil && checkDataAvailability && blockVersion < clparams.FuluVersion && blobCommitmentCount > 0 && !f.peerDas.IsArchivedMode() { + blobs, proofs, err := f.engine.GetBlobs(ctx, versionedHashes, block.Version()) if err != nil { log.Warn("OnBlock: GetBlobs failed", "blockRoot", common.Hash(blockRoot), "err", err) } - elHasBlobs = err == nil && len(blobsWithProof) == len(versionedHashes) && len(proofs) == len(versionedHashes) - log.Trace("OnBlock: EL blob data availability", "blockRoot", common.Hash(blockRoot), "elHasBlobs", elHasBlobs) + elHasBlobData = err == nil && hasCompleteBlobDataForAllCommitments(blobs, proofs, blobCommitmentCount) + log.Trace("OnBlock: EL blob data availability", "blockRoot", common.Hash(blockRoot), "elHasBlobData", elHasBlobData) } - // Check if blob data is available (skip if blobs are in txpool) - if checkDataAvaiability && block.Block.Body.BlobKzgCommitments.Len() > 0 && !elHasBlobs { - if block.Version() >= clparams.FuluVersion && f.peerDas != nil { - available, err := f.peerDas.IsDataAvailable(block.Block.Slot, blockRoot) - if err != nil { + if checkDataAvailability && blobCommitmentCount > 0 && !elHasBlobData { + if blockVersion == clparams.FuluVersion { + if err := f.requireDataColumnAvailability(block, blockRoot); err != nil { return err } - if !available { - if f.syncedDataManager.Syncing() { - return ErrEIP7594ColumnDataNotAvailable - } else { - if err := f.peerDas.SyncColumnDataLater(block); err != nil { - log.Warn("failed to schedule deferred column data sync", "slot", block.Block.Slot, "blockRoot", blockRoot, "err", err) - } - } - } } else if block.Version() >= clparams.DenebVersion { if err := f.isDataAvailable(ctx, block.Block.Slot, blockRoot, block.Block.Body.BlobKzgCommitments); err != nil { if errors.Is(err, ErrEIP4844DataNotAvailable) { @@ -316,10 +313,6 @@ func (f *ForkChoiceStore) OnBlock(ctx context.Context, block *cltypes.SignedBeac f.eth2Roots.Add(blockRoot, parentHash) } } - // Note: highestSeen was already updated before AddChainSegment (line ~216) - // so aggregates/attestations for this slot are accepted promptly. No second - // update needed here. - // Remove the parent from the head set delete(f.headSet, block.Block.ParentRoot) f.headSet[blockRoot] = struct{}{} @@ -370,7 +363,7 @@ func (f *ForkChoiceStore) OnBlock(ctx context.Context, block *cltypes.SignedBeac // Always validate payload with EL for pending envelopes, regardless of the caller's newPayload flag. // During forward sync newPayload is false, but the envelope still needs to reach the EL; // otherwise the EL never learns about this block and the chain stalls. - applied, applyErr := f.applyEnvelopeLocked(ctx, pending, checkDataAvaiability, true) + applied, applyErr := f.applyEnvelopeLocked(ctx, pending, checkDataAvailability, true) if applyErr != nil { log.Warn("OnBlock: failed to process pending envelope", "blockRoot", common.Hash(blockRoot), "err", applyErr) } else if applied { diff --git a/cl/phase1/forkchoice/on_block_test.go b/cl/phase1/forkchoice/on_block_test.go new file mode 100644 index 00000000000..5d513850b94 --- /dev/null +++ b/cl/phase1/forkchoice/on_block_test.go @@ -0,0 +1,137 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package forkchoice + +import ( + "context" + "fmt" + "testing" + + "github.com/stretchr/testify/require" + "go.uber.org/mock/gomock" + + "github.com/erigontech/erigon/cl/clparams" + "github.com/erigontech/erigon/cl/cltypes" + "github.com/erigontech/erigon/cl/cltypes/solid" + dasmock "github.com/erigontech/erigon/cl/das/mock_services" + "github.com/erigontech/erigon/common" +) + +func TestOnBlockRequiresRecentFuluDataAvailability(t *testing.T) { + for _, checkDataAvailability := range []bool{false, true} { + t.Run(fmt.Sprintf("caller_check_%t", checkDataAvailability), func(t *testing.T) { + store := buildExAnteStore(t) + cfg := clparams.MainnetBeaconConfig + cfg.AltairForkEpoch = 0 + cfg.BellatrixForkEpoch = 0 + cfg.CapellaForkEpoch = 0 + cfg.DenebForkEpoch = 0 + cfg.ElectraForkEpoch = 0 + cfg.FuluForkEpoch = 1 + cfg.InitializeForkSchedule() + require.Equal(t, clparams.FuluVersion, cfg.GetCurrentStateVersion(cfg.FuluForkEpoch)) + store.beaconCfg = &cfg + + parentRoot, _, err := store.GetHead(nil) + require.NoError(t, err) + block := cltypes.NewSignedBeaconBlock(&cfg, clparams.FuluVersion) + block.Block.Slot = cfg.SlotsPerEpoch + block.Block.ParentRoot = parentRoot + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + store.OnTick(store.genesisTime + block.Block.Slot*cfg.SecondsPerSlot) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + ctrl := gomock.NewController(t) + peerDas := dasmock.NewMockPeerDas(ctrl) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(false, nil) + peerDas.EXPECT().SyncColumnDataLater(block).Return(nil) + store.InitPeerDas(peerDas) + + err = store.OnBlock(context.Background(), block, false, false, checkDataAvailability) + require.ErrorIs(t, err, ErrEIP7594ColumnDataNotAvailable) + }) + } +} + +func TestHasCompleteBlobDataForAllCommitments(t *testing.T) { + blob := make([]byte, cltypes.BYTES_PER_BLOB) + proof := make([]byte, cltypes.BYTES_KZG_PROOF) + tests := []struct { + name string + blobs [][]byte + proofs [][][]byte + expectedCount int + want bool + }{ + { + name: "missing blob", + blobs: make([][]byte, 1), + proofs: [][][]byte{{proof}}, + expectedCount: 1, + }, + { + name: "missing proofs", + blobs: [][]byte{blob}, + proofs: make([][][]byte, 1), + expectedCount: 1, + }, + { + name: "truncated blob", + blobs: [][]byte{{1}}, + proofs: [][][]byte{{proof}}, + expectedCount: 1, + }, + { + name: "truncated proof", + blobs: [][]byte{blob}, + proofs: [][][]byte{{{1}}}, + expectedCount: 1, + }, + { + name: "multiple proofs", + blobs: [][]byte{blob}, + proofs: [][][]byte{{proof, proof}}, + expectedCount: 1, + }, + { + name: "unexpected count", + blobs: [][]byte{blob}, + proofs: [][][]byte{{proof}}, + expectedCount: 2, + }, + { + name: "complete entries", + blobs: [][]byte{blob, blob}, + proofs: [][][]byte{{proof}, {proof}}, + expectedCount: 2, + want: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + require.Equal(t, tt.want, hasCompleteBlobDataForAllCommitments(tt.blobs, tt.proofs, tt.expectedCount)) + }) + } +} + +func TestKzgCommitmentsToVersionedHashesPreservesEmptyList(t *testing.T) { + commitments := solid.NewStaticListSSZ[*cltypes.KZGCommitment](1, 48) + require.NotNil(t, kzgCommitmentsToVersionedHashes(commitments)) + require.Empty(t, kzgCommitmentsToVersionedHashes(commitments)) +} diff --git a/cl/phase1/forkchoice/on_execution_payload.go b/cl/phase1/forkchoice/on_execution_payload.go index 6de1ef11431..f8f4be1a9b9 100644 --- a/cl/phase1/forkchoice/on_execution_payload.go +++ b/cl/phase1/forkchoice/on_execution_payload.go @@ -25,14 +25,12 @@ import ( "github.com/erigontech/erigon/cl/abstract" "github.com/erigontech/erigon/cl/clparams" "github.com/erigontech/erigon/cl/cltypes" - "github.com/erigontech/erigon/cl/cltypes/solid" "github.com/erigontech/erigon/cl/fork" "github.com/erigontech/erigon/cl/monitor" "github.com/erigontech/erigon/cl/persistence/beacon_indicies" "github.com/erigontech/erigon/cl/phase1/core/state" "github.com/erigontech/erigon/cl/phase1/execution_client" "github.com/erigontech/erigon/cl/transition" - "github.com/erigontech/erigon/cl/utils" "github.com/erigontech/erigon/cl/utils/bls" "github.com/erigontech/erigon/common" "github.com/erigontech/erigon/common/hexutil" @@ -178,49 +176,21 @@ func (f *ForkChoiceStore) verifyEnvelopeBuilderSignature( return nil } -// checkDataAvailability checks if blob data is available for the execution payload. -// For GLOAS, blob_kzg_commitments are in the committed bid, not directly in BeaconBlock. -// Returns nil if data is available, ErrEIP7594ColumnDataNotAvailable if not available yet. func (f *ForkChoiceStore) checkDataAvailability( - ctx context.Context, + _ context.Context, block *cltypes.SignedBeaconBlock, beaconBlockRoot common.Hash, ) error { - // Get committed bid from the block committedBid := block.Block.Body.GetSignedExecutionPayloadBid() if committedBid == nil || committedBid.Message == nil { - // No bid means no blobs to check return nil } blobCommitments := &committedBid.Message.BlobKzgCommitments if blobCommitments.Len() == 0 { - // No blobs to check return nil } - - // Check PeerDAS data availability - // Note: Unlike OnBlock, we don't skip this check even if EL has blobs, - // because we need to ensure blobs are stored in CL's blob storage for beacon API. - available, err := f.peerDas.IsDataAvailable(block.Block.Slot, beaconBlockRoot) - if err != nil { - return fmt.Errorf("checkDataAvailability: failed to check data availability: %w", err) - } - if !available { - if f.syncedDataManager.Syncing() { - // During sync, return error immediately to retry later - return ErrEIP7594ColumnDataNotAvailable - } - // Not syncing - schedule deferred column data sync - if err := f.peerDas.SyncColumnDataLater(block); err != nil { - log.Warn("checkDataAvailability: failed to schedule deferred column data sync", - "slot", block.Block.Slot, "beaconBlockRoot", beaconBlockRoot, "err", err) - } - // Return error so envelope can be queued for later processing - return ErrEIP7594ColumnDataNotAvailable - } - - return nil + return f.requireDataColumnAvailability(block, beaconBlockRoot) } // validatePayloadWithEL validates the execution payload with the execution layer engine. @@ -241,22 +211,8 @@ func (f *ForkChoiceStore) validatePayloadWithEL( return errors.New("validatePayloadWithEL: block missing execution payload bid") } - // Calculate versioned hashes from committed bid's blob_kzg_commitments - versionedHashes := make([]common.Hash, 0) blobCommitments := &committedBid.Message.BlobKzgCommitments - if blobCommitments.Len() > 0 { - versionedHashes = make([]common.Hash, 0, blobCommitments.Len()) - if err := solid.RangeErr[*cltypes.KZGCommitment](blobCommitments, func(_ int, k *cltypes.KZGCommitment, _ int) error { - versionedHash, err := utils.KzgCommitmentToVersionedHash(common.Bytes48(*k)) - if err != nil { - return err - } - versionedHashes = append(versionedHashes, versionedHash) - return nil - }); err != nil { - return fmt.Errorf("validatePayloadWithEL: failed to compute versioned hashes: %w", err) - } - } + versionedHashes := kzgCommitmentsToVersionedHashes(blobCommitments) // Get execution requests list var executionRequestsList []hexutil.Bytes diff --git a/cl/phase1/network/services/block_service.go b/cl/phase1/network/services/block_service.go index 990a11a95ff..c21fcbf662c 100644 --- a/cl/phase1/network/services/block_service.go +++ b/cl/phase1/network/services/block_service.go @@ -51,8 +51,10 @@ type proposerIndexAndSlot struct { } type blockJob struct { - block *cltypes.SignedBeaconBlock - creationTime time.Time + block *cltypes.SignedBeaconBlock + creationTime time.Time + retryAfter time.Time + dataAvailabilityAttempts uint8 } type blockService struct { @@ -64,11 +66,10 @@ type blockService struct { // reference: https://github.com/ethereum/consensus-specs/blob/dev/specs/phase0/p2p-interface.md#beacon_block seenBlocksCache *lru.Cache[proposerIndexAndSlot, struct{}] - // blocks that should be scheduled for later execution (e.g missing blobs). emitter *beaconevents.EventEmitter blocksScheduledForLaterExecution sync.Map - // store the block in db - db kv.RwDB + db kv.RwDB + now func() time.Time } // NewBlockService creates a new block service @@ -93,6 +94,7 @@ func NewBlockService( seenBlocksCache: seenBlocksCache, emitter: emitter, db: db, + now: time.Now, } go b.loop(ctx) return b @@ -159,7 +161,6 @@ func (b *blockService) ProcessMessage(ctx context.Context, _ *uint64, msg *cltyp } return err } - // [IGNORE] The block's parent (defined by block.parent_root) has been seen (via both gossip and non-gossip sources) (a client MAY queue blocks for processing once the parent block is retrieved). parentHeader, ok := b.forkchoiceStore.GetHeader(msg.Block.ParentRoot) if !ok { @@ -217,10 +218,11 @@ func (b *blockService) ProcessMessage(ctx context.Context, _ *uint64, msg *cltyp // i.e. validate that len(body.signed_beacon_block.message.blob_kzg_commitments) <= MAX_BLOBS_PER_BLOCK return ErrInvalidCommitmentsCount } + b.seenBlocksCache.Add(seenCacheKey, struct{}{}) b.publishBlockGossipEvent(msg) // the rest of the validation is done in the forkchoice store if err := b.processAndStoreBlock(ctx, msg); err != nil { - if errors.Is(err, forkchoice.ErrEIP4844DataNotAvailable) || errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) || errors.Is(err, forkchoice.ErrParentEnvelopePending) { + if isDataAvailabilityError(err) || errors.Is(err, forkchoice.ErrParentEnvelopePending) { b.scheduleBlockForLaterProcessing(msg) return nil } @@ -260,12 +262,25 @@ func (b *blockService) scheduleBlockForLaterProcessing(block *cltypes.SignedBeac return } - b.blocksScheduledForLaterExecution.Store(blockRoot, &blockJob{ + now := b.currentTime() + b.blocksScheduledForLaterExecution.LoadOrStore(blockRoot, &blockJob{ block: block, - creationTime: time.Now(), + creationTime: now, + retryAfter: now.Add(blockRetryInterval), }) } +func isDataAvailabilityError(err error) bool { + return errors.Is(err, forkchoice.ErrEIP4844DataNotAvailable) || errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) +} + +func (b *blockService) currentTime() time.Time { + if b.now == nil { + return time.Now() + } + return b.now() +} + // processAndStoreBlock processes and stores a block func (b *blockService) processAndStoreBlock(ctx context.Context, block *cltypes.SignedBeaconBlock) error { blockRoot, err := block.Block.HashSSZ() @@ -330,19 +345,35 @@ func (b *blockService) loop(ctx context.Context) { return case <-ticker.C: } - b.blocksScheduledForLaterExecution.Range(func(key, value any) bool { - blockJob := value.(*blockJob) - // check if it has expired - if time.Since(blockJob.creationTime) > blockJobExpiry { - b.blocksScheduledForLaterExecution.Delete(key.([32]byte)) - return true - } - if err := b.processAndStoreBlock(ctx, blockJob.block); err != nil { - log.Trace("Failed to process and store block", "block", blockJob.block, "error", err) - return true - } + b.processScheduledBlocks(ctx, b.currentTime()) + } +} + +func (b *blockService) processScheduledBlocks(ctx context.Context, now time.Time) { + b.blocksScheduledForLaterExecution.Range(func(key, value any) bool { + blockJob := value.(*blockJob) + if now.Sub(blockJob.creationTime) > blockJobExpiry { b.blocksScheduledForLaterExecution.Delete(key.([32]byte)) return true - }) - } + } + if now.Before(blockJob.retryAfter) { + return true + } + if err := b.processAndStoreBlock(ctx, blockJob.block); err != nil { + if isDataAvailabilityError(err) { + blockJob.dataAvailabilityAttempts++ + if blockJob.dataAvailabilityAttempts >= maxDataAvailabilityRetries { + b.blocksScheduledForLaterExecution.Delete(key.([32]byte)) + } else { + blockJob.retryAfter = b.currentTime().Add(blockRetryInterval) + } + } else { + blockJob.retryAfter = b.currentTime().Add(blockRetryInterval) + } + log.Trace("Failed to process and store block", "block", blockJob.block, "error", err) + return true + } + b.blocksScheduledForLaterExecution.Delete(key.([32]byte)) + return true + }) } diff --git a/cl/phase1/network/services/block_service_test.go b/cl/phase1/network/services/block_service_test.go index 6ce009d23f5..59b88ed2b37 100644 --- a/cl/phase1/network/services/block_service_test.go +++ b/cl/phase1/network/services/block_service_test.go @@ -21,6 +21,7 @@ import ( "context" "errors" "testing" + "time" "github.com/stretchr/testify/require" "go.uber.org/mock/gomock" @@ -47,6 +48,46 @@ func (s attesterSlashingErrorStore) OnAttesterSlashing(*cltypes.AttesterSlashing return s.err } +type dataUnavailableStore struct { + *mock_services.ForkChoiceStorageMock + onBlockCalls int + onBlockResult error + afterOnBlock func() +} + +func (s *dataUnavailableStore) OnBlock(context.Context, *cltypes.SignedBeaconBlock, bool, bool, bool) error { + s.onBlockCalls++ + if s.afterOnBlock != nil { + s.afterOnBlock() + } + return s.onBlockResult +} + +type blockServiceTestClock struct { + now time.Time +} + +func (c *blockServiceTestClock) currentTime() time.Time { + return c.now +} + +func newDataUnavailableBlockJob(t *testing.T) (*blockService, *dataUnavailableStore, *cltypes.SignedBeaconBlock, *blockServiceTestClock) { + t.Helper() + store := &dataUnavailableStore{ + ForkChoiceStorageMock: mock_services.NewForkChoiceStorageMock(t), + onBlockResult: forkchoice.ErrEIP7594ColumnDataNotAvailable, + } + clock := &blockServiceTestClock{now: time.Unix(1, 0)} + service := &blockService{ + forkchoiceStore: store, + db: mdbxtest.NewTestDB(t, dbcfg.ChainDB), + now: clock.currentTime, + } + block := cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.FuluVersion) + service.scheduleBlockForLaterProcessing(block) + return service, store, block, clock +} + func setupBlockService(t *testing.T, ctrl *gomock.Controller) (BlockService, *synced_data.SyncedDataManager, *eth_clock.MockEthereumClock, *mock_services.ForkChoiceStorageMock) { db := mdbxtest.NewTestDB(t, dbcfg.ChainDB) cfg := &clparams.MainnetBeaconConfig @@ -155,7 +196,7 @@ func TestBlockServiceSuccess(t *testing.T) { blocks, _, post := tests.GetBellatrixRandom() - blockService, syncedData, ethClock, fcu := setupBlockService(t, ctrl) + blockServiceAPI, syncedData, ethClock, fcu := setupBlockService(t, ctrl) syncedData.OnHeadState(post) ethClock.EXPECT().GetCurrentSlot().Return(uint64(0)).AnyTimes() ethClock.EXPECT().IsSlotCurrentSlotWithMaximumClockDisparity(gomock.Any()).Return(true).AnyTimes() @@ -163,7 +204,92 @@ func TestBlockServiceSuccess(t *testing.T) { fcu.Headers[blocks[1].Block.ParentRoot] = blocks[0].SignedBeaconBlockHeader().Header.Copy() blocks[1].Block.Body.BlobKzgCommitments = solid.NewStaticListSSZ[*cltypes.KZGCommitment](100, 48) - require.NoError(t, blockService.ProcessMessage(context.Background(), nil, blocks[1])) + require.NoError(t, blockServiceAPI.ProcessMessage(context.Background(), nil, blocks[1])) + service := blockServiceAPI.(*blockService) + require.True(t, service.seenBlocksCache.Contains(proposerIndexAndSlot{ + proposerIndex: blocks[1].Block.ProposerIndex, + slot: blocks[1].Block.Slot, + })) +} + +func TestDataAvailabilityRetriesAreDelayed(t *testing.T) { + service, store, _, clock := newDataUnavailableBlockJob(t) + clock.now = clock.now.Add(blockRetryInterval) + service.processScheduledBlocks(t.Context(), clock.now) + clock.now = clock.now.Add(blockRetryInterval - time.Nanosecond) + service.processScheduledBlocks(t.Context(), clock.now) + + require.Equal(t, 1, store.onBlockCalls) +} + +func TestFirstScheduledBlockRetryIsDelayed(t *testing.T) { + service, store, block, clock := newDataUnavailableBlockJob(t) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + service.blocksScheduledForLaterExecution.Delete(blockRoot) + + now := clock.now + service.scheduleBlockForLaterProcessing(block) + clock.now = now.Add(blockRetryInterval - time.Nanosecond) + service.processScheduledBlocks(t.Context(), clock.now) + require.Zero(t, store.onBlockCalls) + + clock.now = now.Add(blockRetryInterval) + service.processScheduledBlocks(t.Context(), clock.now) + require.Equal(t, 1, store.onBlockCalls) +} + +func TestDataAvailabilityRetryDelayStartsAfterAttempt(t *testing.T) { + service, store, block, clock := newDataUnavailableBlockJob(t) + startedAt := clock.now.Add(blockRetryInterval) + finishedAt := startedAt.Add(2 * time.Second) + clock.now = startedAt + store.afterOnBlock = func() { clock.now = finishedAt } + + service.processScheduledBlocks(t.Context(), startedAt) + + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + value, ok := service.blocksScheduledForLaterExecution.Load(blockRoot) + require.True(t, ok) + require.Equal(t, finishedAt.Add(blockRetryInterval), value.(*blockJob).retryAfter) +} + +func TestDataAvailabilityRetriesAreBounded(t *testing.T) { + const retryLimit = 4 + for _, availabilityErr := range []error{ + forkchoice.ErrEIP4844DataNotAvailable, + forkchoice.ErrEIP7594ColumnDataNotAvailable, + } { + t.Run(availabilityErr.Error(), func(t *testing.T) { + service, store, block, clock := newDataUnavailableBlockJob(t) + store.onBlockResult = availabilityErr + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + for range retryLimit + 1 { + clock.now = clock.now.Add(blockRetryInterval + time.Nanosecond) + service.processScheduledBlocks(t.Context(), clock.now) + } + + require.Equal(t, retryLimit, store.onBlockCalls) + _, pending := service.blocksScheduledForLaterExecution.Load(blockRoot) + require.False(t, pending) + }) + } +} + +func TestSchedulingSameBlockPreservesDataAvailabilityRetryBudget(t *testing.T) { + service, _, block, clock := newDataUnavailableBlockJob(t) + clock.now = clock.now.Add(blockRetryInterval) + service.processScheduledBlocks(t.Context(), clock.now) + service.scheduleBlockForLaterProcessing(block) + + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + value, ok := service.blocksScheduledForLaterExecution.Load(blockRoot) + require.True(t, ok) + require.Equal(t, uint8(1), value.(*blockJob).dataAvailabilityAttempts) } func TestImportBlockOperationsAttesterSlashingLogging(t *testing.T) { diff --git a/cl/phase1/network/services/constants.go b/cl/phase1/network/services/constants.go index 5977b35aa5a..9fcf543ae1c 100644 --- a/cl/phase1/network/services/constants.go +++ b/cl/phase1/network/services/constants.go @@ -30,6 +30,8 @@ const ( operationSeenCacheSize = 16_384 seenBlockCacheSize = 1000 // SeenBlockCacheSize is the size of the cache for seen blocks. blockJobsIntervalTick = 50 * time.Millisecond + blockRetryInterval = time.Second + maxDataAvailabilityRetries = 4 blobJobsIntervalTick = 5 * time.Millisecond singleAttestationIntervalTick = 10 * time.Millisecond attestationJobsIntervalTick = 100 * time.Millisecond diff --git a/cl/phase1/stages/chain_tip_sync.go b/cl/phase1/stages/chain_tip_sync.go index 02dffddf278..78afa2b0023 100644 --- a/cl/phase1/stages/chain_tip_sync.go +++ b/cl/phase1/stages/chain_tip_sync.go @@ -106,9 +106,7 @@ func waitForExecutionEngineToBeFinished(ctx context.Context, cfg *Cfg) (ready bo } } -// fetchBlocksFromReqResp retrieves blocks starting from a specified block number and continues for a given count. -// It sends a request to fetch the blocks, verifies the associated blobs, and inserts them into the blob store. -// It returns a PeeredObject containing the blocks and the peer ID, or an error if something goes wrong. +// fetchBlocksFromReqResp requests a block range and returns the response sorted by slot. func fetchBlocksFromReqResp(ctx context.Context, cfg *Cfg, from uint64, count uint64) (*peers.PeeredObject[[]*cltypes.SignedBeaconBlock], error) { blocks, pid, err := cfg.rpc.SendBeaconBlocksByRangeReq(ctx, from, count) for err != nil { @@ -180,8 +178,7 @@ func startFetchingBlocksMissedByGossipAfterSomeTime(ctx context.Context, cfg *Cf // Fetch blocks from the specified range blocks, err := fetchBlocksFromReqResp(ctx, cfg, from, count) if err != nil { - // Send error to the error channel and return - errCh <- err + sendFetchError(ctx, errCh, err) return } @@ -195,6 +192,17 @@ func startFetchingBlocksMissedByGossipAfterSomeTime(ctx context.Context, cfg *Cf } } +func sendFetchError(ctx context.Context, errCh chan<- error, err error) { + select { + case errCh <- err: + case <-ctx.Done(): + } +} + +func newChainTipBlockResponseChannel() chan *peers.PeeredObject[[]*cltypes.SignedBeaconBlock] { + return make(chan *peers.PeeredObject[[]*cltypes.SignedBeaconBlock], 1) +} + // listenToIncomingBlocksUntilANewBlockIsReceived listens for incoming blocks until a new block with a slot greater than or equal to the target slot is received. // It processes blocks, checks their validity, and publishes them. It also handles context cancellation and logs progress periodically. func listenToIncomingBlocksUntilANewBlockIsReceived(ctx context.Context, logger log.Logger, cfg *Cfg, args Args, respCh <-chan *peers.PeeredObject[[]*cltypes.SignedBeaconBlock], errCh chan error) error { @@ -208,6 +216,7 @@ func listenToIncomingBlocksUntilANewBlockIsReceived(ctx context.Context, logger // Map to keep track of seen block roots seenBlockRoots := make(map[common.Hash]struct{}) + dataAvailabilityRetryAfter := make(map[common.Hash]time.Time) MainLoop: for { select { @@ -253,6 +262,9 @@ MainLoop: if _, ok := seenBlockRoots[blockRoot]; ok { continue } + if retryAfter, ok := dataAvailabilityRetryAfter[blockRoot]; ok && time.Now().Before(retryAfter) { + continue + } // [GLOAS] Apply parent's envelope before processBlock so that // latestBlockHash is up-to-date for bid validation. @@ -265,9 +277,23 @@ MainLoop: } } - // Process the block - DA can be downloaded later if we are behind (see blobHistoryDownloader) - if err := processBlock(ctx, cfg, cfg.indiciesDB, block, true, true, false); err != nil { - log.Debug("bad blocks segment received", "err", err, "blockSlot", block.Block.Slot) + if err := acquireRecentBlockDataAvailability(ctx, cfg, block); err != nil { + log.Debug("failed to acquire chain-tip data columns", "err", err, "blockSlot", block.Block.Slot) + if errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) { + dataAvailabilityRetryAfter[blockRoot] = time.Now().Add(time.Second) + continue + } + seenBlockRoots[blockRoot] = struct{}{} + continue + } + + checkDataAvailability := block.Version() < clparams.GloasVersion + if err := processBlock(ctx, cfg, cfg.indiciesDB, block, true, true, checkDataAvailability); err != nil { + log.Debug("failed to process chain-tip block", "err", err, "blockSlot", block.Block.Slot) + if errors.Is(err, forkchoice.ErrEIP4844DataNotAvailable) || errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) { + dataAvailabilityRetryAfter[blockRoot] = time.Now().Add(time.Second) + continue + } seenBlockRoots[blockRoot] = struct{}{} continue } @@ -726,7 +752,7 @@ func chainTipSync(ctx context.Context, logger log.Logger, cfg *Cfg, args Args) e "targetSlot", args.targetSlot, "requestedSlots", totalRequest, ) - respCh := make(chan *peers.PeeredObject[[]*cltypes.SignedBeaconBlock], 1024) + respCh := newChainTipBlockResponseChannel() errCh := make(chan error) // 25 seconds is a good timeout for this diff --git a/cl/phase1/stages/chain_tip_sync_test.go b/cl/phase1/stages/chain_tip_sync_test.go new file mode 100644 index 00000000000..bf44185b447 --- /dev/null +++ b/cl/phase1/stages/chain_tip_sync_test.go @@ -0,0 +1,153 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package stages + +import ( + "context" + "errors" + "math" + "testing" + "time" + + "github.com/spf13/afero" + "github.com/stretchr/testify/require" + "go.uber.org/mock/gomock" + + "github.com/erigontech/erigon/cl/beacon/beacon_router_configuration" + "github.com/erigontech/erigon/cl/beacon/beaconevents" + "github.com/erigontech/erigon/cl/beacon/synced_data" + "github.com/erigontech/erigon/cl/clparams" + "github.com/erigontech/erigon/cl/cltypes" + dasmock "github.com/erigontech/erigon/cl/das/mock_services" + "github.com/erigontech/erigon/cl/persistence/blob_storage" + "github.com/erigontech/erigon/cl/phase1/core/state" + "github.com/erigontech/erigon/cl/phase1/forkchoice" + "github.com/erigontech/erigon/cl/phase1/forkchoice/fork_graph" + "github.com/erigontech/erigon/cl/phase1/forkchoice/public_keys_registry" + "github.com/erigontech/erigon/cl/pool" + "github.com/erigontech/erigon/cl/sentinel/peers" + "github.com/erigontech/erigon/cl/utils/eth_clock" + "github.com/erigontech/erigon/cl/validator/validator_params" + "github.com/erigontech/erigon/common" + "github.com/erigontech/erigon/common/log/v3" + "github.com/erigontech/erigon/db/kv/dbcfg" + "github.com/erigontech/erigon/db/kv/mdbx/mdbxtest" +) + +func TestSendFetchErrorReturnsAfterCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + cancel() + + done := make(chan struct{}) + go func() { + sendFetchError(ctx, make(chan error), errors.New("fetch failed")) + close(done) + }() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("error sender blocked after cancellation") + } +} + +func TestChainTipBlockResponseQueueHoldsOneBatch(t *testing.T) { + responses := newChainTipBlockResponseChannel() + first := &peers.PeeredObject[[]*cltypes.SignedBeaconBlock]{} + second := &peers.PeeredObject[[]*cltypes.SignedBeaconBlock]{} + + select { + case responses <- first: + default: + t.Fatal("first block batch was not buffered") + } + select { + case responses <- second: + t.Fatal("more than one block batch was buffered") + default: + } +} + +func TestChainTipSyncChecksFuluDataAvailability(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + clparams.ApplyMinimalPreset(&cfg) + cfg.AltairForkEpoch = 0 + cfg.BellatrixForkEpoch = 0 + cfg.CapellaForkEpoch = 0 + cfg.DenebForkEpoch = 0 + cfg.ElectraForkEpoch = 0 + cfg.FuluForkEpoch = 0 + cfg.InitializeForkSchedule() + + anchorState := state.New(&cfg) + anchorState.SetVersion(clparams.FuluVersion) + anchorRoot, err := anchorState.BlockRoot() + require.NoError(t, err) + + db := mdbxtest.NewTestDB(t, dbcfg.ChainDB) + clock := eth_clock.NewEthereumClock(0, common.Hash{}, &cfg) + store, err := forkchoice.NewForkChoiceStore( + clock, + anchorState, + nil, + pool.NewOperationsPool(&cfg), + fork_graph.NewForkGraphDisk(anchorState, nil, afero.NewMemMapFs(), beacon_router_configuration.RouterConfiguration{}), + beaconevents.NewEventEmitter(), + synced_data.NewSyncedDataManager(&cfg, true), + blob_storage.NewBlobStore(db, afero.NewMemMapFs(), math.MaxUint64, &cfg, clock), + public_keys_registry.NewInMemoryPublicKeysRegistry(), + validator_params.NewValidatorParams(), + false, + nil, + ) + require.NoError(t, err) + + block := cltypes.NewSignedBeaconBlock(&cfg, clparams.FuluVersion) + block.Block.Slot = 1 + block.Block.ParentRoot = anchorRoot + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + store.OnTick(block.Block.Slot * cfg.SecondsPerSlot) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + ctrl := gomock.NewController(t) + stageClock := eth_clock.NewMockEthereumClock(ctrl) + stageClock.EXPECT().GetCurrentSlot().Return(block.Block.Slot).AnyTimes() + peerDas := dasmock.NewMockPeerDas(ctrl) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(false, nil) + peerDas.EXPECT().IsArchivedMode().Return(false) + peerDas.EXPECT().DownloadOnlyCustodyColumns(gomock.Any(), gomock.Any()).Return(nil) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(false, nil) + store.InitPeerDas(peerDas) + + respCh := make(chan *peers.PeeredObject[[]*cltypes.SignedBeaconBlock], 1) + respCh <- &peers.PeeredObject[[]*cltypes.SignedBeaconBlock]{Data: []*cltypes.SignedBeaconBlock{block}} + errCh := make(chan error) + ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond) + defer cancel() + + err = listenToIncomingBlocksUntilANewBlockIsReceived(ctx, log.Root(), &Cfg{ + beaconCfg: &cfg, + ethClock: stageClock, + forkChoice: store, + indiciesDB: db, + peerDas: peerDas, + }, Args{targetSlot: block.Block.Slot}, respCh, errCh) + require.ErrorIs(t, err, context.DeadlineExceeded) + _, imported := store.GetHeader(common.Hash(blockRoot)) + require.False(t, imported) +} diff --git a/cl/phase1/stages/data_availability.go b/cl/phase1/stages/data_availability.go new file mode 100644 index 00000000000..004b1d724a7 --- /dev/null +++ b/cl/phase1/stages/data_availability.go @@ -0,0 +1,84 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package stages + +import ( + "context" + "fmt" + "time" + + "github.com/erigontech/erigon/cl/cltypes" + "github.com/erigontech/erigon/cl/das" + "github.com/erigontech/erigon/cl/phase1/forkchoice" + "github.com/erigontech/erigon/common" +) + +func acquireBlockDataAvailability(ctx context.Context, peerDas das.PeerDas, block *cltypes.SignedBeaconBlock) error { + if peerDas == nil { + return forkchoice.ErrEIP7594ColumnDataNotAvailable + } + blockRoot, err := block.Block.HashSSZ() + if err != nil { + return err + } + available, err := peerDas.IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)) + if err != nil { + return fmt.Errorf("check data column availability: %w", err) + } + if available { + return nil + } + blocks := []cltypes.ColumnSyncableSignedBlock{block} + if peerDas.IsArchivedMode() { + err = peerDas.DownloadColumnsAndRecoverBlobs(ctx, blocks) + } else { + err = peerDas.DownloadOnlyCustodyColumns(ctx, blocks) + } + if err != nil { + return fmt.Errorf("download data columns: %w", err) + } + available, err = peerDas.IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)) + if err != nil { + return fmt.Errorf("recheck data column availability: %w", err) + } + if !available { + return forkchoice.ErrEIP7594ColumnDataNotAvailable + } + return nil +} + +func acquireRecentBlockDataAvailability(ctx context.Context, cfg *Cfg, block *cltypes.SignedBeaconBlock) error { + commitments := block.GetBlobKzgCommitments() + if commitments == nil || commitments.Len() == 0 || !requiresRecentBlockDataAvailability(cfg, block) { + return nil + } + timeout := time.Duration(cfg.beaconCfg.SecondsPerSlot) * time.Second + if timeout <= 0 { + timeout = time.Second + } + downloadCtx, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + return acquireBlockDataAvailability(downloadCtx, cfg.peerDas, block) +} + +func requiresRecentBlockDataAvailability(cfg *Cfg, block *cltypes.SignedBeaconBlock) bool { + blockVersion := block.Version() + if cfg.beaconCfg.SlotsPerEpoch != 0 { + blockVersion = cfg.beaconCfg.GetCurrentStateVersion(block.Block.Slot / cfg.beaconCfg.SlotsPerEpoch) + } + return das.IsDataAvailabilityRequired(cfg.beaconCfg, cfg.ethClock.GetCurrentSlot(), block.Block.Slot, blockVersion) +} diff --git a/cl/phase1/stages/data_availability_test.go b/cl/phase1/stages/data_availability_test.go new file mode 100644 index 00000000000..e194bde8a3a --- /dev/null +++ b/cl/phase1/stages/data_availability_test.go @@ -0,0 +1,120 @@ +// Copyright 2026 The Erigon Authors +// This file is part of Erigon. +// +// Erigon is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// Erigon is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with Erigon. If not, see . + +package stages + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + "go.uber.org/mock/gomock" + + "github.com/erigontech/erigon/cl/clparams" + "github.com/erigontech/erigon/cl/cltypes" + dasmock "github.com/erigontech/erigon/cl/das/mock_services" + "github.com/erigontech/erigon/cl/phase1/forkchoice" + "github.com/erigontech/erigon/cl/utils/eth_clock" + "github.com/erigontech/erigon/common" +) + +func TestAcquireBlockDataAvailability(t *testing.T) { + tests := []struct { + name string + afterDownload bool + wantErr error + }{ + {name: "available after download", afterDownload: true}, + {name: "still unavailable", wantErr: forkchoice.ErrEIP7594ColumnDataNotAvailable}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + block := cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.FuluVersion) + block.Block.Slot = 1 + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + ctrl := gomock.NewController(t) + peerDas := dasmock.NewMockPeerDas(ctrl) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(false, nil) + peerDas.EXPECT().IsArchivedMode().Return(false) + peerDas.EXPECT().DownloadOnlyCustodyColumns(gomock.Any(), gomock.Any()).DoAndReturn( + func(_ context.Context, blocks []cltypes.ColumnSyncableSignedBlock) error { + require.Equal(t, []cltypes.ColumnSyncableSignedBlock{block}, blocks) + return nil + }, + ) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(tt.afterDownload, nil) + + err = acquireBlockDataAvailability(t.Context(), peerDas, block) + require.ErrorIs(t, err, tt.wantErr) + }) + } +} + +func TestAcquireBlockDataAvailabilityInArchiveMode(t *testing.T) { + block := cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.FuluVersion) + block.Block.Slot = 1 + block.Block.Body.BlobKzgCommitments.Append(&cltypes.KZGCommitment{}) + blockRoot, err := block.Block.HashSSZ() + require.NoError(t, err) + + ctrl := gomock.NewController(t) + peerDas := dasmock.NewMockPeerDas(ctrl) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(false, nil) + peerDas.EXPECT().IsArchivedMode().Return(true) + peerDas.EXPECT().DownloadColumnsAndRecoverBlobs(gomock.Any(), []cltypes.ColumnSyncableSignedBlock{block}).Return(nil) + peerDas.EXPECT().IsDataAvailable(block.Block.Slot, common.Hash(blockRoot)).Return(true, nil) + + require.NoError(t, acquireBlockDataAvailability(t.Context(), peerDas, block)) +} + +func TestRequiresRecentBlockDataAvailabilityUsesConfiguredFork(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + cfg.AltairForkEpoch = 0 + cfg.BellatrixForkEpoch = 0 + cfg.CapellaForkEpoch = 0 + cfg.DenebForkEpoch = 0 + cfg.ElectraForkEpoch = 0 + cfg.FuluForkEpoch = 2 + cfg.InitializeForkSchedule() + + tests := []struct { + name string + blockSlot uint64 + decodedVersion clparams.StateVersion + want bool + }{ + {name: "Fulu slot decoded as Electra", blockSlot: 2 * cfg.SlotsPerEpoch, decodedVersion: clparams.ElectraVersion, want: true}, + {name: "Electra slot decoded as Fulu", blockSlot: cfg.SlotsPerEpoch, decodedVersion: clparams.FuluVersion}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + clock := eth_clock.NewMockEthereumClock(gomock.NewController(t)) + clock.EXPECT().GetCurrentSlot().Return(tt.blockSlot) + block := cltypes.NewSignedBeaconBlock(&cfg, tt.decodedVersion) + block.Block.Slot = tt.blockSlot + + require.Equal(t, tt.want, requiresRecentBlockDataAvailability(&Cfg{ + beaconCfg: &cfg, + ethClock: clock, + }, block)) + }) + } +} diff --git a/cl/phase1/stages/forward_sync.go b/cl/phase1/stages/forward_sync.go index d9a22b0ae09..10a6624a6b3 100644 --- a/cl/phase1/stages/forward_sync.go +++ b/cl/phase1/stages/forward_sync.go @@ -62,14 +62,19 @@ func processDownloadedBlockBatches(ctx context.Context, logger log.Logger, cfg * return } + if err = acquireRecentBlockDataAvailability(ctx, cfg, block); err != nil { + if errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) { + logger.Trace("[Caplin] forward sync data not available", "blockSlot", block.Block.Slot) + return newHighestBlockProcessed, nil + } + return newHighestBlockProcessed, err + } + // Process the block - if err = processBlock(ctx, cfg, cfg.indiciesDB, block, false, true, false); err != nil { + if err = processBlock(ctx, cfg, cfg.indiciesDB, block, false, true, requiresRecentBlockDataAvailability(cfg, block)); err != nil { if errors.Is(err, forkchoice.ErrEIP4844DataNotAvailable) || errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) { logger.Trace("[Caplin] forward sync data not available", "blockSlot", block.Block.Slot, "err", err) - if newHighestBlockProcessed == 0 { - return 0, nil - } - return newHighestBlockProcessed - 1, nil + return newHighestBlockProcessed, nil } if errors.Is(err, forkchoice.ErrMissingSegment) { // Parent state not available — likely peer returned incomplete chain. diff --git a/cl/spectest/consensus_tests/fork_choice.go b/cl/spectest/consensus_tests/fork_choice.go index 705748d3df6..18f8e41218e 100644 --- a/cl/spectest/consensus_tests/fork_choice.go +++ b/cl/spectest/consensus_tests/fork_choice.go @@ -107,8 +107,8 @@ func (forkChoiceSpectestEngine) GetAssembledBlock(context.Context, []byte, clpar return nil, nil, nil, nil, nil } -func (forkChoiceSpectestEngine) GetBlobs(context.Context, []common.Hash, clparams.StateVersion) ([][]byte, [][][]byte, error) { - return nil, nil, nil +func (forkChoiceSpectestEngine) GetBlobs(_ context.Context, versionedHashes []common.Hash, _ clparams.StateVersion) ([][]byte, [][][]byte, error) { + return make([][]byte, len(versionedHashes)), make([][][]byte, len(versionedHashes)), nil } func (forkChoiceSpectestEngine) GetClientVersionV1(context.Context, *engine_types.ClientVersionV1) ([]engine_types.ClientVersionV1, error) { @@ -294,12 +294,9 @@ func (b *ForkChoice) Run(t *testing.T, root fs.FS, c spectest.TestCase) (err err ctx, cancel := context.WithCancel(context.Background()) t.Cleanup(cancel) // cancel PeerDas worker goroutines when the test finishes - anchorBlock, err := spectest.ReadAnchorBlock(root, c.Version(), "anchor_block.ssz_snappy") + _, err = spectest.ReadAnchorBlock(root, c.Version(), "anchor_block.ssz_snappy") require.NoError(t, err) - // TODO: what to do with anchor block ? - _ = anchorBlock - anchorState, err := spectest.ReadBeaconState(root, c.Version(), "anchor_state.ssz_snappy") require.NoError(t, err)