eth/tracers: declare muxTracer as public

Signed-off-by: jsvisa <delweng@gmail.com>

fix(eth/tracers): clear callstack in OnTxStart

Signed-off-by: jsvisa <delweng@gmail.com>

eth/trace: add filter storage

Signed-off-by: jsvisa <delweng@gmail.com>

eth/tracers: register trace API

Signed-off-by: jsvisa <delweng@gmail.com>

eth/tracers: add filter_api

Signed-off-by: jsvisa <delweng@gmail.com>

eth/tracers: add trace_block rpc

Signed-off-by: jsvisa <delweng@gmail.com>

eth/tracers/native: MuxTracer as public

Signed-off-by: jsvisa <delweng@gmail.com>

eth/tracers: put/get trace by tracer name

Signed-off-by: jsvisa <delweng@gmail.com>

fix: the opposit blknum cmp f.latest

Signed-off-by: jsvisa <delweng@gmail.com>

feat: no need to handle reorg if blknum is the same

Signed-off-by: jsvisa <delweng@gmail.com>

fix: restore f.latest if offset

Signed-off-by: jsvisa <delweng@gmail.com>

feat: always update f.latest

Signed-off-by: jsvisa <delweng@gmail.com>

fix blknum=latest

Signed-off-by: jsvisa <delweng@gmail.com>

fix: latest use atomic.Uint64

Signed-off-by: jsvisa <delweng@gmail.com>

fix: f.latest use atomic

Signed-off-by: jsvisa <delweng@gmail.com>

fix: ignore if block < latest

Signed-off-by: jsvisa <delweng@gmail.com>

ignore all later ops is ignored

Signed-off-by: jsvisa <delweng@gmail.com>

ignored -> skipped

Signed-off-by: jsvisa <delweng@gmail.com>

db -> frdb

Signed-off-by: jsvisa <delweng@gmail.com>

add kvdb

Signed-off-by: jsvisa <delweng@gmail.com>

no skipper

Signed-off-by: jsvisa <delweng@gmail.com>

wip: save into kvdb

Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
jsvisa 2024-07-28 15:52:19 +00:00
parent 9ae1444d48
commit d5632f55fb
4 changed files with 387 additions and 37 deletions

307
eth/tracers/live/filter.go Normal file
View file

