Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion cl/beacon/handler/epbs.go
Original file line number Diff line number Diff line change
Expand Up @@ -816,7 +816,7 @@ func (a *ApiHandler) PostEthV1BeaconExecutionPayloadEnvelope(w http.ResponseWrit
// checkBlobData=false because gossip validation handles it; validatePayload=true
// so the EL receives NewPayload for the execution payload.
if err := a.forkchoiceStore.OnExecutionPayload(r.Context(), signedEnvelope, false, true); err != nil {
if errors.Is(err, forkchoice.ErrIgnore) || errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) {
if errors.Is(err, forkchoice.ErrIgnore) || errors.Is(err, forkchoice.ErrExecutionPayloadAlreadyStored) || errors.Is(err, forkchoice.ErrEIP7594ColumnDataNotAvailable) {
a.logger.Debug("[Beacon REST] OnExecutionPayload queued or ignored", "err", err)
} else {
beaconhttp.WrapEndpointError(err).WriteTo(w)
Expand Down
14 changes: 14 additions & 0 deletions cl/beacon/handler/epbs_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,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/forkchoice"
"github.com/erigontech/erigon/cl/phase1/network/services"
mock_services "github.com/erigontech/erigon/cl/phase1/network/services/mock_services"
"github.com/erigontech/erigon/cl/pool"
Expand Down Expand Up @@ -165,6 +166,19 @@ func TestPostExecutionPayloadEnvelopeReturnsForkchoiceError(t *testing.T) {
require.Contains(t, recorder.Body.String(), "invalid execution payload")
}

func TestPostExecutionPayloadEnvelopeAcceptsAlreadyStored(t *testing.T) {
_, _, _, _, _, handler, _, _, fcu, _ := setupTestingHandler(t, clparams.BellatrixVersion, log.Root(), true)
fcu.OnExecutionPayloadErr = forkchoice.ErrExecutionPayloadAlreadyStored

request := httptest.NewRequest(http.MethodPost, "/eth/v1/beacon/execution_payload_envelope", strings.NewReader(`{}`))
request.Header.Set("Content-Type", "application/json; charset=utf-8")
recorder := httptest.NewRecorder()

handler.PostEthV1BeaconExecutionPayloadEnvelope(recorder, request)

require.Equal(t, http.StatusOK, recorder.Code, recorder.Body.String())
}

func TestPostPtcDutiesDoesNotCapValidatorCount(t *testing.T) {
_, _, _, _, _, handler, _, _, _, _ := setupTestingHandler(t, clparams.BellatrixVersion, log.Root(), true)
handler.beaconChainCfg.GloasForkEpoch = 0
Expand Down
85 changes: 71 additions & 14 deletions cl/phase1/forkchoice/fork_graph/fork_graph_disk_fs.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ package fork_graph

import (
"encoding/binary"
"errors"
"fmt"
"io"
"os"
Expand Down Expand Up @@ -193,6 +194,11 @@ func (f *forkGraphDisk) HasEnvelope(blockRoot common.Hash) bool {
if _, ok := f.envelopeExists.Load(blockRoot); ok {
return true
}
f.stateDumpLock.Lock()
defer f.stateDumpLock.Unlock()
if _, ok := f.envelopeExists.Load(blockRoot); ok {
return true
}
// Slow path: fall back to disk and populate cache on hit
exists, err := afero.Exists(f.fs, getEnvelopeFilename(blockRoot))
if err == nil && exists {
Expand All @@ -206,57 +212,91 @@ func (f *forkGraphDisk) HasEnvelope(blockRoot common.Hash) bool {
// [New in Gloas:EIP7732]
func (f *forkGraphDisk) ReadEnvelopeFromDisk(blockRoot common.Hash) (envelope *cltypes.SignedExecutionPayloadEnvelope, err error) {
var file afero.File
var corrupt bool
f.stateDumpLock.Lock()
defer f.stateDumpLock.Unlock()

file, err = f.fs.Open(getEnvelopeFilename(blockRoot))
filename := getEnvelopeFilename(blockRoot)
file, err = f.fs.Open(filename)
if err != nil {
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 {
if removeErr := f.fs.Remove(filename); removeErr != nil && !os.IsNotExist(removeErr) {
log.Warn("failed to remove corrupt envelope", "root", blockRoot, "err", removeErr)
}
}
if err != nil {
f.envelopeExists.Delete(blockRoot)
}
}()

readTracker := &envelopeReadTracker{Reader: file}
if f.sszSnappyReader == nil {
f.sszSnappyReader = snappy.NewReader(file)
f.sszSnappyReader = snappy.NewReader(readTracker)
} else {
f.sszSnappyReader.Reset(file)
f.sszSnappyReader.Reset(readTracker)
}

// Read the length
lengthBytes := make([]byte, 8)
var n int
n, err = io.ReadFull(f.sszSnappyReader, 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 {
corrupt = true
return nil, fmt.Errorf("failed to read length: %d, want 8, root: %x", n, blockRoot)
}

envelopeLength := binary.BigEndian.Uint64(lengthBytes)
if envelopeLength > maxSSZObjectSize {
corrupt = true
return nil, fmt.Errorf("corrupt envelope file: length %d exceeds max %d, root: %x", envelopeLength, maxSSZObjectSize, 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)
ownedBuffer := make([]byte, envelopeLength)
n, err = io.ReadFull(f.sszSnappyReader, ownedBuffer)
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]
ownedBuffer = ownedBuffer[:n]

envelope = &cltypes.SignedExecutionPayloadEnvelope{
Message: cltypes.NewExecutionPayloadEnvelope(f.beaconCfg),
}
if err = envelope.DecodeSSZ(f.sszBuffer, int(clparams.GloasVersion)); err != nil {
if err = envelope.DecodeSSZ(ownedBuffer, int(clparams.GloasVersion)); err != nil {
corrupt = true
return nil, fmt.Errorf("failed to decode envelope: %w, root: %x, len: %d", err, blockRoot, n)
}

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) {
Expand All @@ -276,11 +316,21 @@ func (f *forkGraphDisk) DumpEnvelopeOnDisk(blockRoot common.Hash, envelope *clty
return
}

dumpedFile, err := f.fs.OpenFile(getEnvelopeFilename(blockRoot), os.O_TRUNC|os.O_CREATE|os.O_RDWR, 0o755)
filename := getEnvelopeFilename(blockRoot)
tempFilename := filename + ".tmp"
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)
Expand Down Expand Up @@ -309,6 +359,13 @@ func (f *forkGraphDisk) DumpEnvelopeOnDisk(blockRoot common.Hash, envelope *clty
log.Error("failed to sync dumped file", "err", err)
return
}
if err = dumpedFile.Close(); err != nil {
return
}
closed = true
if err = f.fs.Rename(tempFilename, filename); err != nil {
return
}

return
}
Loading
Loading