diff --git a/eth/tracers/live/filter.go b/eth/tracers/live/filter.go new file mode 100644 index 0000000000..3cffb45eaf --- /dev/null +++ b/eth/tracers/live/filter.go @@ -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) + } +} diff --git a/eth/tracers/live/filter_api.go b/eth/tracers/live/filter_api.go new file mode 100644 index 0000000000..6c2bc58eae --- /dev/null +++ b/eth/tracers/live/filter_api.go @@ -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) +} diff --git a/eth/tracers/native/call.go b/eth/tracers/native/call.go index 1b94dd7b67..13840d42e1 100644 --- a/eth/tracers/native/call.go +++ b/eth/tracers/native/call.go @@ -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() } diff --git a/eth/tracers/native/mux.go b/eth/tracers/native/mux.go index c3b1d9f8ca..2d62e02677 100644 --- a/eth/tracers/native/mux.go +++ b/eth/tracers/native/mux.go @@ -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 - } + t, err := NewMuxTracer(cfg) + if 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) - 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 +}