freeze into frdb

Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
jsvisa 2024-08-13 04:46:39 +00:00
parent 68b63d48a2
commit 3e598cb740
2 changed files with 233 additions and 61 deletions

View file

@ -41,11 +41,14 @@ type filter struct {
backend tracers.Backend backend tracers.Backend
kvdb ethdb.Database kvdb ethdb.Database
frdb *rawdb.Freezer frdb *rawdb.Freezer
blockCh chan uint64
stopCh chan struct{}
tables map[string]bool tables map[string]bool
traces map[string][]*traceResult traces map[string][]*traceResult
tracer *native.MuxTracer tracer *native.MuxTracer
latest atomic.Uint64 latest atomic.Uint64
offset atomic.Uint64 offset atomic.Uint64
finalized atomic.Uint64
hash common.Hash hash common.Hash
once sync.Once once sync.Once
offsetFile string offsetFile string
@ -132,6 +135,8 @@ func newFilter(cfg json.RawMessage, backend tracers.Backend) (*tracing.Hooks, []
backend: backend, backend: backend,
kvdb: kvdb, kvdb: kvdb,
frdb: frdb, frdb: frdb,
blockCh: make(chan uint64, 100),
stopCh: make(chan struct{}),
tables: tables, tables: tables,
traces: traces, traces: traces,
tracer: t, tracer: t,
@ -176,20 +181,25 @@ func newFilter(cfg json.RawMessage, backend tracers.Backend) (*tracing.Hooks, []
Service: &filterAPI{backend: backend, filter: f}, Service: &filterAPI{backend: backend, filter: f},
}, },
} }
// Initialize head if it doesn't exist
head, _ := f.getFreezerHeadTail()
if head == 0 {
f.updateFreezerHead(f.latest.Load())
}
go f.freeze()
return hooks, apis, nil return hooks, apis, nil
} }
func (f *filter) OnBlockStart(ev tracing.BlockEvent) { func (f *filter) OnBlockStart(ev tracing.BlockEvent) {
// track the latest block number // track the latest block number
blknum := ev.Block.NumberU64() blknum := ev.Block.NumberU64()
// latest := f.latest.Load()
// if blknum < latest {
// // TODO: handle the case of setHead
// log.Error("OnBlockStart received an old block", "latest", latest, "number", blknum)
// return
// }
f.latest.Store(blknum) f.latest.Store(blknum)
f.hash = ev.Block.Hash() f.hash = ev.Block.Hash()
if ev.Finalized != nil {
f.finalized.Store(ev.Finalized.Number.Uint64())
}
// reset local cache // reset local cache
txs := ev.Block.Transactions().Len() txs := ev.Block.Transactions().Len()
@ -204,18 +214,6 @@ func (f *filter) OnBlockStart(ev tracing.BlockEvent) {
os.WriteFile(f.offsetFile, []byte(fmt.Sprintf("%d", blknum)), 0666) os.WriteFile(f.offsetFile, []byte(fmt.Sprintf("%d", blknum)), 0666)
} }
}) })
// // truncate the freezer db if the block number is less than the latest
// if blknum <= latest {
// frozen, _ := f.frdb.Ancients()
// offset := f.offset.Load()
// log.Info("Reorg detected", "number", blknum, "latest", latest, "offset", offset, "frozen", frozen)
// if _, err := f.frdb.TruncateHead(blknum - offset); err != nil {
// log.Error("Failed to truncate filter tracer db", "error", err)
// // TODO: how to handle this error?
// return
// }
// }
} }
func (f *filter) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) { func (f *filter) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) {
@ -240,7 +238,7 @@ func (f *filter) OnTxEnd(receipt *types.Receipt, err error) {
func (f *filter) OnBlockEnd(err error) { func (f *filter) OnBlockEnd(err error) {
if err != nil { if err != nil {
log.Warn("OnBlockEnd", "latest", f.latest.Load(), "err", err) log.Warn("OnBlockEnd", "latest", f.latest.Load(), "error", err)
} }
batch := f.kvdb.NewBatch() batch := f.kvdb.NewBatch()
@ -250,32 +248,60 @@ func (f *filter) OnBlockEnd(err error) {
data, err := json.Marshal(traces) data, err := json.Marshal(traces)
if err != nil { if err != nil {
log.Error("Failed to marshal traces", "error", err) log.Error("Failed to marshal traces", "error", err)
break
} }
batch.Put(toKVKey(name, number, hash), data) batch.Put(toKVKey(name, number, hash), data)
} }
if err := batch.Write(); err != nil { if err := batch.Write(); err != nil {
log.Error("Failed to write", "err", err) log.Error("Failed to write", "error", err)
return
} }
// f.frdb.ModifyAncients(func(w ethdb.AncientWriteOp) error { select {
// latest := f.latest.Load() case f.blockCh <- f.finalized.Load():
// offset := f.offset.Load() default:
// number := latest - offset // Channel is full, log a warning
// for name, traces := range f.traces { log.Warn("Block channel is full, skipping finalized block notification")
// data, err := json.Marshal(traces) }
// if err != nil {
// log.Error("Failed to marshal traces", "error", err)
// }
// table := toTraceTable(name)
// if err := w.AppendRaw(table, number, data); err != nil {
// log.Error("Failed to write block traces", "number", latest, "table", table, "error", err)
// }
// }
// return nil
// })
} }
func (f *filter) readBlockTraces(ctx context.Context, name string, blknum uint64) ([]*traceResult, error) { func (f *filter) readBlockTraces(ctx context.Context, name string, blknum uint64) ([]*traceResult, error) {
if blknum > f.latest.Load() {
return nil, errors.New("notfound")
}
if blknum < f.offset.Load() {
return nil, errors.New("historical data not available")
}
_, tail := f.getFreezerHeadTail()
// If tail is 0 (not found in kvdb), use the offset
if tail == 0 {
tail = f.offset.Load()
}
// Determine whether to read from kvdb or frdb
var (
data []byte
err error
)
if blknum >= tail {
// Data is in kvdb
data, err = f.readFromKVDB(ctx, name, blknum)
} else {
// Data is in frdb
data, err = f.readFromFRDB(name, blknum)
}
if err != nil {
return nil, err
}
var traces []*traceResult
err = json.Unmarshal(data, &traces)
return traces, err
}
func (f *filter) readFromKVDB(ctx context.Context, name string, blknum uint64) ([]byte, error) {
header, err := f.backend.HeaderByNumber(ctx, rpc.BlockNumber(blknum)) header, err := f.backend.HeaderByNumber(ctx, rpc.BlockNumber(blknum))
if err != nil { if err != nil {
return nil, err return nil, err
@ -284,36 +310,28 @@ func (f *filter) readBlockTraces(ctx context.Context, name string, blknum uint64
kvKey := toKVKey(name, blknum, header.Hash()) kvKey := toKVKey(name, blknum, header.Hash())
data, err := f.kvdb.Get(kvKey) data, err := f.kvdb.Get(kvKey)
if err != nil { if err != nil {
return nil, err return nil, fmt.Errorf("traces not found in kvdb for block %d: %w", blknum, err)
}
return data, err
} }
var traces []*traceResult
err = json.Unmarshal(data, &traces)
return traces, err
// table := toTraceTable(name) func (f *filter) readFromFRDB(name string, blknum uint64) ([]byte, error) {
// if _, ok := f.tables[table]; !ok { table := toTraceTable(name)
// return nil, errors.New("tracer not found") var data []byte
// } err := f.frdb.ReadAncients(func(reader ethdb.AncientReaderOp) error {
// var err error
// if blknum < f.offset.Load() || blknum > f.latest.Load() { data, err = reader.Ancient(table, blknum-f.offset.Load())
// return nil, nil return err
// } })
// if err != nil {
// var data []byte return nil, fmt.Errorf("traces not found in frdb for block %d: %w", blknum, err)
// err := f.frdb.ReadAncients(func(reader ethdb.AncientReaderOp) error { }
// var err error return data, nil
// data, err = reader.Ancient(table, blknum-f.offset.Load())
// return err
// })
// if err != nil {
// return nil, err
// }
// var traces []*traceResult
// err = json.Unmarshal(data, &traces)
// return traces, err
} }
func (f *filter) Close() { func (f *filter) Close() {
close(f.stopCh)
if err := f.kvdb.Close(); err != nil { if err := f.kvdb.Close(); err != nil {
log.Error("Close kvdb failed", "err", err) log.Error("Close kvdb failed", "err", err)
} }

View file

@ -0,0 +1,154 @@
package live
import (
"context"
"encoding/binary"
"fmt"
"github.com/ethereum/go-ethereum/ethdb"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/rpc"
)
const (
freezeThreshold = 64
kvdbHeadKey = "FilterFreezerHead"
kvdbTailKey = "FilterFreezerTail"
)
func (f *filter) freeze() {
var lastFinalized uint64
for {
select {
case <-f.stopCh:
return
case finalizedBlock := <-f.blockCh:
if finalizedBlock <= lastFinalized {
continue
}
lastFinalized = finalizedBlock
head, tail := f.getFreezerHeadTail()
// If tail is 0 (not found in kvdb), use the offset
if tail == 0 {
tail = f.offset.Load()
}
// Freeze at most freezeThreshold blocks
freezeUpTo := finalizedBlock
freezeUpTo = min(freezeUpTo, tail+freezeThreshold)
log.Info("Move traces from kvdb to frdb", "from", tail, "to", freezeUpTo)
for blknum := tail; blknum < freezeUpTo; blknum++ {
if err := f.moveBlockToFreezer(blknum); err != nil {
log.Error("Failed to move block to freezer", "block", blknum, "error", err)
break
}
}
// Update head and tail
if freezeUpTo > tail {
if err := f.updateFreezerTail(freezeUpTo); err != nil {
log.Error("Failed to update freezer tail", "error", err)
}
}
if freezeUpTo > head {
if err := f.updateFreezerHead(freezeUpTo); err != nil {
log.Error("Failed to update freezer head", "error", err)
}
}
}
}
}
func (f *filter) getFreezerHeadTail() (head, tail uint64) {
headBytes, _ := f.kvdb.Get([]byte(kvdbHeadKey))
tailBytes, _ := f.kvdb.Get([]byte(kvdbTailKey))
if len(headBytes) > 0 {
head = binary.BigEndian.Uint64(headBytes)
}
if len(tailBytes) > 0 {
tail = binary.BigEndian.Uint64(tailBytes)
}
return
}
func (f *filter) updateFreezerHead(head uint64) error {
headBytes := make([]byte, 8)
binary.BigEndian.PutUint64(headBytes, head)
return f.kvdb.Put([]byte(kvdbHeadKey), headBytes)
}
func (f *filter) updateFreezerTail(tail uint64) error {
tailBytes := make([]byte, 8)
binary.BigEndian.PutUint64(tailBytes, tail)
return f.kvdb.Put([]byte(kvdbTailKey), tailBytes)
}
func (f *filter) moveBlockToFreezer(blknum uint64) error {
header, err := f.backend.HeaderByNumber(context.Background(), rpc.BlockNumber(blknum))
if err != nil {
return err
}
offset := f.offset.Load()
for name := range f.tracer.Tracers() {
kvKey := toKVKey(name, blknum, header.Hash())
data, err := f.kvdb.Get(kvKey)
if err != nil {
return err
}
table := toTraceTable(name)
n, err := f.frdb.ModifyAncients(func(op ethdb.AncientWriteOp) error {
return op.AppendRaw(table, blknum-offset, data)
})
if err != nil {
return err
}
log.Info("Move from kvdb to frdb", "blknum", blknum, "size", n)
// Delete all entries for this prefix from kvdb, ignore error
prefix := append([]byte(name), encodeBlockNumber(blknum)...)
if err := f.deleteKVDBEntriesWithPrefix(prefix); err != nil {
log.Error("Failed to delete entries from kvdb", "error", err)
}
}
return nil
}
func (f *filter) deleteKVDBEntriesWithPrefix(prefix []byte) error {
batch := f.kvdb.NewBatch()
it := f.kvdb.NewIterator(prefix, nil)
defer it.Release()
for it.Next() {
if err := batch.Delete(it.Key()); err != nil {
return fmt.Errorf("failed to add delete operation to batch: %w", err)
}
// Write batch if it's getting too large
if batch.ValueSize() > ethdb.IdealBatchSize {
if err := batch.Write(); err != nil {
return fmt.Errorf("failed to write batch: %w", err)
}
batch.Reset()
}
}
if err := it.Error(); err != nil {
return fmt.Errorf("iterator error: %w", err)
}
// Write any remaining batch operations
if batch.ValueSize() > 0 {
if err := batch.Write(); err != nil {
return fmt.Errorf("failed to write final batch: %w", err)
}
}
return nil
}