Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 72 additions & 5 deletions rpc/jsonrpc/call_traces_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@
package jsonrpc

import (
"bytes"
"context"
"encoding/json"
"sync"
"testing"

Expand Down Expand Up @@ -383,9 +385,9 @@ func TestFilterRejectedBlockOverrideReturnsError(t *testing.T) {
// TestFilterSignerReflectsBlockOverridesNumber is filterV3's analogue of
// TestReplayTransactionSignerReflectsBlockOverridesNumber: filterV3 derives
// fork rules (lastRules) from the overridden BlockContext but must also
// recompute lastSigner from it, not from the block's real number. filterV3
// reports per-transaction failures as an "error" field inside the stream
// rather than as a Go error, so the assertion inspects the stream contents.
// recompute lastSigner from it, not from the block's real number. A
// transaction that cannot be traced must fail the whole request instead of
// mixing an error object into the result array.
func TestFilterSignerReflectsBlockOverridesNumber(t *testing.T) {
if testing.Short() {
t.Skip("slow test")
Expand All @@ -408,6 +410,71 @@ func TestFilterSignerReflectsBlockOverridesNumber(t *testing.T) {
err := api.Filter(context.Background(), traceReq, new(bool), &config.TraceConfig{
BlockOverrides: &ethapi.BlockOverrides{Number: (*hexutil.U256)(uint256.NewInt(1))},
}, stream)
require.NoError(t, err)
require.Contains(t, string(stream.Buffer()), "protected txn is not supported by signer")
require.ErrorContains(t, err, "protected txn is not supported by signer")
require.Empty(t, string(stream.Buffer()))
}

// TestFilterErrorAfterExportedTracesKeepsValidJSON covers the other half of
// filterV3's error contract: when a transaction fails after earlier traces were
// already streamed, the request still fails, and the envelope stays valid JSON
// whose result array holds only TraceEntry items. Filtering blocks 1-3 with no
// address filter exports both empty blocks' reward traces before the protected
// transaction in block 3 is rejected by the overridden pre-EIP-155 signer.
// The envelope is assembled the way runMethod does it, since sealing the
// half-written array is the handler's job, not filterV3's.
func TestFilterErrorAfterExportedTracesKeepsValidJSON(t *testing.T) {
if testing.Short() {
t.Skip("slow test")
}

c := newBaseFeeTestChain(t, delayedSpuriousDragonConfig())
c.mineProtectedTxAtBlock3(t)
api := c.traceAPI()

from, to := rpc.BlockNumber(1), rpc.BlockNumber(3)
traceReq := TraceFilterRequest{
FromBlock: &rpc.BlockNumberOrHash{BlockNumber: &from},
ToBlock: &rpc.BlockNumberOrHash{BlockNumber: &to},
}

var buf bytes.Buffer
stream := jsonstream.New(&buf)
stream.WriteObjectStart()
stream.WriteObjectField("jsonrpc")
stream.WriteString("2.0")
stream.WriteMore()
stream.WriteObjectField("id")
stream.WriteInt(1)
stream.WriteMore()
result := jsonstream.NewLazyFieldStream(stream, "result", false)

err := api.Filter(context.Background(), traceReq, new(bool), &config.TraceConfig{
BlockOverrides: &ethapi.BlockOverrides{Number: (*hexutil.U256)(uint256.NewInt(1))},
}, result)
require.ErrorContains(t, err, "protected txn is not supported by signer")
require.True(t, result.Written(), "test needs traces exported before the failure")

result.CloseIfOpen()
stream.WriteMore()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
require.NoError(t, stream.Flush())

var envelope map[string]json.RawMessage
require.NoError(t, json.Unmarshal(buf.Bytes(), &envelope), "envelope is not valid JSON: %s", buf.String())
require.Contains(t, envelope, "error")

var traces []json.RawMessage
require.NoError(t, json.Unmarshal(envelope["result"], &traces), "result array was left unsealed: %s", envelope["result"])
require.NotEmpty(t, traces)
for i, trace := range traces {
var entry map[string]json.RawMessage
require.NoError(t, json.Unmarshal(trace, &entry), "item %d is not a JSON object", i)
require.Contains(t, entry, "type", "item %d is not a TraceEntry", i)
if reason, ok := entry["error"]; ok {
var failure string
require.NoError(t, json.Unmarshal(reason, &failure),
"item %d: error must be a TraceEntry failure reason, not an RPC error object", i)
}
}
}
156 changes: 31 additions & 125 deletions rpc/jsonrpc/trace_filtering.go
Original file line number Diff line number Diff line change
Expand Up @@ -412,8 +412,17 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
engine := api.engine()

var json = jsoniter.ConfigCompatibleWithStandardLibrary
stream.WriteArrayStart()
// The result array opens on the first exported trace so that an error
// before it yields an error-only response instead of "result" plus "error".
first := true
beginItem := func() {
if first {
stream.WriteArrayStart()
first = false
} else {
stream.WriteMore()
}
}
// Execute all transactions in picked blocks

count := uint64(^uint(0)) // this just makes it easier to use below
Expand All @@ -438,20 +447,7 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
noop := state.NewNoopWriter()
isPos := false

// traceFilterTxn traces one transaction. A nil result means the error was
// already written to the stream; the transaction state is closed either way.
traceFilterTxn := func(blockNum, txNum uint64, txIndex int, txn types.Transaction, msg *types.Message) (*TraceCallResult, error) {
writeErr := func(err error) {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
}

stateCache := shards.NewStateCache(32, 0 /* no limit */) // this cache living only during current RPC call, but required to store state writes
cachedReader := state.NewCachedReader(state.NewHistoryReaderV3(dbtx, txNum), stateCache)
cachedWriter := state.NewCachedWriter(noop, stateCache)
Expand Down Expand Up @@ -501,30 +497,23 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
if ot.Tracer() != nil && ot.Tracer().Hooks.OnTxEnd != nil {
ot.Tracer().OnTxEnd(nil, timeoutErr)
}
// Safe to skip FinalizeTx/CommitBlock: each iteration creates a fresh
// stateCache, cachedReader and ibs from the next txNum, and writes go
// to a noop writer, so no partial state escapes this scope.
writeErr(timeoutErr)
return nil, nil
return nil, timeoutErr
}
if execErr != nil {
if ot.Tracer() != nil && ot.Tracer().Hooks.OnTxEnd != nil {
ot.Tracer().OnTxEnd(nil, execErr)
}
writeErr(execErr)
return nil, nil
return nil, execErr
}
if ot.Tracer() != nil && ot.Tracer().Hooks.OnTxEnd != nil {
ot.Tracer().OnTxEnd(&types.Receipt{GasUsed: execResult.ReceiptGasUsed}, nil)
}
traceResult.Output = bytes.Clone(execResult.ReturnData)
if err := ibs.FinalizeTx(evm.ChainRules(), noop); err != nil {
writeErr(err)
return nil, nil
return nil, err
}
if err := ibs.CommitBlock(evm.ChainRules(), cachedWriter); err != nil {
writeErr(err)
return nil, nil
return nil, err
}
return traceResult, nil
}
Expand All @@ -535,39 +524,15 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
}
txNum, blockNum, txIndex, isFnalTxn, blockNumChanged, err := it.Next()
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}

if blockNumChanged {
if lastHeader, err = api._blockReader.HeaderByNumber(ctx, dbtx, blockNum); err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}
if lastHeader == nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(fmt.Errorf("header not found: %d", blockNum), stream)
stream.WriteObjectEnd()
continue
return fmt.Errorf("header not found: %d", blockNum)
}

