diff --git a/cl/beacon/handler/epbs.go b/cl/beacon/handler/epbs.go index 84422149945..b5afd686eb2 100644 --- a/cl/beacon/handler/epbs.go +++ b/cl/beacon/handler/epbs.go @@ -797,7 +797,7 @@ func (a *ApiHandler) PostEthV1BeaconExecutionPayloadEnvelope(w http.ResponseWrit beaconhttp.NewEndpointError(http.StatusBadRequest, err).WriteTo(w) return } - if err := signedEnvelope.DecodeSSZ(octect, int(clparams.GloasVersion)); err != nil { + if err := signedEnvelope.DecodeSSZStrict(octect, int(clparams.GloasVersion)); err != nil { beaconhttp.NewEndpointError(http.StatusBadRequest, err).WriteTo(w) return } diff --git a/cl/cltypes/epbs_payload.go b/cl/cltypes/epbs_payload.go index 857b003e234..caaca97a714 100644 --- a/cl/cltypes/epbs_payload.go +++ b/cl/cltypes/epbs_payload.go @@ -20,6 +20,7 @@ import ( "bytes" "encoding/json" "errors" + "fmt" "github.com/erigontech/erigon/cl/clparams" "github.com/erigontech/erigon/cl/cltypes/solid" @@ -495,20 +496,25 @@ func (e *ExecutionPayloadEnvelope) EncodeSSZ(buf []byte) ([]byte, error) { } func (e *ExecutionPayloadEnvelope) DecodeSSZ(buf []byte, version int) error { + return e.decodeSSZ(buf, version, false) +} + +func (e *ExecutionPayloadEnvelope) DecodeSSZStrict(buf []byte, version int) error { + return e.decodeSSZ(buf, version, true) +} + +func (e *ExecutionPayloadEnvelope) decodeSSZ(buf []byte, version int, strict bool) error { if e.Payload == nil { e.Payload = NewEth1Block(clparams.StateVersion(version), e.beaconCfg) } if e.ExecutionRequests == nil { e.ExecutionRequests = NewExecutionRequestsWithVersion(e.beaconCfg, clparams.StateVersion(version)) } - return ssz2.UnmarshalSSZ( - buf, version, - e.Payload, - e.ExecutionRequests, - &e.BuilderIndex, - e.BeaconBlockRoot[:], - e.ParentBeaconBlockRoot[:], - ) + schema := []any{e.Payload, e.ExecutionRequests, &e.BuilderIndex, e.BeaconBlockRoot[:], e.ParentBeaconBlockRoot[:]} + if strict { + return ssz2.UnmarshalSSZStrict(buf, version, schema...) + } + return ssz2.UnmarshalSSZ(buf, version, schema...) } func (e *ExecutionPayloadEnvelope) EncodingSizeSSZ() int { @@ -546,6 +552,84 @@ type SignedExecutionPayloadEnvelope struct { beaconCfg *clparams.BeaconChainConfig } +// ValidateForConfig checks structural and protocol constraints before hashing an envelope. +func (s *SignedExecutionPayloadEnvelope) ValidateForConfig(cfg *clparams.BeaconChainConfig) error { + if s == nil { + return errors.New("nil execution payload envelope") + } + if cfg == nil { + return errors.New("nil beacon chain config") + } + if s.Message == nil { + return errors.New("nil execution payload envelope message") + } + payload := s.Message.Payload + if payload == nil { + return errors.New("execution payload envelope has nil payload") + } + if payload.Extra == nil { + return errors.New("execution payload envelope has nil extra data") + } + if err := payload.Extra.ValidateBounds(); err != nil { + return fmt.Errorf("invalid execution payload extra data: %w", err) + } + if payload.Transactions == nil { + return errors.New("execution payload envelope has nil transactions") + } + if payload.Withdrawals == nil { + return errors.New("execution payload envelope has nil withdrawals") + } + if err := payload.Withdrawals.ValidateBounds(int(cfg.MaxWithdrawalsPerPayload)); err != nil { + return fmt.Errorf("invalid execution payload withdrawals: %w", err) + } + if err := solid.RangeErr(payload.Withdrawals, func(i int, withdrawal *Withdrawal, _ int) error { + if withdrawal == nil { + return fmt.Errorf("nil withdrawal at index %d", i) + } + return nil + }); err != nil { + return err + } + if payload.BlockAccessList == nil { + return errors.New("execution payload envelope has nil block access list") + } + requests := s.Message.ExecutionRequests + if requests == nil { + return errors.New("execution payload envelope has nil execution requests") + } + if payload.Version() < clparams.GloasVersion { + return fmt.Errorf("execution payload version %d predates Gloas", payload.Version()) + } + if requests.Version() < clparams.GloasVersion { + return fmt.Errorf("execution requests version %d predates Gloas", requests.Version()) + } + if payload.Version() != requests.Version() { + return fmt.Errorf("payload and execution requests versions differ: %d != %d", payload.Version(), requests.Version()) + } + if err := requests.validateForConfig(cfg); err != nil { + return fmt.Errorf("invalid execution requests: %w", err) + } + return nil +} + +// ValidateForPersistence checks that the configured decoder can read the encoded envelope. +func (s *SignedExecutionPayloadEnvelope) ValidateForPersistence(cfg *clparams.BeaconChainConfig) error { + if err := s.ValidateForConfig(cfg); err != nil { + return err + } + payload := s.Message.Payload + if err := payload.Transactions.ValidateBounds(cfg.MaxTransactionsPerPayload, cfg.MaxBytesPerTransaction); err != nil { + return fmt.Errorf("transactions exceed decoder resource limit: %w", err) + } + if err := payload.BlockAccessList.ValidateBounds(cfg.MaxBytesPerTransaction); err != nil { + return fmt.Errorf("block access list exceeds decoder resource limit: %w", err) + } + if err := s.Message.ExecutionRequests.validateForPersistence(cfg); err != nil { + return fmt.Errorf("execution requests exceed decoder resource limit: %w", err) + } + return nil +} + func (s *SignedExecutionPayloadEnvelope) HashSSZ() ([32]byte, error) { return merkle_tree.HashTreeRoot(s.Message, s.Signature[:]) } @@ -559,9 +643,20 @@ func (s *SignedExecutionPayloadEnvelope) EncodeSSZ(buf []byte) ([]byte, error) { } func (s *SignedExecutionPayloadEnvelope) DecodeSSZ(buf []byte, version int) error { + return s.decodeSSZ(buf, version, false) +} + +func (s *SignedExecutionPayloadEnvelope) DecodeSSZStrict(buf []byte, version int) error { + return s.decodeSSZ(buf, version, true) +} + +func (s *SignedExecutionPayloadEnvelope) decodeSSZ(buf []byte, version int, strict bool) error { if s.Message == nil { s.Message = NewExecutionPayloadEnvelope(s.beaconCfg) } + if strict { + return ssz2.UnmarshalSSZStrict(buf, version, s.Message, s.Signature[:]) + } return ssz2.UnmarshalSSZ(buf, version, s.Message, s.Signature[:]) } diff --git a/cl/cltypes/epbs_payload_test.go b/cl/cltypes/epbs_payload_test.go index 165dbfa018b..4666f6543ad 100644 --- a/cl/cltypes/epbs_payload_test.go +++ b/cl/cltypes/epbs_payload_test.go @@ -1,17 +1,31 @@ package cltypes import ( + "encoding/binary" "encoding/json" "errors" "testing" "github.com/stretchr/testify/require" + "github.com/erigontech/erigon/cl/clparams" "github.com/erigontech/erigon/cl/cltypes/solid" "github.com/erigontech/erigon/common" "github.com/erigontech/erigon/common/ssz" ) +func TestExecutionRequestsStrictDecodeRejectsNonCanonicalOffset(t *testing.T) { + requests := NewExecutionRequestsWithVersion(&clparams.MainnetBeaconConfig, clparams.GloasVersion) + encoded, err := requests.EncodeSSZ(nil) + require.NoError(t, err) + firstOffset := binary.LittleEndian.Uint32(encoded) + binary.LittleEndian.PutUint32(encoded, firstOffset+1) + encoded = append(encoded[:firstOffset], append([]byte{0}, encoded[firstOffset:]...)...) + + decoded := NewExecutionRequestsWithVersion(&clparams.MainnetBeaconConfig, clparams.GloasVersion) + require.Error(t, decoded.DecodeSSZStrict(encoded, int(clparams.GloasVersion))) +} + func TestSignedExecutionPayloadEnvelopeCloneNilMessage(t *testing.T) { envelope := &SignedExecutionPayloadEnvelope{ Signature: common.Bytes96{1, 2, 3}, @@ -22,6 +36,84 @@ func TestSignedExecutionPayloadEnvelopeCloneNilMessage(t *testing.T) { require.Equal(t, envelope.Signature, cloned.Signature) } +func TestExecutionPayloadEnvelopeValidationSeparatesProtocolAndPersistenceBounds(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + cfg.MaxWithdrawalsPerPayload = 1 + cfg.MaxWithdrawalRequestsPerPayload = 1 + cfg.MaxConsolidationRequestsPerPayload = 1 + cfg.MaxBuilderDepositRequestsPerPayload = 1 + cfg.MaxBuilderExitRequestsPerPayload = 1 + cfg.MaxTransactionsPerPayload = 1 + cfg.MaxBytesPerTransaction = 1 + + for _, test := range []struct { + name string + mutate func(*SignedExecutionPayloadEnvelope) + }{ + {"payload withdrawals", func(e *SignedExecutionPayloadEnvelope) { + e.Message.Payload.Withdrawals.Append(&Withdrawal{}) + e.Message.Payload.Withdrawals.Append(&Withdrawal{}) + }}, + {"withdrawal requests", func(e *SignedExecutionPayloadEnvelope) { + e.Message.ExecutionRequests.Withdrawals.Append(&solid.WithdrawalRequest{}) + e.Message.ExecutionRequests.Withdrawals.Append(&solid.WithdrawalRequest{}) + }}, + {"consolidation requests", func(e *SignedExecutionPayloadEnvelope) { + e.Message.ExecutionRequests.Consolidations.Append(&solid.ConsolidationRequest{}) + e.Message.ExecutionRequests.Consolidations.Append(&solid.ConsolidationRequest{}) + }}, + {"builder deposit requests", func(e *SignedExecutionPayloadEnvelope) { + e.Message.ExecutionRequests.BuilderDeposits.Append(&solid.BuilderDepositRequest{}) + e.Message.ExecutionRequests.BuilderDeposits.Append(&solid.BuilderDepositRequest{}) + }}, + {"builder exit requests", func(e *SignedExecutionPayloadEnvelope) { + e.Message.ExecutionRequests.BuilderExits.Append(&solid.BuilderExitRequest{}) + e.Message.ExecutionRequests.BuilderExits.Append(&solid.BuilderExitRequest{}) + }}, + } { + t.Run(test.name, func(t *testing.T) { + envelope := validTestExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + test.mutate(envelope) + require.Error(t, envelope.ValidateForConfig(&cfg)) + }) + } + + envelope := validTestExecutionPayloadEnvelope(&cfg) + for range 16_385 { + envelope.Message.ExecutionRequests.Deposits.Append(&solid.DepositRequest{}) + } + require.NoError(t, envelope.ValidateForConfig(&cfg)) + require.Error(t, envelope.ValidateForPersistence(&cfg)) + + for _, test := range []struct { + name string + mutate func(*SignedExecutionPayloadEnvelope) + }{ + {"transactions", func(e *SignedExecutionPayloadEnvelope) { + e.Message.Payload.Transactions = solid.NewTransactionsSSZFromTransactions([][]byte{{1}, {2}}) + }}, + {"block access list", func(e *SignedExecutionPayloadEnvelope) { + require.NoError(t, e.Message.Payload.BlockAccessList.SetBytes([]byte{1, 2})) + }}, + } { + t.Run(test.name+" are resource bounded only", func(t *testing.T) { + envelope := validTestExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + test.mutate(envelope) + require.NoError(t, envelope.ValidateForConfig(&cfg)) + require.Error(t, envelope.ValidateForPersistence(&cfg)) + }) + } +} + +func validTestExecutionPayloadEnvelope(cfg *clparams.BeaconChainConfig) *SignedExecutionPayloadEnvelope { + message := NewExecutionPayloadEnvelope(cfg) + message.Payload.Extra = solid.NewExtraData() + message.Payload.Transactions = solid.NewTransactionsSSZFromTransactions(nil) + message.Payload.Withdrawals = solid.NewStaticListSSZ[*Withdrawal](int(cfg.MaxWithdrawalsPerPayload), 44) + message.Payload.BlockAccessList = solid.NewByteListSSZ(cfg.MaxBytesPerTransaction) + return &SignedExecutionPayloadEnvelope{Message: message} +} + func TestBuilderPendingPaymentSSZIncludesProposerIndex(t *testing.T) { payment := &BuilderPendingPayment{ Weight: 123, diff --git a/cl/cltypes/execution_requests.go b/cl/cltypes/execution_requests.go index 5ddfb690aad..0f624a09c37 100644 --- a/cl/cltypes/execution_requests.go +++ b/cl/cltypes/execution_requests.go @@ -53,6 +53,10 @@ func (e *ExecutionRequests) effectiveVersion() clparams.StateVersion { return e.version } +func (e *ExecutionRequests) Version() clparams.StateVersion { + return e.effectiveVersion() +} + func (e *ExecutionRequests) ensureLists() { if e.cfg == nil { panic("execution requests beacon config is nil") @@ -111,6 +115,14 @@ func (e *ExecutionRequests) EncodeSSZ(buf []byte) ([]byte, error) { } func (e *ExecutionRequests) DecodeSSZ(buf []byte, version int) error { + return e.decodeSSZ(buf, version, false) +} + +func (e *ExecutionRequests) DecodeSSZStrict(buf []byte, version int) error { + return e.decodeSSZ(buf, version, true) +} + +func (e *ExecutionRequests) decodeSSZ(buf []byte, version int, strict bool) error { decodedVersion := clparams.StateVersion(version) if (e.effectiveVersion() >= clparams.GloasVersion) != (decodedVersion >= clparams.GloasVersion) { e.Deposits = nil @@ -121,10 +133,14 @@ func (e *ExecutionRequests) DecodeSSZ(buf []byte, version int) error { } e.version = decodedVersion e.ensureLists() - if e.effectiveVersion() < clparams.GloasVersion { - return ssz2.UnmarshalSSZ(buf, version, e.Deposits, e.Withdrawals, e.Consolidations) + schema := []any{e.Deposits, e.Withdrawals, e.Consolidations} + if e.effectiveVersion() >= clparams.GloasVersion { + schema = append(schema, e.BuilderDeposits, e.BuilderExits) } - return ssz2.UnmarshalSSZ(buf, version, e.Deposits, e.Withdrawals, e.Consolidations, e.BuilderDeposits, e.BuilderExits) + if strict { + return ssz2.UnmarshalSSZStrict(buf, version, schema...) + } + return ssz2.UnmarshalSSZ(buf, version, schema...) } func (e *ExecutionRequests) Clone() clonable.Clonable { @@ -205,6 +221,68 @@ func (e *ExecutionRequests) Static() bool { return false } +func (e *ExecutionRequests) validateForConfig(cfg *clparams.BeaconChainConfig) error { + if e.Deposits == nil { + return fmt.Errorf("nil deposit requests") + } + if e.Withdrawals == nil { + return fmt.Errorf("nil withdrawal requests") + } + if e.Consolidations == nil { + return fmt.Errorf("nil consolidation requests") + } + if e.BuilderDeposits == nil { + return fmt.Errorf("nil builder deposit requests") + } + if e.BuilderExits == nil { + return fmt.Errorf("nil builder exit requests") + } + if err := e.Withdrawals.ValidateBounds(int(cfg.MaxWithdrawalRequestsPerPayload)); err != nil { + return fmt.Errorf("withdrawals: %w", err) + } + if err := e.Consolidations.ValidateBounds(int(cfg.MaxConsolidationRequestsPerPayload)); err != nil { + return fmt.Errorf("consolidations: %w", err) + } + if err := e.BuilderDeposits.ValidateBounds(int(cfg.MaxBuilderDepositRequestsPerPayload)); err != nil { + return fmt.Errorf("builder deposits: %w", err) + } + if err := e.BuilderExits.ValidateBounds(int(cfg.MaxBuilderExitRequestsPerPayload)); err != nil { + return fmt.Errorf("builder exits: %w", err) + } + if err := solid.RangeErr(e.Deposits, rejectNilRequest("deposit", func(request *solid.DepositRequest) bool { return request == nil })); err != nil { + return err + } + if err := solid.RangeErr(e.Withdrawals, rejectNilRequest("withdrawal", func(request *solid.WithdrawalRequest) bool { return request == nil })); err != nil { + return err + } + if err := solid.RangeErr(e.Consolidations, rejectNilRequest("consolidation", func(request *solid.ConsolidationRequest) bool { return request == nil })); err != nil { + return err + } + if err := solid.RangeErr(e.BuilderDeposits, rejectNilRequest("builder deposit", func(request *solid.BuilderDepositRequest) bool { return request == nil })); err != nil { + return err + } + return solid.RangeErr(e.BuilderExits, rejectNilRequest("builder exit", func(request *solid.BuilderExitRequest) bool { return request == nil })) +} + +func (e *ExecutionRequests) validateForPersistence(cfg *clparams.BeaconChainConfig) error { + if err := e.validateForConfig(cfg); err != nil { + return err + } + if err := e.Deposits.ValidateProgressiveDecodeBounds(int(cfg.MaxDepositRequestsPerPayload)); err != nil { + return fmt.Errorf("deposits exceed decoder resource limit: %w", err) + } + return nil +} + +func rejectNilRequest[T solid.EncodableHashableSSZ](name string, isNil func(T) bool) func(int, T, int) error { + return func(i int, request T, _ int) error { + if isNil(request) { + return fmt.Errorf("nil %s request at index %d", name, i) + } + return nil + } +} + func (e *ExecutionRequests) UnmarshalJSON(b []byte) error { e.ensureLists() newDeposits := solid.NewStaticListSSZ[*solid.DepositRequest](int(e.cfg.MaxDepositRequestsPerPayload), solid.SizeDepositRequest) diff --git a/cl/cltypes/solid/byte_list.go b/cl/cltypes/solid/byte_list.go index 1e9356471d1..04179dae488 100644 --- a/cl/cltypes/solid/byte_list.go +++ b/cl/cltypes/solid/byte_list.go @@ -149,3 +149,10 @@ func (b *ByteListSSZ) SetBytes(buf []byte) error { func (b *ByteListSSZ) Len() int { return len(b.data) } + +func (b *ByteListSSZ) ValidateBounds(limit uint64) error { + if uint64(len(b.data)) > limit { + return fmt.Errorf("data length %d exceeds limit %d", len(b.data), limit) + } + return nil +} diff --git a/cl/cltypes/solid/extra_data.go b/cl/cltypes/solid/extra_data.go index f53e0bcfeb1..8c3fa185751 100644 --- a/cl/cltypes/solid/extra_data.go +++ b/cl/cltypes/solid/extra_data.go @@ -124,3 +124,10 @@ func (e *ExtraData) SetBytes(buf []byte) { e.l = len(e.data) } } + +func (e *ExtraData) ValidateBounds() error { + if e.l > maxExtraDataBytes { + return fmt.Errorf("extra data length %d exceeds limit %d", e.l, maxExtraDataBytes) + } + return nil +} diff --git a/cl/cltypes/solid/list_ssz.go b/cl/cltypes/solid/list_ssz.go index f8f68985136..52e899c0a72 100644 --- a/cl/cltypes/solid/list_ssz.go +++ b/cl/cltypes/solid/list_ssz.go @@ -19,6 +19,7 @@ package solid import ( "bytes" "encoding/json" + "fmt" "github.com/erigontech/erigon/cl/merkle_tree" "github.com/erigontech/erigon/common" @@ -233,6 +234,17 @@ func (l *ListSSZ[T]) Len() int { return len(l.list) } +func (l *ListSSZ[T]) ValidateBounds(limit int) error { + if len(l.list) > limit { + return fmt.Errorf("list has %d elements, max %d", len(l.list), limit) + } + return nil +} + +func (l *ListSSZ[T]) ValidateProgressiveDecodeBounds(configuredLimit int) error { + return l.ValidateBounds(progressiveDecodeLimit(configuredLimit)) +} + func (l *ListSSZ[T]) Set(index int, value T) { l.list[index] = value l.root = common.Hash{} diff --git a/cl/cltypes/solid/transactions.go b/cl/cltypes/solid/transactions.go index 5eaeeba8096..75cfc6e7ae5 100644 --- a/cl/cltypes/solid/transactions.go +++ b/cl/cltypes/solid/transactions.go @@ -138,6 +138,24 @@ func (t *TransactionsSSZ) EncodeSSZ(buf []byte) (dst []byte, err error) { return dst, nil } +func (t *TransactionsSSZ) ValidateBounds(maxTransactions, maxBytesPerTransaction uint64) error { + if maxTransactions == 0 { + maxTransactions = clparams.MaxTransactionsPerPayloadDefault + } + if maxBytesPerTransaction == 0 { + maxBytesPerTransaction = clparams.MaxBytesPerTransactionDefault + } + if uint64(len(t.underlying)) > maxTransactions { + return fmt.Errorf("too many transactions: got %d, max %d", len(t.underlying), maxTransactions) + } + for i, transaction := range t.underlying { + if uint64(len(transaction)) > maxBytesPerTransaction { + return fmt.Errorf("transaction %d is too large: got %d bytes, max %d", i, len(transaction), maxBytesPerTransaction) + } + } + return nil +} + func (t *TransactionsSSZ) HashSSZ() ([32]byte, error) { var err error if t.root != (common.Hash{}) { diff --git a/cl/phase1/core/checkpoint_sync/checkpoint_sync_test.go b/cl/phase1/core/checkpoint_sync/checkpoint_sync_test.go index b7b0c3397d7..d4b36114ece 100644 --- a/cl/phase1/core/checkpoint_sync/checkpoint_sync_test.go +++ b/cl/phase1/core/checkpoint_sync/checkpoint_sync_test.go @@ -174,6 +174,17 @@ func TestNormalizeCheckpointURL(t *testing.T) { } } +func TestRemoteCheckpointSyncRejectsOversizedEnvelopeResponse(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + _, _ = w.Write(make([]byte, clparams.MaxChunkSize+1)) + })) + defer server.Close() + syncer := &RemoteCheckpointSync{&clparams.MainnetBeaconConfig, chainspec.MainnetChainID, time.Second} + + _, err := syncer.fetchEnvelope(context.Background(), server.URL+beaconStatePath) + require.ErrorContains(t, err, "too large") +} + func TestRemoteCheckpointSyncRejectsHTML(t *testing.T) { mockHTMLServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "text/html; charset=utf-8") diff --git a/cl/phase1/core/checkpoint_sync/remote_checkpoint_sync.go b/cl/phase1/core/checkpoint_sync/remote_checkpoint_sync.go index 9e61b409d8a..3f2b921c8de 100644 --- a/cl/phase1/core/checkpoint_sync/remote_checkpoint_sync.go +++ b/cl/phase1/core/checkpoint_sync/remote_checkpoint_sync.go @@ -181,7 +181,7 @@ func (r *RemoteCheckpointSync) fetchEnvelope(ctx context.Context, stateURI strin return nil, fmt.Errorf("finalized envelope fetch failed, status %d", resp.StatusCode) } - marshaled, err := io.ReadAll(resp.Body) + marshaled, err := readEnvelopeHTTPBody(resp.Body) if err != nil { return nil, fmt.Errorf("finalized envelope read failed: %w", err) } @@ -189,13 +189,24 @@ func (r *RemoteCheckpointSync) fetchEnvelope(ctx context.Context, stateURI strin envelope := &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(r.beaconConfig), } - if err := envelope.DecodeSSZ(marshaled, int(clparams.GloasVersion)); err != nil { + if err := envelope.DecodeSSZStrict(marshaled, int(clparams.GloasVersion)); err != nil { return nil, fmt.Errorf("finalized envelope decode failed: %w", err) } log.Info("[Checkpoint Sync] Finalized envelope retrieved", "beaconBlockRoot", envelope.Message.BeaconBlockRoot) return envelope, nil } +func readEnvelopeHTTPBody(r io.Reader) ([]byte, error) { + body, err := io.ReadAll(io.LimitReader(r, int64(clparams.MaxChunkSize)+1)) + if err != nil { + return nil, err + } + if uint64(len(body)) > clparams.MaxChunkSize { + return nil, fmt.Errorf("execution payload envelope response too large: max %d bytes", clparams.MaxChunkSize) + } + return body, nil +} + // normalizeCheckpointURL ensures the URL includes the beacon state API path. // Users often provide just the base URL (e.g. https://checkpoint-sync.example.io) // which serves an HTML landing page instead of SSZ data. diff --git a/cl/phase1/forkchoice/fork_graph/fork_graph_disk.go b/cl/phase1/forkchoice/fork_graph/fork_graph_disk.go index 065b0e5f944..790d983b200 100644 --- a/cl/phase1/forkchoice/fork_graph/fork_graph_disk.go +++ b/cl/phase1/forkchoice/fork_graph/fork_graph_disk.go @@ -19,6 +19,7 @@ package fork_graph import ( "errors" "fmt" + "io/fs" "slices" "sync" "sync/atomic" @@ -128,7 +129,8 @@ type forkGraphDisk struct { lightClientUpdates sync.Map // period -> lightclientupdate // in-memory cache of block roots that have envelopes on disk [Optimization for Gloas:EIP7732] - envelopeExists sync.Map // common.Hash -> struct{} + envelopeExists sync.Map // common.Hash -> struct{} + invalidEnvelopes sync.Map // common.Hash -> struct{} // reusable buffers sszBuffer []byte @@ -199,7 +201,6 @@ func NewForkGraphDisk(anchorState *state.CachingBeaconState, syncedData synced_d f.lowestAvailableBlock.Store(anchorState.Slot()) f.headers.Store(common.Hash(anchorRoot), &anchorHeader) f.sszBuffer = make([]byte, 0, (anchorState.EncodingSizeSSZ()*3)/2) - f.DumpBeaconStateOnDisk(anchorRoot, anchorState, true) // preallocate buffer return f @@ -607,16 +608,25 @@ func (f *forkGraphDisk) Prune(pruneSlot uint64) (err error) { } for _, root := range oldRoots { f.badBlocks.Delete(root) - f.blocks.Delete(root) f.lightclientBootstraps.Delete(root) f.currentJustifiedCheckpoints.Delete(root) f.finalizedCheckpoints.Delete(root) f.headers.Delete(root) f.blockRewards.Delete(root) + } + for _, root := range oldRoots { + f.stateDumpLock.Lock() + f.blocks.Delete(root) f.fs.Remove(getBeaconStateFilename(root)) - // [New in Gloas:EIP7732] Also remove envelope files f.envelopeExists.Delete(root) - f.fs.Remove(getEnvelopeFilename(root)) + f.invalidEnvelopes.Delete(root) + f.stateDumpLock.Unlock() + if removeErr := f.fs.Remove(getEnvelopeFilename(root)); removeErr != nil && !errors.Is(removeErr, fs.ErrNotExist) { + err = errors.Join(err, fmt.Errorf("remove envelope for root %x: %w", root, removeErr)) + } + if removeErr := f.fs.Remove(getEnvelopeTempFilename(root)); removeErr != nil && !errors.Is(removeErr, fs.ErrNotExist) { + err = errors.Join(err, fmt.Errorf("remove envelope temp for root %x: %w", root, removeErr)) + } } log.Debug("Pruned old blocks", "pruneSlot", pruneSlot) return diff --git a/cl/phase1/forkchoice/fork_graph/fork_graph_disk_fs.go b/cl/phase1/forkchoice/fork_graph/fork_graph_disk_fs.go index f8b728db142..ad0973f8b89 100644 --- a/cl/phase1/forkchoice/fork_graph/fork_graph_disk_fs.go +++ b/cl/phase1/forkchoice/fork_graph/fork_graph_disk_fs.go @@ -17,7 +17,9 @@ package fork_graph import ( + "bytes" "encoding/binary" + "errors" "fmt" "io" "os" @@ -27,6 +29,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/phase1/core/state" "github.com/erigontech/erigon/common" "github.com/erigontech/erigon/common/log/v3" @@ -48,6 +51,10 @@ func getEnvelopeFilename(blockRoot common.Hash) string { return fmt.Sprintf("%x.envelope.snappy_ssz", blockRoot) } +func getEnvelopeTempFilename(blockRoot common.Hash) string { + return getEnvelopeFilename(blockRoot) + ".tmp" +} + func (f *forkGraphDisk) readBeaconStateFromDisk(blockRoot common.Hash) (bs *state.CachingBeaconState, err error) { var file afero.File f.stateDumpLock.Lock() @@ -186,87 +193,170 @@ func (f *forkGraphDisk) DumpBeaconStateOnDisk(blockRoot common.Hash, bs *state.C } // HasEnvelope checks if an envelope exists for the given block root. -// Uses an in-memory cache populated by DumpEnvelopeOnDisk to avoid repeated disk stats. +// Only envelopes successfully persisted or validated by this process are reported. // [New in Gloas:EIP7732] func (f *forkGraphDisk) HasEnvelope(blockRoot common.Hash) bool { - // Fast path: check in-memory cache - if _, ok := f.envelopeExists.Load(blockRoot); ok { - return true + if _, invalid := f.invalidEnvelopes.Load(blockRoot); invalid { + return false + } + if !f.knowsBlockRoot(blockRoot) { + return false } - // Slow path: fall back to disk and populate cache on hit - exists, err := afero.Exists(f.fs, getEnvelopeFilename(blockRoot)) - if err == nil && exists { - f.envelopeExists.Store(blockRoot, struct{}{}) + _, ok := f.envelopeExists.Load(blockRoot) + return ok +} + +func (f *forkGraphDisk) knowsBlockRoot(blockRoot common.Hash) bool { + if blockRoot == f.anchorRoot { return true } - return false + _, ok := f.blocks.Load(blockRoot) + return ok } // ReadEnvelopeFromDisk reads an execution payload envelope from disk. // [New in Gloas:EIP7732] func (f *forkGraphDisk) ReadEnvelopeFromDisk(blockRoot common.Hash) (envelope *cltypes.SignedExecutionPayloadEnvelope, err error) { - var file afero.File f.stateDumpLock.Lock() defer f.stateDumpLock.Unlock() + return f.readEnvelopeFromDiskLocked(blockRoot) +} - file, err = f.fs.Open(getEnvelopeFilename(blockRoot)) +func (f *forkGraphDisk) readEnvelopeFromDiskLocked(blockRoot common.Hash) (envelope *cltypes.SignedExecutionPayloadEnvelope, err error) { + var file afero.File + var corrupt bool + _, wasCached := f.envelopeExists.Load(blockRoot) + if _, invalid := f.invalidEnvelopes.Load(blockRoot); invalid { + return nil, fmt.Errorf("cannot read known invalid envelope for root %x", blockRoot) + } + if !f.knowsBlockRoot(blockRoot) { + return nil, fmt.Errorf("cannot read envelope for unknown block root %x", blockRoot) + } + + filename := getEnvelopeFilename(blockRoot) + file, err = f.fs.Open(filename) if err != nil { + if !wasCached || errors.Is(err, os.ErrNotExist) { + f.envelopeExists.Delete(blockRoot) + } return } - defer file.Close() + defer func() { + if closeErr := file.Close(); closeErr != nil { + log.Warn("failed to close envelope after read", "root", blockRoot, "err", closeErr) + } + if corrupt { + f.invalidEnvelopes.Store(blockRoot, struct{}{}) + f.envelopeExists.Delete(blockRoot) + } else if err != nil && !wasCached { + f.envelopeExists.Delete(blockRoot) + } + }() - if f.sszSnappyReader == nil { - f.sszSnappyReader = snappy.NewReader(file) - } else { - f.sszSnappyReader.Reset(file) + readTracker := &envelopeReadTracker{Reader: file} + snappyReader := snappy.NewReader(readTracker) + versionBytes := []byte{0} + if _, err = io.ReadFull(snappyReader, versionBytes); err != nil { + corrupt = isCorruptEnvelopeReadError(err, readTracker.err) + return nil, fmt.Errorf("failed to read envelope version: %w, root: %x", err, blockRoot) + } + version := clparams.StateVersion(versionBytes[0]) + if version < clparams.GloasVersion { + corrupt = true + return nil, fmt.Errorf("corrupt envelope file: version %d predates Gloas, root: %x", version, blockRoot) } // Read the length lengthBytes := make([]byte, 8) - var n int - n, err = io.ReadFull(f.sszSnappyReader, lengthBytes) + _, err = io.ReadFull(snappyReader, lengthBytes) if err != nil { + corrupt = isCorruptEnvelopeReadError(err, readTracker.err) return nil, fmt.Errorf("failed to read length: %w, root: %x", err, blockRoot) } - if n != 8 { - return nil, fmt.Errorf("failed to read length: %d, want 8, root: %x", n, blockRoot) - } - envelopeLength := binary.BigEndian.Uint64(lengthBytes) - if envelopeLength > maxSSZObjectSize { - return nil, fmt.Errorf("corrupt envelope file: length %d exceeds max %d, root: %x", envelopeLength, maxSSZObjectSize, blockRoot) + if envelopeLength > clparams.MaxChunkSize { + corrupt = true + return nil, fmt.Errorf("corrupt envelope file: length %d exceeds max %d, root: %x", envelopeLength, clparams.MaxChunkSize, blockRoot) } if envelopeLength > uint64(cap(f.sszBuffer)) { f.sszBuffer = make([]byte, envelopeLength) } else { f.sszBuffer = f.sszBuffer[:envelopeLength] } - n, err = io.ReadFull(f.sszSnappyReader, f.sszBuffer) + n, err := io.ReadFull(snappyReader, f.sszBuffer) if err != nil { + corrupt = isCorruptEnvelopeReadError(err, readTracker.err) return nil, fmt.Errorf("failed to read snappy buffer: %w, root: %x", err, blockRoot) } - f.sszBuffer = f.sszBuffer[:n] - envelope = &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(f.beaconCfg), } - if err = envelope.DecodeSSZ(f.sszBuffer, int(clparams.GloasVersion)); err != nil { + if err = envelope.DecodeSSZStrict(f.sszBuffer, int(version)); err != nil { + corrupt = true return nil, fmt.Errorf("failed to decode envelope: %w, root: %x, len: %d", err, blockRoot, n) } + if err = envelope.ValidateForPersistence(f.beaconCfg); err != nil { + corrupt = true + return nil, fmt.Errorf("invalid persisted envelope: %w, root: %x", err, blockRoot) + } + if envelope.Message.BeaconBlockRoot != blockRoot { + corrupt = true + return nil, fmt.Errorf("corrupt envelope file: embedded root %x does not match filename root %x", envelope.Message.BeaconBlockRoot, blockRoot) + } + transactions := envelope.Message.Payload.Transactions.UnderlyngReference() + ownedTransactions := make([][]byte, len(transactions)) + for i, transaction := range transactions { + ownedTransactions[i] = bytes.Clone(transaction) + } + // TransactionsSSZ decode aliases its input, so detach it before reusing the shared buffer. + envelope.Message.Payload.Transactions = solid.NewTransactionsSSZFromTransactions(ownedTransactions) + f.envelopeExists.Store(blockRoot, struct{}{}) + f.invalidEnvelopes.Delete(blockRoot) return } +type envelopeReadTracker struct { + io.Reader + err error +} + +func (r *envelopeReadTracker) Read(p []byte) (int, error) { + n, err := r.Reader.Read(p) + if err != nil && !errors.Is(err, io.EOF) { + r.err = err + } + return n, err +} + +func isCorruptEnvelopeReadError(err, sourceErr error) bool { + return sourceErr == nil || !errors.Is(err, sourceErr) +} + // DumpEnvelopeOnDisk dumps an execution payload envelope to disk. // [New in Gloas:EIP7732] func (f *forkGraphDisk) DumpEnvelopeOnDisk(blockRoot common.Hash, envelope *cltypes.SignedExecutionPayloadEnvelope) (err error) { + if validateErr := envelope.ValidateForPersistence(f.beaconCfg); validateErr != nil { + return fmt.Errorf("cannot persist invalid envelope: %w", validateErr) + } + if envelope.Message.BeaconBlockRoot != blockRoot { + return fmt.Errorf("cannot persist envelope for root %x with embedded root %x", blockRoot, envelope.Message.BeaconBlockRoot) + } + envelopeSize := envelope.EncodingSizeSSZ() + if envelopeSize < 0 || uint64(envelopeSize) > clparams.MaxChunkSize { + return fmt.Errorf("cannot persist envelope: length %d exceeds max %d", envelopeSize, clparams.MaxChunkSize) + } f.stateDumpLock.Lock() defer f.stateDumpLock.Unlock() + if !f.knowsBlockRoot(blockRoot) { + return fmt.Errorf("cannot persist envelope for unknown block root %x", blockRoot) + } // Populate in-memory cache on successful write defer func() { if err == nil { f.envelopeExists.Store(blockRoot, struct{}{}) + f.invalidEnvelopes.Delete(blockRoot) } }() @@ -275,12 +365,25 @@ func (f *forkGraphDisk) DumpEnvelopeOnDisk(blockRoot common.Hash, envelope *clty if err != nil { return } + if uint64(len(f.sszBuffer)) > clparams.MaxChunkSize { + return fmt.Errorf("cannot persist envelope: length %d exceeds max %d", len(f.sszBuffer), clparams.MaxChunkSize) + } - dumpedFile, err := f.fs.OpenFile(getEnvelopeFilename(blockRoot), os.O_TRUNC|os.O_CREATE|os.O_RDWR, 0o755) + filename := getEnvelopeFilename(blockRoot) + tempFilename := getEnvelopeTempFilename(blockRoot) + dumpedFile, err := f.fs.OpenFile(tempFilename, os.O_TRUNC|os.O_CREATE|os.O_RDWR, 0o644) if err != nil { return err } - defer dumpedFile.Close() + closed := false + defer func() { + if !closed { + _ = dumpedFile.Close() + } + if err != nil { + _ = f.fs.Remove(tempFilename) + } + }() if f.sszSnappyWriter == nil { f.sszSnappyWriter = snappy.NewBufferedWriter(dumpedFile) @@ -291,6 +394,10 @@ func (f *forkGraphDisk) DumpEnvelopeOnDisk(blockRoot common.Hash, envelope *clty // Write the length length := make([]byte, 8) binary.BigEndian.PutUint64(length, uint64(len(f.sszBuffer))) + if _, err := f.sszSnappyWriter.Write([]byte{byte(envelope.Message.Payload.Version())}); err != nil { + log.Error("failed to write envelope version", "err", err) + return err + } if _, err := f.sszSnappyWriter.Write(length); err != nil { log.Error("failed to write length", "err", err) return err @@ -309,6 +416,14 @@ func (f *forkGraphDisk) DumpEnvelopeOnDisk(blockRoot common.Hash, envelope *clty log.Error("failed to sync dumped file", "err", err) return } + err = dumpedFile.Close() + closed = true + if err != nil { + return + } + if err = f.fs.Rename(tempFilename, filename); err != nil { + return + } return } diff --git a/cl/phase1/forkchoice/fork_graph/fork_graph_test.go b/cl/phase1/forkchoice/fork_graph/fork_graph_test.go index 55cbbd11d83..c01352ebf80 100644 --- a/cl/phase1/forkchoice/fork_graph/fork_graph_test.go +++ b/cl/phase1/forkchoice/fork_graph/fork_graph_test.go @@ -17,15 +17,25 @@ package fork_graph import ( + "bytes" _ "embed" + "encoding/binary" + "errors" + "os" + "strings" + "sync" + "sync/atomic" "testing" + "time" "github.com/erigontech/erigon/cl/beacon/beacon_router_configuration" "github.com/erigontech/erigon/cl/phase1/core/state" + "github.com/golang/snappy" "github.com/spf13/afero" "github.com/erigontech/erigon/cl/clparams" "github.com/erigontech/erigon/cl/cltypes" + "github.com/erigontech/erigon/cl/cltypes/solid" "github.com/erigontech/erigon/cl/utils" "github.com/erigontech/erigon/common" "github.com/stretchr/testify/require" @@ -40,6 +50,163 @@ var block2 []byte //go:embed test_data/anchor_state.ssz_snappy var anchor []byte +var errTestEnvelopeIO = errors.New("test envelope I/O error") + +type envelopeCloseErrorFile struct{ afero.File } + +func (f envelopeCloseErrorFile) Close() error { + _ = f.File.Close() + return errTestEnvelopeIO +} + +type envelopeCloseErrorFs struct{ afero.Fs } + +func (f envelopeCloseErrorFs) Open(name string) (afero.File, error) { + file, err := f.Fs.Open(name) + if err != nil { + return nil, err + } + return envelopeCloseErrorFile{File: file}, nil +} + +type envelopeReadErrorFile struct{ afero.File } + +func (envelopeReadErrorFile) Read([]byte) (int, error) { return 0, errTestEnvelopeIO } + +type envelopeReadErrorFs struct{ afero.Fs } + +func (f envelopeReadErrorFs) Open(name string) (afero.File, error) { + file, err := f.Fs.Open(name) + if err != nil { + return nil, err + } + return envelopeReadErrorFile{File: file}, nil +} + +type envelopeOpenErrorFs struct{ afero.Fs } + +func (f envelopeOpenErrorFs) Open(name string) (afero.File, error) { + if strings.HasSuffix(name, ".envelope.snappy_ssz") { + return nil, errTestEnvelopeIO + } + return f.Fs.Open(name) +} + +type envelopeOpenCountingFs struct { + afero.Fs + opens atomic.Int32 +} + +func (f *envelopeOpenCountingFs) Open(name string) (afero.File, error) { + if strings.HasSuffix(name, ".envelope.snappy_ssz") { + f.opens.Add(1) + } + return f.Fs.Open(name) +} + +type envelopeWriteFailureFs struct { + afero.Fs + stage string + closes *atomic.Int32 +} + +func (f envelopeWriteFailureFs) OpenFile(name string, flag int, perm os.FileMode) (afero.File, error) { + if f.stage == "open" && strings.HasSuffix(name, ".tmp") { + return nil, errTestEnvelopeIO + } + file, err := f.Fs.OpenFile(name, flag, perm) + if err != nil { + return nil, err + } + return envelopeWriteFailureFile{File: file, stage: f.stage, closes: f.closes}, nil +} + +func (f envelopeWriteFailureFs) Rename(oldname, newname string) error { + if f.stage == "rename" { + return errTestEnvelopeIO + } + return f.Fs.Rename(oldname, newname) +} + +type envelopeWriteFailureFile struct { + afero.File + stage string + closes *atomic.Int32 +} + +type envelopeBlockingRenameFs struct { + afero.Fs + reached chan struct{} + release chan struct{} + once sync.Once +} + +type envelopeRemoveFailureFs struct { + afero.Fs + suffix string +} + +type envelopeBlockingPruneFs struct { + afero.Fs + firstReached chan struct{} + releaseFirst chan struct{} + secondReached chan struct{} + releaseSecond chan struct{} + removes atomic.Int32 +} + +func (f envelopeRemoveFailureFs) Remove(name string) error { + if strings.HasSuffix(name, f.suffix) { + return errTestEnvelopeIO + } + return f.Fs.Remove(name) +} + +func (f *envelopeBlockingPruneFs) Remove(name string) error { + switch f.removes.Add(1) { + case 1: + close(f.firstReached) + <-f.releaseFirst + case 2: + close(f.secondReached) + <-f.releaseSecond + } + return f.Fs.Remove(name) +} + +func (f *envelopeBlockingRenameFs) Rename(oldname, newname string) error { + if strings.HasSuffix(oldname, ".tmp") { + f.once.Do(func() { close(f.reached) }) + <-f.release + } + return f.Fs.Rename(oldname, newname) +} + +func (f envelopeWriteFailureFile) Write(p []byte) (int, error) { + if f.stage == "write" { + return 0, errTestEnvelopeIO + } + return f.File.Write(p) +} + +func (f envelopeWriteFailureFile) Sync() error { + if f.stage == "sync" { + return errTestEnvelopeIO + } + return f.File.Sync() +} + +func (f envelopeWriteFailureFile) Close() error { + if f.closes != nil { + f.closes.Add(1) + } + if f.stage == "close" { + _ = f.File.Close() + return errTestEnvelopeIO + } + return f.File.Close() +} + func TestForkGraphInDisk(t *testing.T) { blockA, blockB, blockC := cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.DenebVersion), cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.DenebVersion), @@ -87,3 +254,677 @@ func TestPruneKeepsLowestAvailableBlockMonotonic(t *testing.T) { require.NoError(t, f.Prune(120)) require.Equal(t, uint64(151), f.LowestAvailableSlot()) } + +func TestReadEnvelopeMarksCorruptFileInvalidWithoutDeletingIt(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, afero.WriteFile(fs, getEnvelopeFilename(root), []byte("truncated"), 0o644)) + + _, err := f.ReadEnvelopeFromDisk(root) + require.Error(t, err) + require.False(t, f.HasEnvelope(root)) + exists, existsErr := afero.Exists(fs, getEnvelopeFilename(root)) + require.NoError(t, existsErr) + require.True(t, exists) +} + +func TestHasEnvelopeDoesNotTrustUnvalidatedDiskFile(t *testing.T) { + baseFs := afero.NewMemMapFs() + root := common.HexToHash("0x1234") + writer := &forkGraphDisk{fs: baseFs, beaconCfg: &clparams.MainnetBeaconConfig} + addEnvelopeTestBlock(writer, root, 1) + require.NoError(t, writer.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1}))) + + fs := &envelopeOpenCountingFs{Fs: baseFs} + restarted := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + addEnvelopeTestBlock(restarted, root, 1) + restarted.stateDumpLock.Lock() + result := make(chan bool, 1) + go func() { result <- restarted.HasEnvelope(root) }() + select { + case hasEnvelope := <-result: + require.False(t, hasEnvelope) + case <-time.After(time.Second): + restarted.stateDumpLock.Unlock() + t.Fatal("HasEnvelope waited for state disk I/O") + } + restarted.stateDumpLock.Unlock() + require.Zero(t, fs.opens.Load()) + + _, err := restarted.ReadEnvelopeFromDisk(root) + require.NoError(t, err) + require.Equal(t, int32(1), fs.opens.Load()) + require.True(t, restarted.HasEnvelope(root)) +} + +func TestReadEnvelopeRemovesUnsupportedSnappyFrames(t *testing.T) { + streamIdentifier := []byte{0xff, 0x06, 0x00, 0x00, 's', 'N', 'a', 'P', 'p', 'Y'} + frame := append(append([]byte{}, streamIdentifier...), 0x02, 0x00, 0x00, 0x00) + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, afero.WriteFile(fs, getEnvelopeFilename(root), frame, 0o644)) + + _, err := f.ReadEnvelopeFromDisk(root) + require.ErrorIs(t, err, snappy.ErrUnsupported) + require.False(t, f.HasEnvelope(root)) +} + +func TestEnvelopeReadClassifiesSnappyStructuralErrors(t *testing.T) { + require.True(t, isCorruptEnvelopeReadError(snappy.ErrCorrupt, nil)) + require.True(t, isCorruptEnvelopeReadError(snappy.ErrUnsupported, nil)) + require.True(t, isCorruptEnvelopeReadError(snappy.ErrTooLarge, nil)) + require.False(t, isCorruptEnvelopeReadError(errTestEnvelopeIO, errTestEnvelopeIO)) +} + +func TestDumpEnvelopeAtomicallyPersistsReadableFile(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1, 2, 3}))) + tempExists, err := afero.Exists(fs, getEnvelopeFilename(root)+".tmp") + require.NoError(t, err) + require.False(t, tempExists) + persisted, err := f.ReadEnvelopeFromDisk(root) + require.NoError(t, err) + require.Equal(t, root, persisted.Message.BeaconBlockRoot) +} + +func TestDumpEnvelopePreservesPostGloasVersion(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + version := clparams.GloasVersion + 1 + envelope := testEnvelopeWithVersion(root, []byte{1, 2, 3}, version) + + require.NoError(t, f.DumpEnvelopeOnDisk(root, envelope)) + persisted, err := f.ReadEnvelopeFromDisk(root) + require.NoError(t, err) + require.Equal(t, version, persisted.Message.Payload.Version()) + require.Equal(t, version, persisted.Message.ExecutionRequests.Version()) +} + +func TestDumpEnvelopeRejectsMismatchedNestedVersions(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := testEnvelopeWithVersion(root, []byte{1}, clparams.GloasVersion+1) + envelope.Message.ExecutionRequests = cltypes.NewExecutionRequestsWithVersion(&clparams.MainnetBeaconConfig, clparams.GloasVersion) + + require.ErrorContains(t, f.DumpEnvelopeOnDisk(root, envelope), "versions differ") + require.False(t, f.HasEnvelope(root)) +} + +func TestDumpEnvelopeRejectsMismatchedRoot(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + rootA := common.HexToHash("0xa") + rootB := common.HexToHash("0xb") + addEnvelopeTestBlock(f, rootA, 1) + + require.Error(t, f.DumpEnvelopeOnDisk(rootA, testEnvelopeWithTransaction(rootB, []byte{1}))) + require.False(t, f.HasEnvelope(rootA)) +} + +func TestReadEnvelopeRejectsMismatchedRoot(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + rootA := common.HexToHash("0xa") + rootB := common.HexToHash("0xb") + addEnvelopeTestBlock(f, rootA, 1) + addEnvelopeTestBlock(f, rootB, 2) + require.NoError(t, f.DumpEnvelopeOnDisk(rootA, testEnvelopeWithTransaction(rootA, []byte{1}))) + require.NoError(t, fs.Rename(getEnvelopeFilename(rootA), getEnvelopeFilename(rootB))) + + _, err := f.ReadEnvelopeFromDisk(rootB) + require.Error(t, err) + require.False(t, f.HasEnvelope(rootB)) + exists, existsErr := afero.Exists(fs, getEnvelopeFilename(rootB)) + require.NoError(t, existsErr) + require.True(t, exists) +} + +func TestReadEnvelopeRejectsLengthAboveGossipLimit(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + var compressed bytes.Buffer + writer := snappy.NewBufferedWriter(&compressed) + length := make([]byte, 8) + binary.BigEndian.PutUint64(length, clparams.MaxChunkSize+1) + _, err := writer.Write([]byte{byte(clparams.GloasVersion)}) + require.NoError(t, err) + _, err = writer.Write(length) + require.NoError(t, err) + require.NoError(t, writer.Close()) + require.NoError(t, afero.WriteFile(fs, getEnvelopeFilename(root), compressed.Bytes(), 0o644)) + + _, err = f.ReadEnvelopeFromDisk(root) + require.ErrorContains(t, err, "exceeds max") + require.False(t, f.HasEnvelope(root)) +} + +func TestReadEnvelopeRejectsNonCanonicalSSZ(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + encoded, err := testEnvelopeWithTransaction(root, []byte{1}).EncodeSSZ(nil) + require.NoError(t, err) + messageOffset := binary.LittleEndian.Uint32(encoded) + binary.LittleEndian.PutUint32(encoded, messageOffset+1) + encoded = append(encoded[:messageOffset], append([]byte{0}, encoded[messageOffset:]...)...) + writeEnvelopeTestFile(t, fs, root, clparams.GloasVersion, encoded) + + _, err = f.ReadEnvelopeFromDisk(root) + require.Error(t, err) + require.False(t, f.HasEnvelope(root)) +} + +func TestReadEnvelopeValidatesDecodedEnvelopeAgainstConfig(t *testing.T) { + fs := afero.NewMemMapFs() + cfg := clparams.MainnetBeaconConfig + cfg.MaxWithdrawalRequestsPerPayload = 1 + f := &forkGraphDisk{fs: fs, beaconCfg: &cfg} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := testEnvelopeWithTransaction(root, []byte{1}) + envelope.Message.ExecutionRequests.Withdrawals.Append(&solid.WithdrawalRequest{}) + envelope.Message.ExecutionRequests.Withdrawals.Append(&solid.WithdrawalRequest{}) + encoded, err := envelope.EncodeSSZ(nil) + require.NoError(t, err) + writeEnvelopeTestFile(t, fs, root, clparams.GloasVersion, encoded) + + _, err = f.ReadEnvelopeFromDisk(root) + require.ErrorContains(t, err, "withdrawals") + require.False(t, f.HasEnvelope(root)) +} + +func TestReadEnvelopeRejectsKnownInvalidFileWithoutOpening(t *testing.T) { + baseFs := afero.NewMemMapFs() + fs := &envelopeWriteFailureFs{Fs: baseFs, stage: "open"} + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + f.invalidEnvelopes.Store(root, struct{}{}) + + _, err := f.ReadEnvelopeFromDisk(root) + require.ErrorContains(t, err, "known invalid") + require.NotErrorIs(t, err, errTestEnvelopeIO) +} + +func TestDumpEnvelopeRejectsLengthAboveGossipLimit(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + + err := f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, make([]byte, clparams.MaxChunkSize))) + require.ErrorContains(t, err, "exceeds max") + require.Zero(t, cap(f.sszBuffer)) + require.False(t, f.HasEnvelope(root)) +} + +func TestDumpEnvelopeRejectsIncompleteInput(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + validMessage := cltypes.NewExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + wrongPayloadVersion := cltypes.NewExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + wrongPayloadVersion.BeaconBlockRoot = root + wrongPayloadVersion.Payload = cltypes.NewEth1Block(clparams.DenebVersion, &clparams.MainnetBeaconConfig) + wrongRequestsVersion := cltypes.NewExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + wrongRequestsVersion.BeaconBlockRoot = root + wrongRequestsVersion.ExecutionRequests = cltypes.NewExecutionRequestsWithVersion(&clparams.MainnetBeaconConfig, clparams.ElectraVersion) + zeroRequests := cltypes.NewExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + zeroRequests.BeaconBlockRoot = root + zeroRequests.ExecutionRequests = &cltypes.ExecutionRequests{} + + for _, tt := range []struct { + name string + envelope *cltypes.SignedExecutionPayloadEnvelope + }{ + {name: "nil envelope"}, + {name: "nil message", envelope: &cltypes.SignedExecutionPayloadEnvelope{}}, + {name: "nil payload", envelope: &cltypes.SignedExecutionPayloadEnvelope{Message: &cltypes.ExecutionPayloadEnvelope{ExecutionRequests: validMessage.ExecutionRequests}}}, + {name: "nil execution requests", envelope: &cltypes.SignedExecutionPayloadEnvelope{Message: &cltypes.ExecutionPayloadEnvelope{Payload: validMessage.Payload}}}, + {name: "wrong payload version", envelope: &cltypes.SignedExecutionPayloadEnvelope{Message: wrongPayloadVersion}}, + {name: "wrong requests version", envelope: &cltypes.SignedExecutionPayloadEnvelope{Message: wrongRequestsVersion}}, + {name: "uninitialized requests", envelope: &cltypes.SignedExecutionPayloadEnvelope{Message: zeroRequests}}, + } { + t.Run(tt.name, func(t *testing.T) { + require.Error(t, f.DumpEnvelopeOnDisk(root, tt.envelope)) + }) + } +} + +func TestDumpEnvelopeRejectsNilNestedInput(t *testing.T) { + tests := []struct { + name string + mutate func(*cltypes.ExecutionPayloadEnvelope) + wantError string + }{ + {name: "payload extra data", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.Payload.Extra = nil }, wantError: "nil extra data"}, + {name: "payload transactions", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.Payload.Transactions = nil }, wantError: "nil transactions"}, + {name: "payload withdrawals", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.Payload.Withdrawals = nil }, wantError: "nil withdrawals"}, + {name: "payload block access list", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.Payload.BlockAccessList = nil }, wantError: "nil block access list"}, + {name: "deposit requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.Deposits = nil }, wantError: "nil deposit requests"}, + {name: "withdrawal requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.Withdrawals = nil }, wantError: "nil withdrawal requests"}, + {name: "consolidation requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.Consolidations = nil }, wantError: "nil consolidation requests"}, + {name: "builder deposit requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.BuilderDeposits = nil }, wantError: "nil builder deposit requests"}, + {name: "builder exit requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.BuilderExits = nil }, wantError: "nil builder exit requests"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := testEnvelopeWithTransaction(root, []byte{1}) + tt.mutate(envelope.Message) + + require.ErrorContains(t, f.DumpEnvelopeOnDisk(root, envelope), tt.wantError) + require.False(t, f.HasEnvelope(root)) + }) + } +} + +func TestDumpEnvelopeAcceptsInitializedEmptyNestedCollections(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := testEnvelopeWithTransaction(root, nil) + envelope.Message.Payload.Transactions = &solid.TransactionsSSZ{} + + require.NoError(t, f.DumpEnvelopeOnDisk(root, envelope)) + require.True(t, f.HasEnvelope(root)) +} + +func TestDumpEnvelopeRejectsNilNestedListMembers(t *testing.T) { + tests := []struct { + name string + mutate func(*cltypes.ExecutionPayloadEnvelope) + wantError string + }{ + {name: "payload withdrawals", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.Payload.Withdrawals.Append(nil) }, wantError: "nil withdrawal at index 0"}, + {name: "deposit requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.Deposits.Append(nil) }, wantError: "nil deposit request at index 0"}, + {name: "withdrawal requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.Withdrawals.Append(nil) }, wantError: "nil withdrawal request at index 0"}, + {name: "consolidation requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.Consolidations.Append(nil) }, wantError: "nil consolidation request at index 0"}, + {name: "builder deposit requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.BuilderDeposits.Append(nil) }, wantError: "nil builder deposit request at index 0"}, + {name: "builder exit requests", mutate: func(e *cltypes.ExecutionPayloadEnvelope) { e.ExecutionRequests.BuilderExits.Append(nil) }, wantError: "nil builder exit request at index 0"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := testEnvelopeWithTransaction(root, []byte{1}) + tt.mutate(envelope.Message) + + require.ErrorContains(t, f.DumpEnvelopeOnDisk(root, envelope), tt.wantError) + require.False(t, f.HasEnvelope(root)) + exists, err := afero.Exists(fs, getEnvelopeFilename(root)) + require.NoError(t, err) + require.False(t, exists) + }) + } +} + +func TestDumpEnvelopeAllowsAnchorRoot(t *testing.T) { + fs := afero.NewMemMapFs() + root := common.HexToHash("0x1234") + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig, anchorRoot: root} + + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1}))) +} + +func TestDumpEnvelopeFailurePreservesExistingFinal(t *testing.T) { + for _, stage := range []string{"open", "write", "sync", "close", "rename"} { + t.Run(stage, func(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1, 2, 3}))) + + var closes atomic.Int32 + f.fs = envelopeWriteFailureFs{Fs: fs, stage: stage, closes: &closes} + require.ErrorIs(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{9, 8, 7})), errTestEnvelopeIO) + if stage == "close" { + require.Equal(t, int32(1), closes.Load()) + } + f.fs = fs + + tempExists, err := afero.Exists(fs, getEnvelopeFilename(root)+".tmp") + require.NoError(t, err) + require.False(t, tempExists) + persisted, err := f.ReadEnvelopeFromDisk(root) + require.NoError(t, err) + require.Equal(t, [][]byte{{1, 2, 3}}, persisted.Message.Payload.Transactions.UnderlyngReference()) + }) + } +} + +func TestPruneDoesNotRaceEnvelopeReplacement(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + oldRoot := common.HexToHash("0x1") + newRoot := common.HexToHash("0x2") + for root, slot := range map[common.Hash]uint64{oldRoot: 1, newRoot: 3} { + block := cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.DenebVersion) + block.Block.Slot = slot + f.blocks.Store(root, block) + require.NoError(t, afero.WriteFile(fs, getBeaconStateFilename(root), []byte{1}, 0o644)) + } + require.NoError(t, f.DumpEnvelopeOnDisk(oldRoot, testEnvelopeWithTransaction(oldRoot, []byte{1}))) + require.NoError(t, afero.WriteFile(fs, getEnvelopeFilename(oldRoot)+".tmp", []byte("stale"), 0o644)) + + blockingFs := &envelopeBlockingRenameFs{Fs: fs, reached: make(chan struct{}), release: make(chan struct{})} + f.fs = blockingFs + dumpDone := make(chan error, 1) + go func() { + dumpDone <- f.DumpEnvelopeOnDisk(oldRoot, testEnvelopeWithTransaction(oldRoot, []byte{2})) + }() + <-blockingFs.reached + pruneDone := make(chan error, 1) + go func() { pruneDone <- f.Prune(2) }() + + pruneCompleted := false + select { + case err := <-pruneDone: + require.NoError(t, err) + pruneCompleted = true + case <-time.After(time.Second): + } + close(blockingFs.release) + require.NoError(t, <-dumpDone) + if !pruneCompleted { + require.NoError(t, <-pruneDone) + } + + finalExists, err := afero.Exists(fs, getEnvelopeFilename(oldRoot)) + require.NoError(t, err) + require.False(t, finalExists) + tempExists, err := afero.Exists(fs, getEnvelopeFilename(oldRoot)+".tmp") + require.NoError(t, err) + require.False(t, tempExists) + require.False(t, f.HasEnvelope(oldRoot)) +} + +func TestDumpEnvelopeRejectsPrunedRoot(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + oldRoot := common.HexToHash("0x1") + newRoot := common.HexToHash("0x2") + addEnvelopeTestBlock(f, oldRoot, 1) + addEnvelopeTestBlock(f, newRoot, 3) + require.NoError(t, afero.WriteFile(fs, getBeaconStateFilename(oldRoot), []byte{1}, 0o644)) + require.NoError(t, afero.WriteFile(fs, getBeaconStateFilename(newRoot), []byte{1}, 0o644)) + require.NoError(t, f.Prune(2)) + + require.Error(t, f.DumpEnvelopeOnDisk(oldRoot, testEnvelopeWithTransaction(oldRoot, []byte{1}))) + exists, err := afero.Exists(fs, getEnvelopeFilename(oldRoot)) + require.NoError(t, err) + require.False(t, exists) +} + +func TestDumpEnvelopeRejectsTooManyTransactionsBeforeWriting(t *testing.T) { + cfg := clparams.MainnetBeaconConfig + cfg.MaxTransactionsPerPayload = 1 + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &cfg} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := cltypes.NewExecutionPayloadEnvelope(&cfg) + envelope.BeaconBlockRoot = root + envelope.Payload.Extra = solid.NewExtraData() + envelope.Payload.Transactions = solid.NewTransactionsSSZFromTransactions([][]byte{{1}, {2}}) + envelope.Payload.Withdrawals = solid.NewStaticListSSZ[*cltypes.Withdrawal](int(cfg.MaxWithdrawalsPerPayload), 44) + envelope.Payload.BlockAccessList = solid.NewByteListSSZ(cfg.MaxBytesPerTransaction) + + err := f.DumpEnvelopeOnDisk(root, &cltypes.SignedExecutionPayloadEnvelope{Message: envelope}) + require.ErrorContains(t, err, "too many transactions") + exists, existsErr := afero.Exists(fs, getEnvelopeFilename(root)) + require.NoError(t, existsErr) + require.False(t, exists) +} + +func TestDumpEnvelopeRejectsDepositRepresentationUnreadableByConfiguredDecoder(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + envelope := testEnvelopeWithTransaction(root, nil) + envelope.Message.ExecutionRequests.Deposits = solid.NewStaticProgressiveListSSZ[*solid.DepositRequest](8193, solid.SizeDepositRequest) + for range 16_385 { + envelope.Message.ExecutionRequests.Deposits.Append(&solid.DepositRequest{}) + } + + err := f.DumpEnvelopeOnDisk(root, envelope) + require.ErrorContains(t, err, "decoder resource limit") + exists, existsErr := afero.Exists(fs, getEnvelopeFilename(root)) + require.NoError(t, existsErr) + require.False(t, exists) +} + +func TestPruneReportsEnvelopeRemovalFailureWithoutRecachingRoot(t *testing.T) { + baseFs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: baseFs, beaconCfg: &clparams.MainnetBeaconConfig} + oldRoot := common.HexToHash("0x1") + newRoot := common.HexToHash("0x2") + addEnvelopeTestBlock(f, oldRoot, 1) + addEnvelopeTestBlock(f, newRoot, 3) + require.NoError(t, afero.WriteFile(baseFs, getBeaconStateFilename(oldRoot), []byte{1}, 0o644)) + require.NoError(t, afero.WriteFile(baseFs, getBeaconStateFilename(newRoot), []byte{1}, 0o644)) + require.NoError(t, f.DumpEnvelopeOnDisk(oldRoot, testEnvelopeWithTransaction(oldRoot, []byte{1}))) + f.fs = envelopeRemoveFailureFs{Fs: baseFs, suffix: ".envelope.snappy_ssz"} + + require.ErrorIs(t, f.Prune(2), errTestEnvelopeIO) + require.False(t, f.HasEnvelope(oldRoot)) + _, err := f.ReadEnvelopeFromDisk(oldRoot) + require.Error(t, err) + _, exists := f.blocks.Load(oldRoot) + require.False(t, exists) +} + +func TestPruneAllowsUnrelatedEnvelopeIOBetweenRoots(t *testing.T) { + baseFs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: baseFs, beaconCfg: &clparams.MainnetBeaconConfig} + oldRootA := common.HexToHash("0x1") + oldRootB := common.HexToHash("0x2") + retainedRoot := common.HexToHash("0x3") + for root, slot := range map[common.Hash]uint64{oldRootA: 1, oldRootB: 2, retainedRoot: 4} { + addEnvelopeTestBlock(f, root, slot) + require.NoError(t, afero.WriteFile(baseFs, getBeaconStateFilename(root), []byte{1}, 0o644)) + } + blockingFs := &envelopeBlockingPruneFs{ + Fs: baseFs, + firstReached: make(chan struct{}), + releaseFirst: make(chan struct{}), + secondReached: make(chan struct{}), + releaseSecond: make(chan struct{}), + } + f.fs = blockingFs + pruneDone := make(chan error, 1) + go func() { pruneDone <- f.Prune(3) }() + <-blockingFs.firstReached + + dumpDone := make(chan error, 1) + go func() { + dumpDone <- f.DumpEnvelopeOnDisk(retainedRoot, testEnvelopeWithTransaction(retainedRoot, []byte{1})) + }() + time.Sleep(10 * time.Millisecond) + close(blockingFs.releaseFirst) + <-blockingFs.secondReached + select { + case err := <-dumpDone: + require.NoError(t, err) + case <-time.After(time.Second): + close(blockingFs.releaseSecond) + require.NoError(t, <-pruneDone) + t.Fatal("prune retained the disk lock while removing envelope files") + } + close(blockingFs.releaseSecond) + require.NoError(t, <-pruneDone) +} + +func TestReadEnvelopeOwnsDecodedTransactions(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + rootA := common.HexToHash("0xa") + rootB := common.HexToHash("0xb") + addEnvelopeTestBlock(f, rootA, 1) + addEnvelopeTestBlock(f, rootB, 2) + require.NoError(t, f.DumpEnvelopeOnDisk(rootA, testEnvelopeWithTransaction(rootA, []byte{1, 2, 3}))) + + persistedA, err := f.ReadEnvelopeFromDisk(rootA) + require.NoError(t, err) + require.NoError(t, f.DumpEnvelopeOnDisk(rootB, testEnvelopeWithTransaction(rootB, []byte{9, 8, 7}))) + _, err = f.ReadEnvelopeFromDisk(rootB) + require.NoError(t, err) + require.Equal(t, [][]byte{{1, 2, 3}}, persistedA.Message.Payload.Transactions.UnderlyngReference()) +} + +func TestReadEnvelopeTransactionsDoNotRaceWithDump(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + rootA := common.HexToHash("0xa") + rootB := common.HexToHash("0xb") + addEnvelopeTestBlock(f, rootA, 1) + addEnvelopeTestBlock(f, rootB, 2) + require.NoError(t, f.DumpEnvelopeOnDisk(rootA, testEnvelopeWithTransaction(rootA, []byte{1, 2, 3}))) + persistedA, err := f.ReadEnvelopeFromDisk(rootA) + require.NoError(t, err) + envelopeB := testEnvelopeWithTransaction(rootB, []byte{9, 8, 7}) + + start := make(chan struct{}) + errCh := make(chan error, 1) + var wg sync.WaitGroup + wg.Go(func() { + <-start + for range 100 { + if err := f.DumpEnvelopeOnDisk(rootB, envelopeB); err != nil { + errCh <- err + return + } + } + }) + close(start) + var observed uint64 + for range 100 { + observed += uint64(persistedA.Message.Payload.Transactions.UnderlyngReference()[0][0]) + } + wg.Wait() + require.Equal(t, uint64(100), observed) + require.Equal(t, [][]byte{{1, 2, 3}}, persistedA.Message.Payload.Transactions.UnderlyngReference()) + close(errCh) + for err := range errCh { + require.NoError(t, err) + } +} + +func TestReadEnvelopeCloseErrorKeepsDecodedFile(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1}))) + f.fs = envelopeCloseErrorFs{Fs: fs} + + envelope, err := f.ReadEnvelopeFromDisk(root) + require.NoError(t, err) + require.NotNil(t, envelope) + require.True(t, f.HasEnvelope(root)) +} + +func TestReadEnvelopeTransientReadErrorKeepsFile(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1}))) + f.fs = envelopeReadErrorFs{Fs: fs} + + _, err := f.ReadEnvelopeFromDisk(root) + require.ErrorIs(t, err, errTestEnvelopeIO) + require.True(t, f.HasEnvelope(root)) + exists, existsErr := afero.Exists(fs, getEnvelopeFilename(root)) + require.NoError(t, existsErr) + require.True(t, exists) +} + +func TestReadEnvelopeTransientOpenErrorKeepsTrustedCache(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1}))) + f.fs = envelopeOpenErrorFs{Fs: fs} + + _, err := f.ReadEnvelopeFromDisk(root) + require.ErrorIs(t, err, errTestEnvelopeIO) + require.True(t, f.HasEnvelope(root)) +} + +func TestReadEnvelopeStructuralCorruptionEvictsTrustedCache(t *testing.T) { + fs := afero.NewMemMapFs() + f := &forkGraphDisk{fs: fs, beaconCfg: &clparams.MainnetBeaconConfig} + root := common.HexToHash("0x1234") + addEnvelopeTestBlock(f, root, 1) + require.NoError(t, f.DumpEnvelopeOnDisk(root, testEnvelopeWithTransaction(root, []byte{1}))) + require.NoError(t, afero.WriteFile(fs, getEnvelopeFilename(root), []byte("corrupt"), 0o644)) + + _, err := f.ReadEnvelopeFromDisk(root) + require.Error(t, err) + require.False(t, f.HasEnvelope(root)) + _, invalid := f.invalidEnvelopes.Load(root) + require.True(t, invalid) +} + +func testEnvelopeWithTransaction(root common.Hash, transaction []byte) *cltypes.SignedExecutionPayloadEnvelope { + return testEnvelopeWithVersion(root, transaction, clparams.GloasVersion) +} + +func testEnvelopeWithVersion(root common.Hash, transaction []byte, version clparams.StateVersion) *cltypes.SignedExecutionPayloadEnvelope { + envelope := cltypes.NewExecutionPayloadEnvelope(&clparams.MainnetBeaconConfig) + envelope.BeaconBlockRoot = root + envelope.Payload = cltypes.NewEth1Block(version, &clparams.MainnetBeaconConfig) + envelope.Payload.Extra = solid.NewExtraData() + envelope.Payload.Transactions = solid.NewTransactionsSSZFromTransactions([][]byte{transaction}) + envelope.Payload.Withdrawals = solid.NewStaticListSSZ[*cltypes.Withdrawal](int(clparams.MainnetBeaconConfig.MaxWithdrawalsPerPayload), 44) + envelope.ExecutionRequests = cltypes.NewExecutionRequestsWithVersion(&clparams.MainnetBeaconConfig, version) + return &cltypes.SignedExecutionPayloadEnvelope{Message: envelope} +} + +func addEnvelopeTestBlock(f *forkGraphDisk, root common.Hash, slot uint64) { + block := cltypes.NewSignedBeaconBlock(&clparams.MainnetBeaconConfig, clparams.DenebVersion) + block.Block.Slot = slot + f.blocks.Store(root, block) +} + +func writeEnvelopeTestFile(t *testing.T, fs afero.Fs, root common.Hash, version clparams.StateVersion, encoded []byte) { + t.Helper() + var compressed bytes.Buffer + writer := snappy.NewBufferedWriter(&compressed) + _, err := writer.Write([]byte{byte(version)}) + require.NoError(t, err) + length := make([]byte, 8) + binary.BigEndian.PutUint64(length, uint64(len(encoded))) + _, err = writer.Write(length) + require.NoError(t, err) + _, err = writer.Write(encoded) + require.NoError(t, err) + require.NoError(t, writer.Close()) + require.NoError(t, afero.WriteFile(fs, getEnvelopeFilename(root), compressed.Bytes(), 0o644)) +} diff --git a/cl/phase1/forkchoice/mock_services/forkchoice_mock.go b/cl/phase1/forkchoice/mock_services/forkchoice_mock.go index 5c283b4a6a8..663d38c870e 100644 --- a/cl/phase1/forkchoice/mock_services/forkchoice_mock.go +++ b/cl/phase1/forkchoice/mock_services/forkchoice_mock.go @@ -352,6 +352,9 @@ func (f *ForkChoiceStorageMock) OnBlock( } func (f *ForkChoiceStorageMock) OnExecutionPayload(ctx context.Context, signedEnvelope *cltypes.SignedExecutionPayloadEnvelope, checkBlobData, validatePayload bool) error { + if f.OnExecutionPayloadErr == nil && signedEnvelope != nil && signedEnvelope.Message != nil { + f.Envelopes[signedEnvelope.Message.BeaconBlockRoot] = signedEnvelope + } return f.OnExecutionPayloadErr } diff --git a/cl/phase1/forkchoice/on_execution_payload.go b/cl/phase1/forkchoice/on_execution_payload.go index 780855a1cbd..5cc45b3ff4c 100644 --- a/cl/phase1/forkchoice/on_execution_payload.go +++ b/cl/phase1/forkchoice/on_execution_payload.go @@ -485,8 +485,8 @@ func (f *ForkChoiceStore) StoreAnchorEnvelope(blockRoot common.Hash, signedEnvel // - checkBlobData: if true, verify blob data availability via PeerDAS before processing // - validatePayload: if true, call engine.NewPayload() to validate with EL before state transition func (f *ForkChoiceStore) OnExecutionPayload(ctx context.Context, signedEnvelope *cltypes.SignedExecutionPayloadEnvelope, checkBlobData, validatePayload bool) error { - if signedEnvelope == nil || signedEnvelope.Message == nil { - return errors.New("nil execution payload envelope") + if err := signedEnvelope.ValidateForConfig(f.beaconCfg); err != nil { + return fmt.Errorf("invalid execution payload envelope: %w", err) } envelope := signedEnvelope.Message @@ -525,8 +525,8 @@ func (f *ForkChoiceStore) OnExecutionPayload(ctx context.Context, signedEnvelope // verifies BLS signatures. // [New in Gloas:EIP7732] func (f *ForkChoiceStore) ApplyLocalSelfBuildEnvelope(ctx context.Context, signedEnvelope *cltypes.SignedExecutionPayloadEnvelope) error { - if signedEnvelope == nil || signedEnvelope.Message == nil { - return errors.New("nil execution payload envelope") + if err := signedEnvelope.ValidateForConfig(f.beaconCfg); err != nil { + return fmt.Errorf("invalid execution payload envelope: %w", err) } envelope := signedEnvelope.Message diff --git a/cl/phase1/forkchoice/on_execution_payload_test.go b/cl/phase1/forkchoice/on_execution_payload_test.go index 85dd3524a94..ec47b7d2e4b 100644 --- a/cl/phase1/forkchoice/on_execution_payload_test.go +++ b/cl/phase1/forkchoice/on_execution_payload_test.go @@ -62,6 +62,22 @@ func TestValidateEnvelopeAgainstBlock_NoBid(t *testing.T) { require.Contains(t, err.Error(), "block missing signed_execution_payload_bid") } +func TestOnExecutionPayloadRejectsNilWithdrawalBeforeForkchoice(t *testing.T) { + cfg := &clparams.MainnetBeaconConfig + envelope := cltypes.NewExecutionPayloadEnvelope(cfg) + envelope.Payload.Extra = solid.NewExtraData() + envelope.Payload.Transactions = solid.NewTransactionsSSZFromTransactions(nil) + envelope.Payload.Withdrawals = solid.NewStaticListSSZ[*cltypes.Withdrawal](int(cfg.MaxWithdrawalsPerPayload), 44) + envelope.Payload.Withdrawals.Append(nil) + envelope.Payload.BlockAccessList = solid.NewByteListSSZ(cfg.MaxBytesPerTransaction) + f := &ForkChoiceStore{beaconCfg: cfg} + + require.NotPanics(t, func() { + err := f.OnExecutionPayload(context.Background(), &cltypes.SignedExecutionPayloadEnvelope{Message: envelope}, false, true) + require.ErrorContains(t, err, "nil withdrawal at index 0") + }) +} + // TestValidateEnvelopeAgainstBlock_SlotNumberMismatch tests that validation fails when // block.slot != envelope.payload.slot_number (EIP-7843 / GLOAS p2p-interface REJECT rule). func TestValidateEnvelopeAgainstBlock_SlotNumberMismatch(t *testing.T) { diff --git a/cl/phase1/network/backward_beacon_downloader.go b/cl/phase1/network/backward_beacon_downloader.go index ef7cfad7315..ea5ca78901b 100644 --- a/cl/phase1/network/backward_beacon_downloader.go +++ b/cl/phase1/network/backward_beacon_downloader.go @@ -716,7 +716,7 @@ func (b *BackwardBeaconDownloader) fetchSingleEnvelope(ctx context.Context, bloc if err != nil { return nil, err } - body, err := io.ReadAll(resp.Body) + body, err := readEnvelopeHTTPBody(resp.Body) resp.Body.Close() if err != nil { return nil, err @@ -731,7 +731,7 @@ func (b *BackwardBeaconDownloader) fetchSingleEnvelope(ctx context.Context, bloc envelope := &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(b.beaconCfg), } - if err := envelope.DecodeSSZ(body, int(clparams.GloasVersion)); err != nil { + if err := envelope.DecodeSSZStrict(body, int(clparams.GloasVersion)); err != nil { return nil, fmt.Errorf("envelope decode: %w", err) } return envelope, nil diff --git a/cl/phase1/network/backward_beacon_downloader_test.go b/cl/phase1/network/backward_beacon_downloader_test.go index 5d5a10e0b14..f9e4b95cbc3 100644 --- a/cl/phase1/network/backward_beacon_downloader_test.go +++ b/cl/phase1/network/backward_beacon_downloader_test.go @@ -224,3 +224,34 @@ func TestBackwardBeaconDownloaderHTTPPreferredEmptyResponseFallsBack(t *testing. t.Fatal("httpPreferred remained true after empty HTTP response") } } + +func TestBackwardBeaconDownloaderRejectsOversizedEnvelopeResponse(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + _, _ = w.Write(make([]byte, clparams.MaxChunkSize+1)) + })) + defer server.Close() + downloader := &BackwardBeaconDownloader{ + httpFallbackURL: server.URL, + beaconCfg: &clparams.MainnetBeaconConfig, + } + + _, err := downloader.fetchSingleEnvelope(context.Background(), makeGloasBlock(1, hash(1), hash(2))) + require.ErrorContains(t, err, "too large") +} + +func TestForwardBeaconDownloaderRejectsOversizedEnvelopeResponse(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + _, _ = w.Write(make([]byte, clparams.MaxChunkSize+1)) + })) + defer server.Close() + block := makeGloasBlock(1, hash(1), hash(2)) + root, err := block.Block.HashSSZ() + require.NoError(t, err) + received := map[common.Hash]*cltypes.SignedExecutionPayloadEnvelope{} + + fetched := fetchEnvelopesFromBeaconAPI( + context.Background(), server.URL, []*cltypes.SignedBeaconBlock{block}, [][32]byte{root}, received, &clparams.MainnetBeaconConfig, + ) + require.Zero(t, fetched) + require.Empty(t, received) +} diff --git a/cl/phase1/network/beacon_downloader.go b/cl/phase1/network/beacon_downloader.go index 259c9c75f85..54455d24949 100644 --- a/cl/phase1/network/beacon_downloader.go +++ b/cl/phase1/network/beacon_downloader.go @@ -612,7 +612,7 @@ func fetchEnvelopesFromBeaconAPI( if err != nil { return } - body, err := io.ReadAll(resp.Body) + body, err := readEnvelopeHTTPBody(resp.Body) resp.Body.Close() if err != nil || resp.StatusCode != http.StatusOK { return @@ -621,7 +621,7 @@ func fetchEnvelopesFromBeaconAPI( envelope := &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(beaconCfg), } - if err := envelope.DecodeSSZ(body, int(clparams.GloasVersion)); err != nil { + if err := envelope.DecodeSSZStrict(body, int(clparams.GloasVersion)); err != nil { log.Debug("[ForwardBeaconDownloader] HTTP envelope decode failed", "slot", slot, "err", err) return } @@ -640,6 +640,17 @@ func fetchEnvelopesFromBeaconAPI( return fetched } +func readEnvelopeHTTPBody(r io.Reader) ([]byte, error) { + body, err := io.ReadAll(io.LimitReader(r, int64(clparams.MaxChunkSize)+1)) + if err != nil { + return nil, err + } + if uint64(len(body)) > clparams.MaxChunkSize { + return nil, fmt.Errorf("execution payload envelope response too large: max %d bytes", clparams.MaxChunkSize) + } + return body, nil +} + // GetHighestProcessedSlot retrieve the highest processed slot we accumulated. func (f *ForwardBeaconDownloader) GetHighestProcessedSlot() uint64 { f.mu.Lock() diff --git a/cl/phase1/network/services/execution_payload_service.go b/cl/phase1/network/services/execution_payload_service.go index cdac0bf5c17..ae7868e553f 100644 --- a/cl/phase1/network/services/execution_payload_service.go +++ b/cl/phase1/network/services/execution_payload_service.go @@ -113,7 +113,7 @@ func (s *executionPayloadService) DecodeGossipMessage(_ peer.ID, data []byte, ve obj := &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(s.beaconCfg), } - if err := obj.DecodeSSZ(data, int(version)); err != nil { + if err := obj.DecodeSSZStrict(data, int(version)); err != nil { return nil, err } return obj, nil @@ -123,8 +123,8 @@ func (s *executionPayloadService) DecodeGossipMessage(_ peer.ID, data []byte, ve // Reference: https://github.com/ethereum/consensus-specs/blob/dev/specs/_features/epbs/p2p-interface.md#execution_payload // [New in Gloas:EIP7732] func (s *executionPayloadService) ProcessMessage(ctx context.Context, _ *uint64, signedEnvelope *cltypes.SignedExecutionPayloadEnvelope) error { - if signedEnvelope == nil || signedEnvelope.Message == nil { - return errors.New("nil execution payload envelope") + if err := signedEnvelope.ValidateForConfig(s.beaconCfg); err != nil { + return fmt.Errorf("invalid execution payload envelope: %w", err) } envelope := signedEnvelope.Message @@ -161,7 +161,7 @@ func (s *executionPayloadService) ProcessMessage(ctx context.Context, _ *uint64, beaconBlockRoot: beaconBlockRoot, builderIndex: builderIndex, } - if s.seenEnvelopesCache.Contains(seenKey) { + if s.seenEnvelopesCache.Contains(seenKey) && s.forkchoiceStore.HasEnvelope(beaconBlockRoot) { return fmt.Errorf("%w: already seen envelope for block %v from builder %d", ErrIgnore, beaconBlockRoot, builderIndex) } diff --git a/cl/phase1/network/services/execution_payload_service_test.go b/cl/phase1/network/services/execution_payload_service_test.go index 9e32176d091..ab2ad55d6eb 100644 --- a/cl/phase1/network/services/execution_payload_service_test.go +++ b/cl/phase1/network/services/execution_payload_service_test.go @@ -18,6 +18,7 @@ package services import ( "context" + "encoding/binary" "errors" "sync" "testing" @@ -49,6 +50,8 @@ func newTestSignedEnvelope(slot uint64, blockRoot common.Hash, builderIndex uint if envelope.Payload != nil { envelope.Payload.Extra = solid.NewExtraData() envelope.Payload.Transactions = &solid.TransactionsSSZ{} + envelope.Payload.Withdrawals = solid.NewStaticListSSZ[*cltypes.Withdrawal](int(clparams.MainnetBeaconConfig.MaxWithdrawalsPerPayload), 44) + envelope.Payload.BlockAccessList = solid.NewByteListSSZ(clparams.MainnetBeaconConfig.MaxBytesPerTransaction) } return &cltypes.SignedExecutionPayloadEnvelope{ Message: envelope, @@ -56,6 +59,57 @@ func newTestSignedEnvelope(slot uint64, blockRoot common.Hash, builderIndex uint } } +func oversizedExtraDataEnvelopeSSZ(t *testing.T, envelope *cltypes.SignedExecutionPayloadEnvelope) []byte { + t.Helper() + encoded, err := envelope.EncodeSSZ(nil) + require.NoError(t, err) + + const ( + signedMessageOffsetPosition = 0 + messageRequestsOffsetPosition = 4 + payloadExtraOffsetPosition = 436 + payloadTransactionsOffsetPosition = 504 + payloadWithdrawalsOffsetPosition = 508 + payloadBlockAccessListOffsetPosition = 528 + ) + messageStart := int(binary.LittleEndian.Uint32(encoded[signedMessageOffsetPosition:])) + payloadStart := messageStart + int(binary.LittleEndian.Uint32(encoded[messageStart:])) + extraStart := payloadStart + int(binary.LittleEndian.Uint32(encoded[payloadStart+payloadExtraOffsetPosition:])) + malformed := append([]byte{}, encoded[:extraStart]...) + malformed = append(malformed, make([]byte, 33)...) + malformed = append(malformed, encoded[extraStart:]...) + + for _, position := range []int{ + messageStart + messageRequestsOffsetPosition, + payloadStart + payloadTransactionsOffsetPosition, + payloadStart + payloadWithdrawalsOffsetPosition, + payloadStart + payloadBlockAccessListOffsetPosition, + } { + offset := binary.LittleEndian.Uint32(malformed[position:]) + binary.LittleEndian.PutUint32(malformed[position:], offset+33) + } + return malformed +} + +func TestExecutionPayloadServiceRejectsOversizedExtraDataSSZ(t *testing.T) { + service, _ := setupExecutionPayloadService(t) + envelope := newTestSignedEnvelope(100, common.HexToHash("0x1234"), 1) + + _, err := service.DecodeGossipMessage("", oversizedExtraDataEnvelopeSSZ(t, envelope), clparams.GloasVersion) + require.Error(t, err) +} + +func TestExecutionPayloadServiceRejectsMalformedEnvelopeBeforePendingHash(t *testing.T) { + service, _ := setupExecutionPayloadService(t) + envelope := newTestSignedEnvelope(100, common.HexToHash("0x1234"), 1) + envelope.Message.Payload.Withdrawals.Append(nil) + + require.NotPanics(t, func() { + err := service.ProcessMessage(context.Background(), nil, envelope) + require.ErrorContains(t, err, "nil withdrawal at index 0") + }) +} + func TestExecutionPayloadServiceNilEnvelope(t *testing.T) { service, _ := setupExecutionPayloadService(t) @@ -121,6 +175,16 @@ func TestExecutionPayloadServiceAlreadySeen(t *testing.T) { require.Contains(t, err.Error(), "already seen envelope") } +func TestExecutionPayloadServiceRetriesSeenEnvelopeWhenPersistenceIsUnavailable(t *testing.T) { + service, fcu := setupExecutionPayloadService(t) + blockRoot := common.HexToHash("0x1234") + envelope := newTestSignedEnvelope(100, blockRoot, 1) + fcu.Blocks[blockRoot] = &cltypes.SignedBeaconBlock{Block: &cltypes.BeaconBlock{Slot: 100}} + service.(*executionPayloadService).seenEnvelopesCache.Add(seenEnvelopeKey{blockRoot, 1}, struct{}{}) + + require.NoError(t, service.ProcessMessage(context.Background(), nil, envelope)) +} + func TestExecutionPayloadServiceSlotBelowFinalized(t *testing.T) { service, fcu := setupExecutionPayloadService(t) diff --git a/cl/rpc/rpc.go b/cl/rpc/rpc.go index 9dac6a11221..e154a0143e7 100644 --- a/cl/rpc/rpc.go +++ b/cl/rpc/rpc.go @@ -195,7 +195,7 @@ func (b *BeaconRpcP2P) SendExecutionPayloadEnvelopesByRangeReq(ctx context.Conte envelope := &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(b.beaconConfig), } - if err := envelope.DecodeSSZ(data.raw, int(data.version)); err != nil { + if err := envelope.DecodeSSZStrict(data.raw, int(data.version)); err != nil { return nil, pid, err } envelopes = append(envelopes, envelope) @@ -232,7 +232,7 @@ func (b *BeaconRpcP2P) SendExecutionPayloadEnvelopesByRootReq(ctx context.Contex envelope := &cltypes.SignedExecutionPayloadEnvelope{ Message: cltypes.NewExecutionPayloadEnvelope(b.beaconConfig), } - if err := envelope.DecodeSSZ(data.raw, int(data.version)); err != nil { + if err := envelope.DecodeSSZStrict(data.raw, int(data.version)); err != nil { return nil, pid, err } envelopes = append(envelopes, envelope)