@ -0,0 +1,307 @@
package live
import (
"encoding/binary"
"encoding/json"
"errors"
"fmt"
"os"
"path"
"strconv"
"sync"
"sync/atomic"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/rawdb"
"github.com/ethereum/go-ethereum/core/tracing"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/eth/tracers"
"github.com/ethereum/go-ethereum/eth/tracers/native"
"github.com/ethereum/go-ethereum/ethdb"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/rpc"
)
func init() {
tracers.LiveDirectory.Register("filter", newFilter)
}
const (
tableSize = 2 * 1024 * 1024 * 1024
)
type traceResult struct {
TxHash common.Hash `json:"txHash,omitempty"` // transaction hash
Result interface{} `json:"result,omitempty"` // Trace results produced by the tracer
Error string `json:"error,omitempty"` // Trace failure produced by the tracer
}
type filter struct {
kvdb ethdb.Database
frdb *rawdb.Freezer
tables map[string]bool
traces map[string][]*traceResult
tracer *native.MuxTracer
latest atomic.Uint64
offset atomic.Uint64
hash common.Hash
once sync.Once
offsetFile string
}
type filterTracerConfig struct {
Path string `json:"path"` // Path to the directory where the tracer logs will be stored
Config json.RawMessage `json:"config"`
}
func toTraceTable(name string) string {
return name + "_tracers"
}
// encodeBlockNumber encodes a block number as big endian uint64
func encodeBlockNumber(number uint64) []byte {
enc := make([]byte, 8)
binary.BigEndian.PutUint64(enc, number)
return enc
}
func toKVKey(name string, number uint64, hash common.Hash) []byte {
var key []byte
switch name {
case "callTracer":
key = []byte("C")
case "flatCallTracer":
key = []byte("P")
}
key = append(append(key, encodeBlockNumber(number)...), hash.Bytes()...)
return key
}
func newFilter(cfg json.RawMessage, backend tracers.Backend) (*tracing.Hooks, []rpc.API, error) {
var config filterTracerConfig
if cfg != nil {
if err := json.Unmarshal(cfg, &config); err != nil {
return nil, nil, fmt.Errorf("failed to parse config: %v", err)
}
}
if config.Path == "" {
return nil, nil, errors.New("filter tracer output path is required")
}
t, err := native.NewMuxTracer(config.Config)
if err != nil {
return nil, nil, err
}
var (
kvpath = path.Join(config.Path, "trace")
frpath = path.Join(config.Path, "freeze")
)
kvdb, err := rawdb.NewPebbleDBDatabase(kvpath, 128, 1024, "trace", false, false)
if err != nil {
return nil, nil, err
}
muxTracers := t.Tracers()
tables := make(map[string]bool, len(muxTracers))
traces := make(map[string][]*traceResult, len(muxTracers))
for name := range muxTracers {
tables[toTraceTable(name)] = false
traces[name] = nil
}
frdb, err := rawdb.NewFreezer(frpath, "trace", false, tableSize, tables)
if err != nil {
return nil, nil, fmt.Errorf("failed to create trace freezer db: %v", err)
}
tail, err := frdb.Tail()
if err != nil {
return nil, nil, fmt.Errorf("failed to read the tail block number from the freezer db: %v", err)
}
frozen, err := frdb.Ancients()
if err != nil {
return nil, nil, fmt.Errorf("failed to read the frozen block numbers from the freezer db: %v", err)
}
f := &filter{
kvdb: kvdb,
frdb: frdb,
tables: tables,
traces: traces,
tracer: t,
offsetFile: path.Join(frpath, "OFFSET"),
}
offset := 0
if _, err := os.Stat(f.offsetFile); err == nil || os.IsExist(err) {
data, err := os.ReadFile(f.offsetFile)
if err != nil {
return nil, 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, nil, fmt.Errorf("failed to convert offset: %v", err)
}
}
log.Info("Open filter tracer", "path", config.Path, "offset", offset, "tables", tables)
f.latest.Store(tail + frozen + uint64(offset))
f.offset.Store(uint64(offset))
hooks := &tracing.Hooks{
OnBlockStart: f.OnBlockStart,
OnBlockEnd: f.OnBlockEnd,
OnTxStart: f.OnTxStart,
OnTxEnd: f.OnTxEnd,
// reuse the mux's hooks
OnEnter: t.OnEnter,
OnExit: t.OnExit,
OnOpcode: t.OnOpcode,
OnFault: t.OnFault,
OnGasChange: t.OnGasChange,
OnBalanceChange: t.OnBalanceChange,
OnNonceChange: t.OnNonceChange,
OnCodeChange: t.OnCodeChange,
OnStorageChange: t.OnStorageChange,
OnLog: t.OnLog,
}
apis := []rpc.API{
{
Namespace: "trace",
Service: &filterAPI{filter: f},
},
}
return hooks, apis, nil
}
func (f *filter) OnBlockStart(ev tracing.BlockEvent) {
// track the latest block number
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.hash = ev.Block.Hash()
// reset local cache
txs := ev.Block.Transactions().Len()
for name := range f.traces {
f.traces[name] = make([]*traceResult, 0, txs)
}
// save the earliest arrived blknum as the offset
f.once.Do(func() {
if _, err := os.Stat(f.offsetFile); err != nil && os.IsNotExist(err) {
f.offset.Store(blknum)
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) {
f.tracer.OnTxStart(env, tx, from)
}
func (f *filter) OnTxEnd(receipt *types.Receipt, err error) {
f.tracer.OnTxEnd(receipt, err)
for name, tt := range f.tracer.Tracers() {
trace := &traceResult{TxHash: receipt.TxHash}
result, err := tt.GetResult()
if err != nil {
log.Error("Failed to get tracer results", "number", f.latest.Load(), "error", err)
trace.Error = err.Error()
} else {
trace.Result = result
}
f.traces[name] = append(f.traces[name], trace)
}
}
func (f *filter) OnBlockEnd(err error) {
if err != nil {
log.Warn("OnBlockEnd", "latest", f.latest.Load(), "err", err)
}
batch := f.kvdb.NewBatch()
number := f.latest.Load()
hash := f.hash
for name, traces := range f.traces {
data, err := json.Marshal(traces)
if err != nil {
log.Error("Failed to marshal traces", "error", err)
}
batch.Put(toKVKey(name, number, hash), data)
}
if err := batch.Write(); err != nil {
log.Error("Failed to write", "err", err)
}
// f.frdb.ModifyAncients(func(w ethdb.AncientWriteOp) error {
// latest := f.latest.Load()
// offset := f.offset.Load()
// number := latest - offset
// for name, traces := range f.traces {
// 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(name string, blknum uint64) ([]*traceResult, error) {
table := toTraceTable(name)
if _, ok := f.tables[table]; !ok {
return nil, errors.New("tracer not found")
}
if blknum < f.offset.Load() || blknum > f.latest.Load() {
return nil, nil
}
var data []byte
err := f.frdb.ReadAncients(func(reader ethdb.AncientReaderOp) error {
var err error
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() {
if err := f.kvdb.Close(); err != nil {
log.Error("Close kvdb failed", "err", err)
}
if err := f.frdb.Close(); err != nil {
log.Error("Close freeze db failed", "err", err)
}
}

View file

@ -0,0 +1,33 @@
package live
import (
"context"
"github.com/ethereum/go-ethereum/rpc"
)
type filterAPI struct {
filter *filter
}
type traceConfig struct {
Tracer string `json:"tracer"`
}
var defaultTraceConfig = &traceConfig{
Tracer: "callTracer",
}
func (api *filterAPI) Block(ctx context.Context, blockNr rpc.BlockNumber, cfg *traceConfig) ([]*traceResult, error) {
blknum := uint64(blockNr.Int64())
if blockNr == rpc.LatestBlockNumber {
blknum = api.filter.latest.Load()
}
tracer := defaultTraceConfig.Tracer
if cfg != nil {
tracer = cfg.Tracer
}
return api.filter.readBlockTraces(tracer, blknum)
}

View file

@ -217,6 +217,7 @@ func (t *callTracer) captureEnd(output []byte, gasUsed uint64, err error, revert
}
func (t *callTracer) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) {
t.callstack = nil
t.gasLimit = tx.Gas()
}

View file

@ -30,33 +30,18 @@ func init() {
tracers.DefaultDirectory.Register("muxTracer", newMuxTracer, false)
}
// muxTracer is a go implementation of the Tracer interface which
// MuxTracer is a go implementation of the Tracer interface which
// runs multiple tracers in one go.
type muxTracer struct {
names []string
tracers []*tracers.Tracer
type MuxTracer struct {
tracers map[string]*tracers.Tracer
}
// newMuxTracer returns a new mux tracer.
func newMuxTracer(ctx *tracers.Context, cfg json.RawMessage) (*tracers.Tracer, error) {
var config map[string]json.RawMessage
if cfg != nil {
if err := json.Unmarshal(cfg, &config); err != nil {
return nil, err
}
}
objects := make([]*tracers.Tracer, 0, len(config))
names := make([]string, 0, len(config))
for k, v := range config {
t, err := tracers.DefaultDirectory.New(k, ctx, v)
t, err := NewMuxTracer(cfg)
if err != nil {
return nil, err
}
objects = append(objects, t)
names = append(names, k)
}
t := &muxTracer{names: names, tracers: objects}
return &tracers.Tracer{
Hooks: &tracing.Hooks{
OnTxStart: t.OnTxStart,
@ -77,7 +62,27 @@ func newMuxTracer(ctx *tracers.Context, cfg json.RawMessage) (*tracers.Tracer, e
}, nil
}
func (t *muxTracer) OnOpcode(pc uint64, op byte, gas, cost uint64, scope tracing.OpContext, rData []byte, depth int, err error) {
// NewMuxTracer returns a new mux tracer.
func NewMuxTracer(cfg json.RawMessage) (*MuxTracer, error) {
var config map[string]json.RawMessage
if cfg != nil {
if err := json.Unmarshal(cfg, &config); err != nil {
return nil, err
}
}
objects := make(map[string]*tracers.Tracer, len(config))
for k, v := range config {
t, err := tracers.DefaultDirectory.New(k, nil, v)
if err != nil {
return nil, err
}
objects[k] = t
}
return &MuxTracer{tracers: objects}, nil
}
func (t *MuxTracer) OnOpcode(pc uint64, op byte, gas, cost uint64, scope tracing.OpContext, rData []byte, depth int, err error) {
for _, t := range t.tracers {
if t.OnOpcode != nil {
t.OnOpcode(pc, op, gas, cost, scope, rData, depth, err)
@ -85,7 +90,7 @@ func (t *muxTracer) OnOpcode(pc uint64, op byte, gas, cost uint64, scope tracing
}
}
func (t *muxTracer) OnFault(pc uint64, op byte, gas, cost uint64, scope tracing.OpContext, depth int, err error) {
func (t *MuxTracer) OnFault(pc uint64, op byte, gas, cost uint64, scope tracing.OpContext, depth int, err error) {
for _, t := range t.tracers {
if t.OnFault != nil {
t.OnFault(pc, op, gas, cost, scope, depth, err)
@ -93,7 +98,7 @@ func (t *muxTracer) OnFault(pc uint64, op byte, gas, cost uint64, scope tracing.
}
}
func (t *muxTracer) OnGasChange(old, new uint64, reason tracing.GasChangeReason) {
func (t *MuxTracer) OnGasChange(old, new uint64, reason tracing.GasChangeReason) {
for _, t := range t.tracers {
if t.OnGasChange != nil {
t.OnGasChange(old, new, reason)
@ -101,7 +106,7 @@ func (t *muxTracer) OnGasChange(old, new uint64, reason tracing.GasChangeReason)
}
}
func (t *muxTracer) OnEnter(depth int, typ byte, from common.Address, to common.Address, input []byte, gas uint64, value *big.Int) {
func (t *MuxTracer) OnEnter(depth int, typ byte, from common.Address, to common.Address, input []byte, gas uint64, value *big.Int) {
for _, t := range t.tracers {
if t.OnEnter != nil {
t.OnEnter(depth, typ, from, to, input, gas, value)
@ -109,7 +114,7 @@ func (t *muxTracer) OnEnter(depth int, typ byte, from common.Address, to common.
}
}
func (t *muxTracer) OnExit(depth int, output []byte, gasUsed uint64, err error, reverted bool) {
func (t *MuxTracer) OnExit(depth int, output []byte, gasUsed uint64, err error, reverted bool) {
for _, t := range t.tracers {
if t.OnExit != nil {
t.OnExit(depth, output, gasUsed, err, reverted)
@ -117,7 +122,7 @@ func (t *muxTracer) OnExit(depth int, output []byte, gasUsed uint64, err error,
}
}
func (t *muxTracer) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) {
func (t *MuxTracer) OnTxStart(env *tracing.VMContext, tx *types.Transaction, from common.Address) {
for _, t := range t.tracers {
if t.OnTxStart != nil {
t.OnTxStart(env, tx, from)
@ -125,7 +130,7 @@ func (t *muxTracer) OnTxStart(env *tracing.VMContext, tx *types.Transaction, fro
}
}
func (t *muxTracer) OnTxEnd(receipt *types.Receipt, err error) {
func (t *MuxTracer) OnTxEnd(receipt *types.Receipt, err error) {
for _, t := range t.tracers {
if t.OnTxEnd != nil {
t.OnTxEnd(receipt, err)
@ -133,7 +138,7 @@ func (t *muxTracer) OnTxEnd(receipt *types.Receipt, err error) {
}
}
func (t *muxTracer) OnBalanceChange(a common.Address, prev, new *big.Int, reason tracing.BalanceChangeReason) {
func (t *MuxTracer) OnBalanceChange(a common.Address, prev, new *big.Int, reason tracing.BalanceChangeReason) {
for _, t := range t.tracers {
if t.OnBalanceChange != nil {
t.OnBalanceChange(a, prev, new, reason)
@ -141,7 +146,7 @@ func (t *muxTracer) OnBalanceChange(a common.Address, prev, new *big.Int, reason
}
}
func (t *muxTracer) OnNonceChange(a common.Address, prev, new uint64) {
func (t *MuxTracer) OnNonceChange(a common.Address, prev, new uint64) {
for _, t := range t.tracers {
if t.OnNonceChange != nil {
t.OnNonceChange(a, prev, new)
@ -149,7 +154,7 @@ func (t *muxTracer) OnNonceChange(a common.Address, prev, new uint64) {
}
}
func (t *muxTracer) OnCodeChange(a common.Address, prevCodeHash common.Hash, prev []byte, codeHash common.Hash, code []byte) {
func (t *MuxTracer) OnCodeChange(a common.Address, prevCodeHash common.Hash, prev []byte, codeHash common.Hash, code []byte) {
for _, t := range t.tracers {
if t.OnCodeChange != nil {
t.OnCodeChange(a, prevCodeHash, prev, codeHash, code)
@ -157,7 +162,7 @@ func (t *muxTracer) OnCodeChange(a common.Address, prevCodeHash common.Hash, pre
}
}
func (t *muxTracer) OnStorageChange(a common.Address, k, prev, new common.Hash) {
func (t *MuxTracer) OnStorageChange(a common.Address, k, prev, new common.Hash) {
for _, t := range t.tracers {
if t.OnStorageChange != nil {
t.OnStorageChange(a, k, prev, new)
@ -165,7 +170,7 @@ func (t *muxTracer) OnStorageChange(a common.Address, k, prev, new common.Hash)
}
}
func (t *muxTracer) OnLog(log *types.Log) {
func (t *MuxTracer) OnLog(log *types.Log) {
for _, t := range t.tracers {
if t.OnLog != nil {
t.OnLog(log)
@ -174,14 +179,14 @@ func (t *muxTracer) OnLog(log *types.Log) {
}
// GetResult returns an empty json object.
func (t *muxTracer) GetResult() (json.RawMessage, error) {
func (t *MuxTracer) GetResult() (json.RawMessage, error) {
resObject := make(map[string]json.RawMessage)
for i, tt := range t.tracers {
for n, tt := range t.tracers {
r, err := tt.GetResult()
if err != nil {
return nil, err
}
resObject[t.names[i]] = r
resObject[n] = r
}
res, err := json.Marshal(resObject)
if err != nil {
@ -191,8 +196,12 @@ func (t *muxTracer) GetResult() (json.RawMessage, error) {
}
// Stop terminates execution of the tracer at the first opportune moment.
func (t *muxTracer) Stop(err error) {
func (t *MuxTracer) Stop(err error) {
for _, t := range t.tracers {
t.Stop(err)
}
}
func (t *MuxTracer) Tracers() map[string]*tracers.Tracer {
return t.tracers
}