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)