refactored and updated e2store framing for objects

This commit is contained in:
shantichanal 2025-07-09 15:48:00 +02:00 committed by lightclient
parent fd11aa11e4
commit 65a08322cb
No known key found for this signature in database
GPG key ID: 657913021EF45A6A
3 changed files with 165 additions and 223 deletions

View file

@ -58,7 +58,6 @@ const (
headerSize uint64 = 8
)
// ProofVariant is an idiomatic “enum”.
type proofvar uint16
const (
@ -102,7 +101,6 @@ type Builder struct {
writtenBytes uint64
}
// NewBuilder returns a new EraE Builder writing into the given io.Writer.
func NewBuilder(w io.Writer) *Builder {
buf := bytes.NewBuffer(nil)
return &Builder{
@ -117,27 +115,27 @@ func (b *Builder) Add(header types.Header, body types.Body, receipts types.Recei
return fmt.Errorf("exceeds MaxEraESize %d", MaxEraESize)
}
var pv proofvar = ProofNone
var pData []byte
if proof != nil {
if proof.Variant == ProofNone || len(proof.Data) == 0 {
return fmt.Errorf("invalid proof: variant=%d len=%d", proof.Variant, len(proof.Data))
pv = proof.Variant
pData = proof.Data
if pv == ProofNone || len(pData) == 0 {
return fmt.Errorf("invalid proof: variant=%d len=%d", pv, len(pData))
}
}
if len(b.headersRLP) == 0 {
if proof == nil {
b.prooftype = ProofNone
} else {
b.prooftype = proof.Variant
}
} else if proof.Variant != b.prooftype {
return fmt.Errorf("cannot mix proof variants: have %d want %d", b.prooftype, proof.Variant)
b.prooftype = pv
} else if pv != b.prooftype {
return fmt.Errorf("cannot mix proof variants: have %d want %d", b.prooftype, pv)
}
hdr, err := rlp.EncodeToBytes(header)
hdr, err := rlp.EncodeToBytes(&header)
if err != nil {
return fmt.Errorf("error encoding block header: %w", err)
}
bod, err := rlp.EncodeToBytes(body)
bod, err := rlp.EncodeToBytes(&body)
if err != nil {
return fmt.Errorf("error encoding block header: %w", err)
}
@ -174,30 +172,24 @@ func (b *Builder) Finalize() (common.Hash, error) {
return common.Hash{}, errors.New("no blocks added, cannot finalize")
}
var err error
b.headeroffsets, err = b.writeSection(TypeCompressedHeader, b.headersRLP, true)
if err != nil {
return common.Hash{}, fmt.Errorf("error writing compressed headers: %w", err)
if err := b.flushKind(TypeCompressedHeader, b.headersRLP, true, &b.headeroffsets); err != nil {
return common.Hash{}, err
}
b.bodyoffsets, err = b.writeSection(TypeCompressedBody, b.bodiesRLP, true)
if err != nil {
return common.Hash{}, fmt.Errorf("error writing compressed bodies: %w", err)
if err := b.flushKind(TypeCompressedBody, b.bodiesRLP, true, &b.bodyoffsets); err != nil {
return common.Hash{}, err
}
b.receiptoffsets, err = b.writeSection(TypeCompressedReceipts, b.receiptsRLP, true)
if err != nil {
return common.Hash{}, fmt.Errorf("error writing compressed receipts: %w", err)
if err := b.flushKind(TypeCompressedReceipts, b.receiptsRLP, true, &b.receiptoffsets); err != nil {
return common.Hash{}, err
}
if len(b.tds) > 0 {
b.tdoff, err = b.writeSection(TypeTotalDifficulty, b.tds, false)
if err != nil {
return common.Hash{}, fmt.Errorf("error writing total difficulties: %w", err)
if err := b.flushKind(TypeTotalDifficulty, b.tds, false, &b.tdoff); err != nil {
return common.Hash{}, err
}
}
if b.prooftype != 0 {
b.proofoffsets, err = b.writeSection(uint16(b.prooftype), b.proofsRLP, true)
if err != nil {
return common.Hash{}, fmt.Errorf("error writing proofs: %w", err)
if b.prooftype != ProofNone {
if err := b.flushKind(uint16(b.prooftype), b.proofsRLP, true, &b.proofoffsets); err != nil {
return common.Hash{}, err
}
}
@ -236,7 +228,6 @@ func decodeBigs(raw [][]byte) []*big.Int {
return out
}
// snappyWrite is a small helper to take care snappy encoding and writing an e2store entry.
func (b *Builder) snappyWrite(typ uint16, in []byte) error {
var (
buf = b.buf
@ -258,29 +249,32 @@ func (b *Builder) snappyWrite(typ uint16, in []byte) error {
return nil
}
func (b *Builder) writeSection(typ uint16, list [][]byte, useSnappy bool) ([]uint64, error) {
if len(list) == 0 {
return nil, errors.New("cannot write empty section")
func (b *Builder) addEntry(typ uint16, payload []byte, snappyIt bool) (uint64, error) {
offset := b.writtenBytes
var err error
if snappyIt {
if err = b.snappyWrite(typ, payload); err != nil {
return 0, err
}
buf := bytes.NewBuffer(nil)
offs := make([]uint64, len(list))
for i, data := range list {
offs[i] = b.writtenBytes + headerSize + uint64(buf.Len())
if useSnappy {
buf.Write(snappy.Encode(nil, data))
} else {
buf.Write(data)
var n int
if n, err = b.w.Write(typ, payload); err != nil {
return 0, err
}
}
if n, err := b.w.Write(typ, buf.Bytes()); err != nil {
return nil, fmt.Errorf("error writing section %d: %w", typ, err)
} else {
b.writtenBytes += uint64(n)
}
return offset, nil
}
return offs, nil
func (b *Builder) flushKind(typ uint16, list [][]byte, useSnappy bool, dst *[]uint64) error {
for _, data := range list {
off, err := b.addEntry(typ, data, useSnappy)
if err != nil {
return fmt.Errorf("entry type %d: %w", typ, err)
}
*dst = append(*dst, off)
}
return nil
}
func (b *Builder) writeIndex() error {
@ -306,7 +300,7 @@ func (b *Builder) writeIndex() error {
binary.LittleEndian.PutUint64(idx[pos:], rel(b.receiptoffsets[i]))
pos += 8
if len(b.tds) > 0 {
binary.LittleEndian.PutUint64(idx[pos+24:], rel(b.tdoff[i]))
binary.LittleEndian.PutUint64(idx[pos:], rel(b.tdoff[i]))
pos += 8
}
if len(b.proofsRLP) > 0 {

View file

@ -61,19 +61,18 @@ type Era2 struct {
s *e2store.Reader
m meta // metadata for the era2 file
mu *sync.Mutex
buf [8]byte // buffer reading entry offset
headeroff, bodyoff, receiptsoff, tdoff, proofsoff []uint64 // offsets for each entry type
rootheader uint64 // offset of the root header in the file if present
indstart int64
rootheader uint64 // offset of the root header in the file if present
prooftype uint16
}
func Open(filename string) (*Era2, error) {
f, err := os.Open(filename)
func Open(path string) (*Era2, error) {
f, err := os.Open(path)
if err != nil {
return nil, err
}
e := &Era2{f: f, s: e2store.NewReader(f), mu: new(sync.Mutex)}
e := &Era2{f: f, s: e2store.NewReader(f)}
if err := e.loadIndex(); err != nil {
f.Close()
return nil, err
@ -82,8 +81,6 @@ func Open(filename string) (*Era2, error) {
}
func (e *Era2) Close() error {
e.mu.Lock()
defer e.mu.Unlock()
if e.f == nil {
return nil
}
@ -93,14 +90,10 @@ func (e *Era2) Close() error {
}
func (e *Era2) Start() uint64 {
e.mu.Lock()
defer e.mu.Unlock()
return e.m.start
}
func (e *Era2) Count() uint64 {
e.mu.Lock()
defer e.mu.Unlock()
return e.m.count
}
@ -117,29 +110,29 @@ func (e *Era2) GetBlockByNumber(blockNum uint64) (*types.Block, error) {
}
func (e *Era2) getHeader(blockNum uint64) (*types.Header, error) {
e.mu.Lock()
defer e.mu.Unlock()
if blockNum < e.m.start || blockNum >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", blockNum, e.m.start, e.m.start+e.m.count)
}
r, err := e.snappyPayload(e.headeroff[blockNum-e.m.start])
r, _, err := e.s.ReaderAt(TypeCompressedHeader, int64(e.headeroff[blockNum-e.m.start]))
if err != nil {
return nil, fmt.Errorf("error reading header for block %d: %w", blockNum, err)
return nil, err
}
r = snappy.NewReader(r)
var h types.Header
return &h, rlp.Decode(r, &h)
}
func (e *Era2) getBody(blockNum uint64) (*types.Body, error) {
e.mu.Lock()
defer e.mu.Unlock()
if blockNum < e.m.start || blockNum >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", blockNum, e.m.start, e.m.start+e.m.count)
}
r, err := e.snappyPayload(e.bodyoff[blockNum-e.m.start])
r, _, err := e.s.ReaderAt(TypeCompressedBody, int64(e.bodyoff[blockNum-e.m.start]))
if err != nil {
return nil, fmt.Errorf("error reading body for block %d: %w", blockNum, err)
return nil, err
}
r = snappy.NewReader(r)
var b types.Body
return &b, rlp.Decode(r, &b)
}
@ -153,57 +146,49 @@ func (e *Era2) getTD(blockNum uint64) (*big.Int, error) {
return nil, fmt.Errorf("total-difficulty section not present")
}
start := e.tdoff[blockNum-e.m.start]
buf := make([]byte, 32)
if _, err := e.f.ReadAt(buf, int64(start)); err != nil {
return nil, fmt.Errorf("error reading total difficulty for block %d: %w",
blockNum, err)
r, _, err := e.s.ReaderAt(TypeTotalDifficulty, int64(e.tdoff[blockNum-e.m.start]))
if err != nil {
return nil, err
}
buf, _ := io.ReadAll(r)
td := new(big.Int).SetBytes(reverseOrder(buf))
return td, nil
}
func (e *Era2) GetRawBodyFrameByNumber(n uint64) ([]byte, error) {
if n < e.m.start || n >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", n, e.m.start, e.m.start+e.m.count)
func (e *Era2) GetRawBodyFrameByNumber(blockNum uint64) ([]byte, error) {
if blockNum < e.m.start || blockNum >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", blockNum, e.m.start, e.m.start+e.m.count)
}
start := e.bodyoff[n-e.m.start]
end := e.nextBoundary(n, e.bodyoff, e.receiptsoff) // receipts section is the next safest fallback
length := end - start
sr := io.NewSectionReader(e.f, int64(start), int64(length))
return io.ReadAll(sr)
r, _, err := e.s.ReaderAt(TypeCompressedBody, int64(e.bodyoff[blockNum-e.m.start]))
if err != nil {
return nil, err
}
return io.ReadAll(r)
}
func (e *Era2) GetRawReceiptsFrameByNumber(n uint64) ([]byte, error) {
if n < e.m.start || n >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", n, e.m.start, e.m.start+e.m.count)
func (e *Era2) GetRawReceiptsFrameByNumber(blockNum uint64) ([]byte, error) {
if blockNum < e.m.start || blockNum >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", blockNum, e.m.start, e.m.start+e.m.count)
}
start := e.receiptsoff[n-e.m.start]
end := e.nextBoundary(n, e.receiptsoff, e.tdoff) // TD sec. is next fallback
length := end - start
sr := io.NewSectionReader(e.f, int64(start), int64(length))
return io.ReadAll(sr)
r, _, err := e.s.ReaderAt(TypeCompressedReceipts, int64(e.receiptsoff[blockNum-e.m.start]))
if err != nil {
return nil, err
}
return io.ReadAll(r)
}
func (e *Era2) GetRawProofFrameByNumber(n uint64) ([]byte, error) {
func (e *Era2) GetRawProofFrameByNumber(blockNum uint64) ([]byte, error) {
if len(e.proofsoff) == 0 {
return nil, fmt.Errorf("proofs section not present")
}
if n < e.m.start || n >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", n, e.m.start, e.m.start+e.m.count)
if blockNum < e.m.start || blockNum >= e.m.start+e.m.count {
return nil, fmt.Errorf("block number %d out of range [%d, %d)", blockNum, e.m.start, e.m.start+e.m.count)
}
start := e.proofsoff[n-e.m.start]
// fallback next proof frame → AccRoot header → BlockIndex
var nextSec []uint64
if e.rootheader != 0 {
nextSec = []uint64{e.rootheader}
r, _, err := e.s.ReaderAt(e.prooftype, int64(e.proofsoff[blockNum-e.m.start]))
if err != nil {
return nil, err
}
end := e.nextBoundary(n, e.proofsoff, nextSec)
length := end - start
sr := io.NewSectionReader(e.f, int64(start), int64(length))
return io.ReadAll(sr)
return io.ReadAll(r)
}
func (e *Era2) rawPayload(abs uint64) ([]byte, error) {
@ -216,20 +201,6 @@ func (e *Era2) snappyPayload(abs uint64) (io.Reader, error) {
return snappy.NewReader(sr), nil
}
func (e *Era2) nextBoundary(idx uint64, sameSec []uint64, nextSec []uint64) uint64 {
local := idx - e.m.start
// next frame in the same section
if local+1 < uint64(len(sameSec)) {
return sameSec[local+1]
}
// first frame of the NEXT section (if present)
if len(nextSec) > 0 {
return nextSec[0]
}
// otherwise clamp to start of BlockIndex
return uint64(e.indstart)
}
func (e *Era2) loadIndex() error {
var err error
e.m.filelen, err = e.f.Seek(0, io.SeekEnd)
@ -291,6 +262,14 @@ func (e *Era2) loadIndex() error {
}
}
if len(e.proofsoff) > 0 {
typ, _, perr := e.s.ReadMetadataAt(int64(e.proofsoff[0]))
if perr != nil {
return fmt.Errorf("read proof header: %w", perr)
}
e.prooftype = typ
}
var off int64 = 0 // start at byte-0 of file
for off < e.indstart { // never enter the Block-Index TLV
@ -307,116 +286,79 @@ func (e *Era2) loadIndex() error {
return nil
}
func (e *Era2) BatchRange(first, count uint64, wantHeaders, wantBodies, wantReceipts, wantProofs bool) (hdrs []*types.Header, bods []*types.Body, recs []types.Receipts, prfs [][]byte, err error) {
func (e *Era2) BatchRange(first, count uint64, wantHdr, wantBody, wantRec, wantPrf bool) (hdrs []*types.Header, bods []*types.Body, recs []types.Receipts, prfs [][]byte, err error) {
if count == 0 {
err = fmt.Errorf("count must be greater than 0")
err = fmt.Errorf("count must be > 0")
return
}
if first < e.m.start || first+count > e.m.start+e.m.count {
err = fmt.Errorf("range [%d,%d) out of bounds [%d,%d)", first, first+count, e.m.start, e.m.start+e.m.count)
err = fmt.Errorf("range [%d,%d) out of bounds", first, first+count)
return
}
id0 := first - e.m.start
if wantHeaders {
idx := first - e.m.start
if wantHdr {
hdrs = make([]*types.Header, count)
}
if wantBodies {
if wantBody {
bods = make([]*types.Body, count)
}
if wantReceipts {
if wantRec {
recs = make([]types.Receipts, count)
}
if wantProofs {
if wantPrf {
prfs = make([][]byte, count)
}
stream := func(startOff uint64, endOff uint64, decode func(io.Reader, uint64) error) error {
length := int64(endOff) - int64(startOff)
r := snappy.NewReader(io.NewSectionReader(e.f, int64(startOff), length))
for i := uint64(0); i < count; i++ {
if err := decode(r, i); err != nil {
return err
}
}
return nil
}
id := idx + i
if wantHeaders {
err = stream(e.headeroff[id0], e.bodyoff[0], func(r io.Reader, i uint64) error {
var hdr types.Header
if err := rlp.Decode(r, &hdr); err != nil {
return fmt.Errorf("error decoding header for block %d: %w", first+i, err)
if wantHdr {
r, _, er := e.s.ReaderAt(TypeCompressedHeader, int64(e.headeroff[id]))
if er != nil {
err = er
return
}
hdrs[i] = &hdr
return nil
})
if err != nil {
if er = rlp.Decode(snappy.NewReader(r), &hdrs[i]); er != nil {
err = er
return
}
}
if wantBodies {
err = stream(e.bodyoff[id0], e.receiptsoff[0], func(r io.Reader, i uint64) error {
var body types.Body
if err := rlp.Decode(r, &body); err != nil {
return fmt.Errorf("error decoding body for block %d: %w", first+i, err)
if wantBody {
r, _, er := e.s.ReaderAt(TypeCompressedBody, int64(e.bodyoff[id]))
if er != nil {
err = er
return
}
bods[i] = &body
return nil
})
if err != nil {
if er = rlp.Decode(snappy.NewReader(r), &bods[i]); er != nil {
err = er
return
}
}
if wantReceipts {
var receiptsEnd uint64
if len(e.tdoff) > 0 {
receiptsEnd = e.tdoff[0]
} else {
receiptsEnd = uint64(e.indstart)
if wantRec {
r, _, er := e.s.ReaderAt(TypeCompressedReceipts, int64(e.receiptsoff[id]))
if er != nil {
err = er
return
}
err = stream(e.receiptsoff[id0], receiptsEnd, func(r io.Reader, i uint64) error {
var rct types.Receipts
if err := rlp.Decode(r, &rct); err != nil {
return fmt.Errorf("error decoding receipts for block %d: %w", first+i, err)
}
recs[i] = rct
return nil
})
if err != nil {
if er = rlp.Decode(snappy.NewReader(r), &recs[i]); er != nil {
err = er
return
}
}
if wantProofs {
if len(prfs) == 0 {
err = fmt.Errorf("proofs section is not present")
if wantPrf {
if len(e.proofsoff) == 0 {
err = fmt.Errorf("proofs section not present")
return
}
for i := uint64(0); i < count; i++ {
start := e.proofsoff[id0+i]
var end uint64
if id0+i+1 < uint64(len(e.proofsoff)) {
end = e.proofsoff[id0+i+1] // next proof frame
} else if e.rootheader != 0 {
end = e.rootheader // AccRoot header
} else {
end = uint64(e.indstart) // BlockIndex header
}
length := int64(end - start)
frame, readErr := io.ReadAll(io.NewSectionReader(e.f, int64(start), length))
if readErr != nil {
err = fmt.Errorf("proof read %d: %w", first+i, readErr)
r, _, er := e.s.ReaderAt(e.prooftype, int64(e.proofsoff[id])) // type already checked when writing
if er != nil {
err = er
return
}
prfs[i] = frame
prfs[i], _ = io.ReadAll(r)
}
}
return
}

View file

@ -19,6 +19,7 @@ package era2
import (
"bytes"
"fmt"
"io"
"math/big"
"os"
"testing"
@ -26,6 +27,7 @@ import (
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/rlp"
"github.com/klauspost/compress/snappy"
)
type testchain struct {
@ -37,7 +39,6 @@ type testchain struct {
func TestEra2Builder(t *testing.T) {
t.Parallel()
// Get temp directory.
f, err := os.CreateTemp(t.TempDir(), "era2-test")
if err != nil {
@ -76,22 +77,19 @@ func TestEra2Builder(t *testing.T) {
}
// 3. open reader
t.Logf("filename: %s", f.Name())
era, err := Open(f.Name())
if err != nil {
t.Fatalf("open era: %v", err)
}
defer era.Close()
// 4. verify every block
for i := 0; i < 128; i++ {
bn := uint64(i)
// -- header + body via GetBlockByNumber ------------------------------
gotBlock, err := era.GetBlockByNumber(bn)
if err != nil {
t.Fatalf("get block %d: %v", i, err)
}
if chain.headers[i].Hash() != gotBlock.Header().Hash() {
t.Fatalf("header %d mismatch", i)
}
@ -99,23 +97,31 @@ func TestEra2Builder(t *testing.T) {
t.Fatalf("body %d mismatch", i)
}
// -- raw body frame --------------------------------------------------
rawBody, err := era.GetRawBodyFrameByNumber(bn)
if err != nil {
t.Fatalf("raw body %d: %v", i, err)
}
expectBody := mustEncode(chain.bodies[i])
if !bytes.Contains(rawBody, expectBody) { // frame may include next
decBody, err := io.ReadAll(
snappy.NewReader(bytes.NewReader(rawBody)),
)
if err != nil {
t.Fatalf("snappy decode body %d: %v", i, err)
}
if !bytes.Equal(decBody, mustEncode(chain.bodies[i])) {
t.Fatalf("body frame %d mismatch", i)
}
// -- raw receipts frame ---------------------------------------------
rawRcpt, err := era.GetRawReceiptsFrameByNumber(bn)
if err != nil {
t.Fatalf("raw receipts %d: %v", i, err)
}
expectRcpt := mustEncode(chain.receipts[i])
if !bytes.Contains(rawRcpt, expectRcpt) {
decRcpt, err := io.ReadAll(
snappy.NewReader(bytes.NewReader(rawRcpt)),
)
if err != nil {
t.Fatalf("snappy decode receipts %d: %v", i, err)
}
if !bytes.Equal(decRcpt, mustEncode(chain.receipts[i])) {
t.Fatalf("receipts frame %d mismatch", i)
}
}