internal/era: add ReaderAt func so entry value can be read lazily

Co-authored-by: lightclient <lightclient@protonmail.com>
Co-authored-by: Martin Holst Swende <martin@swende.se>
This commit is contained in:
lightclient@protonmail.com 2023-06-05 20:20:25 +02:00 committed by lightclient
parent 64513d7d3d
commit b8fc70e8f4
No known key found for this signature in database
GPG key ID: 75C916AFEE20183E
3 changed files with 39 additions and 32 deletions

View file

@ -118,6 +118,24 @@ func (r *Reader) ReadAt(entry *Entry, off int64) (int, error) {
return int(headerSize + length), nil return int(headerSize + length), nil
} }
// ReaderAt returns an io.Reader delivering value data for the entry at
// the specified offset. If the entry type does not match the expected type, an
// error is returned.
func (r *Reader) ReaderAt(expectedType uint16, off int64) (io.Reader, int, error) {
// problem = need to return length+headerSize not just value length via section reader
typ, length, err := r.ReadMetadataAt(off)
if err != nil {
return nil, headerSize, err
}
if typ != expectedType {
return nil, headerSize, fmt.Errorf("wrong type, want %d have %d", expectedType, typ)
}
if length > valueSizeLimit {
return nil, headerSize, fmt.Errorf("item larger than item size limit %d: have %d", valueSizeLimit, length)
}
return io.NewSectionReader(r.r, off+headerSize, int64(length)), headerSize + int(length), nil
}
// LengthAt reads the header at off and returns the total length of the entry, // LengthAt reads the header at off and returns the total length of the entry,
// including header. // including header.
func (r *Reader) LengthAt(off int64) (int64, error) { func (r *Reader) LengthAt(off int64) (int64, error) {

View file

@ -17,7 +17,6 @@
package era package era
import ( import (
"bytes"
"encoding/binary" "encoding/binary"
"fmt" "fmt"
"io" "io"
@ -134,7 +133,7 @@ func (e *Era) GetBlockByNumber(num uint64) (*types.Block, error) {
if err != nil { if err != nil {
return nil, err return nil, err
} }
r, n, err := newSnappyReader(e.s, off) r, n, err := newSnappyReader(e.s, TypeCompressedHeader, off)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@ -142,8 +141,8 @@ func (e *Era) GetBlockByNumber(num uint64) (*types.Block, error) {
if err := rlp.Decode(r, &header); err != nil { if err := rlp.Decode(r, &header); err != nil {
return nil, err return nil, err
} }
off += int64(n) off += n
r, _, err = newSnappyReader(e.s, off) r, _, err = newSnappyReader(e.s, TypeCompressedBody, off)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@ -170,7 +169,7 @@ func (e *Era) InitialTD() (*big.Int, error) {
r io.Reader r io.Reader
header types.Header header types.Header
rawTd []byte rawTd []byte
n int n int64
off int64 off int64
err error err error
) )
@ -179,13 +178,13 @@ func (e *Era) InitialTD() (*big.Int, error) {
if off, err = e.readOffset(e.m.start); err != nil { if off, err = e.readOffset(e.m.start); err != nil {
return nil, err return nil, err
} }
if r, n, err = newSnappyReader(e.s, off); err != nil { if r, n, err = newSnappyReader(e.s, TypeCompressedHeader, off); err != nil {
return nil, err return nil, err
} }
if err := rlp.Decode(r, header); err != nil { if err := rlp.Decode(r, header); err != nil {
return nil, err return nil, err
} }
off += int64(n) off += n
// Skip over next two records. // Skip over next two records.
for i := 0; i < 2; i++ { for i := 0; i < 2; i++ {
@ -197,7 +196,7 @@ func (e *Era) InitialTD() (*big.Int, error) {
} }
// Read total difficulty after first block. // Read total difficulty after first block.
if r, _, err = newReader(e.s, off); err != nil { if r, _, err = e.s.ReaderAt(TypeTotalDifficulty, off); err != nil {
return nil, err return nil, err
} }
if err := rlp.Decode(r, rawTd); err != nil { if err := rlp.Decode(r, rawTd); err != nil {
@ -237,23 +236,13 @@ func (e *Era) readOffset(n uint64) (int64, error) {
return offOffset + 8 + int64(binary.LittleEndian.Uint64(e.buf[:])), nil return offOffset + 8 + int64(binary.LittleEndian.Uint64(e.buf[:])), nil
} }
// newReader returns an io.Reader for the e2store entry value at off.
func newReader(e *e2store.Reader, off int64) (io.Reader, int, error) {
var (
entry e2store.Entry
n int
err error
)
if n, err = e.ReadAt(&entry, off); err != nil {
return nil, n, err
}
return bytes.NewReader(entry.Value), n, nil
}
// newReader returns a snappy.Reader for the e2store entry value at off. // newReader returns a snappy.Reader for the e2store entry value at off.
func newSnappyReader(e *e2store.Reader, off int64) (io.Reader, int, error) { func newSnappyReader(e *e2store.Reader, expectedType uint16, off int64) (io.Reader, int64, error) {
r, n, err := newReader(e, off) r, n, err := e.ReaderAt(expectedType, off)
return snappy.NewReader(r), n, err if err != nil {
return nil, 0, err
}
return snappy.NewReader(r), int64(n), err
} }
// clearBuffer zeroes out the buffer. // clearBuffer zeroes out the buffer.

View file

@ -129,20 +129,20 @@ func (it *RawIterator) Next() bool {
it.err = err it.err = err
return false return false
} }
var n int var n int64
if it.Header, n, it.err = newSnappyReader(it.e.s, off); it.err != nil { if it.Header, n, it.err = newSnappyReader(it.e.s, TypeCompressedHeader, off); it.err != nil {
return true return true
} }
off += int64(n) off += n
if it.Body, n, it.err = newSnappyReader(it.e.s, off); it.err != nil { if it.Body, n, it.err = newSnappyReader(it.e.s, TypeCompressedBody, off); it.err != nil {
return true return true
} }
off += int64(n) off += n
if it.Receipts, n, it.err = newSnappyReader(it.e.s, off); it.err != nil { if it.Receipts, n, it.err = newSnappyReader(it.e.s, TypeCompressedReceipts, off); it.err != nil {
return true return true
} }
off += int64(n) off += n
if it.TotalDifficulty, _, it.err = newReader(it.e.s, off); it.err != nil { if it.TotalDifficulty, _, it.err = it.e.s.ReaderAt(TypeTotalDifficulty, off); it.err != nil {
return true return true
} }
it.next += 1 it.next += 1