if !isPos && chainConfig.TerminalTotalDifficulty != nil {
Expand Down Expand Up @@ -595,15 +560,7 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB

body, _, err := api._blockReader.Body(ctx, dbtx, lastBlockHash, blockNum)
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}
// Block reward section, handle specially
minerReward, uncleRewards := ethash.AccumulateRewards(chainConfig, lastHeader, body.Uncles)
Expand All @@ -612,22 +569,10 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
tr := newRewardTrace(lastBlockHash, blockNum, lastHeader.Coinbase, rewardTypeBlock, minerReward.ToBig())
b, err := json.Marshal(tr)
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}
if nSeen > after && nExported < count {
if first {
first = false
} else {
stream.WriteMore()
}
beginItem()
if _, err := stream.Write(b); err != nil {
return err
}
Expand All @@ -641,22 +586,10 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
tr := newRewardTrace(lastBlockHash, blockNum, uncle.Coinbase, rewardTypeUncle, uncleRewards[i].ToBig())
b, err := json.Marshal(tr)
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}
if nSeen > after && nExported < count {
if first {
first = false
} else {
stream.WriteMore()
}
beginItem()
if _, err := stream.Write(b); err != nil {
return err
}
Expand All @@ -674,40 +607,21 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
//fmt.Printf("txNum=%d, blockNum=%d, txIndex=%d\n", txNum, blockNum, txIndex)
txn, ok, err := api._txnReader.TxnByIdxInBlock(ctx, dbtx, blockNum, txIndex)
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}
if !ok {
continue //guess block doesn't have transactions
}
txHash := txn.Hash()
msg, err := txn.AsMessage(*lastSigner, &lastBaseFee, lastRules)
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}

traceResult, err := traceFilterTxn(blockNum, txNum, txIndex, txn, msg)
if err != nil {
return err
}
if traceResult == nil {
continue
}
isIntersectionMode := req.Mode == TraceFilterModeIntersection
for _, pt := range traceResult.Trace {
if includeAll || filterTrace(pt, fromAddresses, toAddresses, isIntersectionMode) {
Expand All @@ -718,22 +632,10 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
pt.TransactionPosition = &txIndexU64
b, err := json.Marshal(pt)
if err != nil {
if first {
first = false
} else {
stream.WriteMore()
}
stream.WriteObjectStart()
rpc.HandleError(err, stream)
stream.WriteObjectEnd()
continue
return err
}
if nSeen > after && nExported < count {
if first {
first = false
} else {
stream.WriteMore()
}
beginItem()
if _, err := stream.Write(b); err != nil {
return err
}
Expand All @@ -742,7 +644,11 @@ func (api *TraceAPIImpl) filterV3(ctx context.Context, dbtx kv.TemporalTx, fromB
}
}
}
stream.WriteArrayEnd()
if first {
stream.WriteEmptyArray()
} else {
stream.WriteArrayEnd()
}
return nil
}

Expand Down
Loading