live: implement maxKeepBlocks

Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
jsvisa 2024-09-28 05:45:16 +00:00
parent d638fd5426
commit c48a2988e2
2 changed files with 109 additions and 95 deletions

View file

@ -7,9 +7,7 @@ import (
"errors" "errors"
"fmt" "fmt"
"io" "io"
"os"
"path" "path"
"strconv"
"sync" "sync"
"sync/atomic" "sync/atomic"
@ -61,17 +59,16 @@ type live struct {
backend tracing.Backend backend tracing.Backend
kvdb ethdb.Database kvdb ethdb.Database
frdb *rawdb.Freezer frdb *rawdb.Freezer
blockCh chan uint64 freezeCh chan uint64
stopCh chan struct{} 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 finalized uint64
hash common.Hash hash common.Hash
once sync.Once once sync.Once
offsetFile string
enableNonceTracer bool enableNonceTracer bool
} }
@ -80,10 +77,11 @@ type liveTracerConfig struct {
Path string `json:"path"` // Path to the directory where the tracer data will be stored Path string `json:"path"` // Path to the directory where the tracer data will be stored
Config json.RawMessage `json:"config"` Config json.RawMessage `json:"config"`
EnableNonceTracer bool `json:"enableNonceTracer"` EnableNonceTracer bool `json:"enableNonceTracer"`
MaxKeepBlocks uint64 `json:"maxKeepBlocks"` // Maximum number of blocks to keep in the freezer db(the unconfirmaed blocks are not included), 0 means no limit
} }
func toTraceTable(name string) string { func toTraceTable(name string) string {
return name + "_tracers" return name + "_traces"
} }
// encodeNumber encodes a number as big endian uint64 // encodeNumber encodes a number as big endian uint64
@ -131,8 +129,8 @@ func newLive(cfg json.RawMessage, stack tracers.LiveApiRegister, backend tracing
} }
var ( var (
kvpath = path.Join(config.Path, "trace") kvpath = path.Join(config.Path, "kvdb")
frpath = path.Join(config.Path, "freeze") frpath = path.Join(config.Path, "frdb")
) )
kvdb, err := rawdb.NewPebbleDBDatabase(kvpath, 128, 1024, "trace", false, false) kvdb, err := rawdb.NewPebbleDBDatabase(kvpath, 128, 1024, "trace", false, false)
@ -159,37 +157,29 @@ func newLive(cfg json.RawMessage, stack tracers.LiveApiRegister, backend tracing
} }
frozen, err := frdb.Ancients() frozen, err := frdb.Ancients()
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to read the frozen block numbers from the freezer db: %v", err) return nil, fmt.Errorf("failed to read the frozen blocks from the freezer db: %v", err)
} }
l := &live{ l := &live{
backend: backend, backend: backend,
kvdb: kvdb, kvdb: kvdb,
frdb: frdb, frdb: frdb,
blockCh: make(chan uint64, 100), freezeCh: make(chan uint64, 100),
stopCh: make(chan struct{}), stopCh: make(chan struct{}),
tables: tables, tables: tables,
traces: traces, traces: traces,
tracer: t, tracer: t,
offsetFile: path.Join(frpath, "OFFSET"),
enableNonceTracer: config.EnableNonceTracer, enableNonceTracer: config.EnableNonceTracer,
} }
offset := 0
if _, err := os.Stat(l.offsetFile); err == nil || os.IsExist(err) {
data, err := os.ReadFile(l.offsetFile)
if err != nil {
return nil, fmt.Errorf("failed to read the offset from the freezer db: %v", err)
}
offset, err = strconv.Atoi(string(data))
if err != nil {
return nil, fmt.Errorf("failed to convert offset: %v", err)
}
}
log.Info("Open live tracer", "path", config.Path, "offset", offset, "tables", tables)
l.latest.Store(tail + frozen + uint64(offset)) latest := l.getFreezerTail()
l.offset.Store(uint64(offset)) offset := latest - tail - frozen
log.Info("Open live tracer", "path", config.Path, "offset", offset, "tail", tail, "frozen", frozen, "latest", latest, "tables", tables)
// Initialize the latest block number as the sum of the tail, frozen, and offset
l.latest.Store(latest)
l.offset.Store(offset)
hooks := &tracing.Hooks{ hooks := &tracing.Hooks{
OnBlockStart: l.OnBlockStart, OnBlockStart: l.OnBlockStart,
OnBlockEnd: l.OnBlockEnd, OnBlockEnd: l.OnBlockEnd,
@ -218,71 +208,68 @@ func newLive(cfg json.RawMessage, stack tracers.LiveApiRegister, backend tracing
} }
stack.RegisterAPIs(apis) stack.RegisterAPIs(apis)
go l.freeze() go l.freeze(config.MaxKeepBlocks)
return hooks, nil return hooks, nil
} }
func (f *live) OnBlockStart(ev tracing.BlockEvent) { func (l *live) OnBlockStart(ev tracing.BlockEvent) {
// track the latest block number // track the latest block number
blknum := ev.Block.NumberU64() blknum := ev.Block.NumberU64()
f.latest.Store(blknum) l.latest.Store(blknum)
f.hash = ev.Block.Hash() l.hash = ev.Block.Hash()
if ev.Finalized != nil { if ev.Finalized != nil {
f.finalized.Store(ev.Finalized.Number.Uint64()) l.finalized = ev.Finalized.Number.Uint64()
} }
// reset local cache // reset local cache
txs := ev.Block.Transactions().Len() txs := ev.Block.Transactions().Len()
for name := range f.traces { for name := range l.traces {
f.traces[name] = make([]*traceResult, 0, txs) l.traces[name] = make([]*traceResult, 0, txs)
} }
// save the earliest arrived blknum as the offset // Save the earliest arrived blknum as the offset only if offset was not set before
f.once.Do(func() { if swapped := l.offset.CompareAndSwap(0, blknum); swapped {
if _, err := os.Stat(f.offsetFile); err != nil && os.IsNotExist(err) { log.Info("Set live tracer offset to new head", "blknum", blknum)
f.offset.Store(blknum)
os.WriteFile(f.offsetFile, []byte(fmt.Sprintf("%d", blknum)), 0666)
} }
})
} }
func (f *live) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) { func (l *live) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) {
if f.enableNonceTracer { if l.enableNonceTracer {
key := append(from.Bytes(), encodeNumber(tx.Nonce())...) key := append(from.Bytes(), encodeNumber(tx.Nonce())...)
val := tx.Hash().Bytes() val := tx.Hash().Bytes()
if err := f.kvdb.Put(key, val); err != nil { if err := l.kvdb.Put(key, val); err != nil {
log.Warn("Failed to put nonce into kvdb", "err", err) log.Warn("Failed to put nonce into kvdb", "err", err)
} }
} }
f.tracer.OnTxStart(env, tx, from) l.tracer.OnTxStart(env, tx, from)
} }
func (f *live) OnTxEnd(receipt *types.Receipt, err error) { func (l *live) OnTxEnd(receipt *types.Receipt, err error) {
f.tracer.OnTxEnd(receipt, err) l.tracer.OnTxEnd(receipt, err)
for name, tt := range f.tracer.Tracers() { for name, tt := range l.tracer.Tracers() {
trace := &traceResult{} trace := &traceResult{}
result, err := tt.GetResult() result, err := tt.GetResult()
if err != nil { if err != nil {
log.Error("Failed to get tracer results", "number", f.latest.Load(), "error", err) log.Error("Failed to get tracer results", "number", l.latest.Load(), "error", err)
trace.Error = err.Error() trace.Error = err.Error()
} else { } else {
trace.Result = result trace.Result = result
} }
f.traces[name] = append(f.traces[name], trace) l.traces[name] = append(l.traces[name], trace)
} }
} }
func (f *live) OnBlockEnd(err error) { func (l *live) OnBlockEnd(err error) {
if err != nil { if err != nil {
log.Warn("OnBlockEnd", "latest", f.latest.Load(), "error", err) log.Warn("OnBlockEnd", "latest", l.latest.Load(), "error", err)
} }
batch := f.kvdb.NewBatch() batch := l.kvdb.NewBatch()
number := f.latest.Load() number := l.latest.Load()
hash := f.hash hash := l.hash
for name, traces := range f.traces { for name, traces := range l.traces {
data, err := rlp.EncodeToBytes(traces) data, err := rlp.EncodeToBytes(traces)
if err != nil { if err != nil {
log.Error("Failed to marshal traces", "error", err) log.Error("Failed to marshal traces", "error", err)
@ -296,22 +283,22 @@ func (f *live) OnBlockEnd(err error) {
} }
select { select {
case f.blockCh <- f.finalized.Load(): case l.freezeCh <- l.finalized:
default: default:
// Channel is full, log a warning // Channel is full, log a warning
log.Warn("Block channel is full, skipping finalized block notification") log.Warn("Block channel is full, skipping finalized block notification")
} }
} }
func (f *live) readBlockTraces(ctx context.Context, name string, blknum uint64) ([]*traceResult, error) { func (l *live) readBlockTraces(ctx context.Context, name string, blknum uint64) ([]*traceResult, error) {
if blknum > f.latest.Load() { if blknum > l.latest.Load() {
return nil, errors.New("notfound") return nil, errors.New("notfound")
} }
if blknum < f.offset.Load() { if blknum < l.offset.Load() {
return nil, errors.New("historical data not available") return nil, errors.New("historical data not available")
} }
tail := f.getFreezerTail() tail := l.getFreezerTail()
// Determine whether to read from kvdb or frdb // Determine whether to read from kvdb or frdb
var ( var (
@ -320,10 +307,10 @@ func (f *live) readBlockTraces(ctx context.Context, name string, blknum uint64)
) )
if blknum >= tail { if blknum >= tail {
// Data is in kvdb // Data is in kvdb
data, err = f.readFromKVDB(ctx, name, blknum) data, err = l.readFromKVDB(ctx, name, blknum)
} else { } else {
// Data is in frdb // Data is in frdb
data, err = f.readFromFRDB(name, blknum) data, err = l.readFromFRDB(name, blknum)
} }
if err != nil { if err != nil {
return nil, err return nil, err
@ -334,26 +321,26 @@ func (f *live) readBlockTraces(ctx context.Context, name string, blknum uint64)
return traces, err return traces, err
} }
func (f *live) readFromKVDB(ctx context.Context, name string, blknum uint64) ([]byte, error) { func (l *live) readFromKVDB(ctx context.Context, name string, blknum uint64) ([]byte, error) {
header, err := f.backend.HeaderByNumber(ctx, rpc.BlockNumber(blknum)) header, err := l.backend.HeaderByNumber(ctx, rpc.BlockNumber(blknum))
if err != nil { if err != nil {
return nil, err return nil, err
} }
kvKey := toKVKey(name, blknum, header.Hash()) kvKey := toKVKey(name, blknum, header.Hash())
data, err := f.kvdb.Get(kvKey) data, err := l.kvdb.Get(kvKey)
if err != nil { if err != nil {
return nil, fmt.Errorf("traces not found in kvdb for block %d: %w", blknum, err) return nil, fmt.Errorf("traces not found in kvdb for block %d: %w", blknum, err)
} }
return data, err return data, err
} }
func (f *live) readFromFRDB(name string, blknum uint64) ([]byte, error) { func (l *live) readFromFRDB(name string, blknum uint64) ([]byte, error) {
table := toTraceTable(name) table := toTraceTable(name)
var data []byte var data []byte
err := f.frdb.ReadAncients(func(reader ethdb.AncientReaderOp) error { err := l.frdb.ReadAncients(func(reader ethdb.AncientReaderOp) error {
var err error var err error
data, err = reader.Ancient(table, blknum-f.offset.Load()) data, err = reader.Ancient(table, blknum-l.offset.Load())
return err return err
}) })
if err != nil { if err != nil {
@ -362,14 +349,16 @@ func (f *live) readFromFRDB(name string, blknum uint64) ([]byte, error) {
return data, nil return data, nil
} }
func (f *live) Close() { // Close the frdb and kvdb
close(f.stopCh) // TODO: when to close it?
func (l *live) Close() {
close(l.stopCh)
if err := f.kvdb.Close(); err != nil { if err := l.kvdb.Close(); err != nil {
log.Error("Close kvdb failed", "err", err) log.Error("Close kvdb failed", "err", err)
} }
if err := f.frdb.Close(); err != nil { if err := l.frdb.Close(); err != nil {
log.Error("Close freeze db failed", "err", err) log.Error("Close freeze db failed", "err", err)
} }
} }

View file

@ -11,17 +11,17 @@ import (
) )
const ( const (
freezeThreshold = 64 freezeThreshold = 64 // the max number of blocks to freeze in one batch
kvdbTailKey = "FilterFreezerTail" kvdbTailKey = "FilterFreezerTail"
) )
func (f *live) freeze() { func (f *live) freeze(maxKeepBlocks uint64) {
var lastFinalized uint64 var lastFinalized uint64
for { for {
select { select {
case <-f.stopCh: case <-f.stopCh:
return return
case finalizedBlock := <-f.blockCh: case finalizedBlock := <-f.freezeCh:
if finalizedBlock <= lastFinalized { if finalizedBlock <= lastFinalized {
continue continue
} }
@ -30,8 +30,7 @@ func (f *live) freeze() {
tail := f.getFreezerTail() tail := f.getFreezerTail()
// Freeze at most freezeThreshold blocks // Freeze at most freezeThreshold blocks
freezeUpTo := finalizedBlock freezeUpTo := min(finalizedBlock, tail+freezeThreshold)
freezeUpTo = min(freezeUpTo, tail+freezeThreshold)
if freezeUpTo <= tail { if freezeUpTo <= tail {
continue continue
} }
@ -45,9 +44,26 @@ func (f *live) freeze() {
} }
// Update the tail of the freezer // Update the tail of the freezer
if err := f.updateFreezerTail(freezeUpTo); err != nil { tail = freezeUpTo
if err := f.updateFreezerTail(tail); err != nil {
log.Error("Failed to update freezer tail", "error", err) log.Error("Failed to update freezer tail", "error", err)
} }
frozen, _ := f.frdb.Ancients()
offset := f.offset.Load()
// No pruning
if maxKeepBlocks == 0 || frozen <= maxKeepBlocks {
continue
}
// Prune old blocks if necessary
itemsToPrune := min(freezeThreshold, frozen-maxKeepBlocks)
head := offset + itemsToPrune - 1
log.Info("Prune old blocks", "pruned", itemsToPrune, "from", offset, "to", head)
if err := f.pruneBlocksFromFreezer(frozen-itemsToPrune, head); err != nil {
log.Error("Failed to prune blocks from freezer", "error", err)
}
} }
} }
} }
@ -140,3 +156,12 @@ func (f *live) deleteKVDBEntriesWithPrefix(blknum uint64) error {
return nil return nil
} }
func (f *live) pruneBlocksFromFreezer(items, head uint64) error {
if _, err := f.frdb.TruncateHead(items); err != nil {
return err
}
// Head should be in sync with the on-mem offset
f.offset.Store(head)
return nil
}