From d0f6dfdf6014f4ef8db0021b6b67bd601a7ab4bc Mon Sep 17 00:00:00 2001 From: David Zhou Date: Fri, 23 May 2025 13:30:22 -0400 Subject: [PATCH] PR fix --- eth/tracers/firehose.go | 42 +++++++++---------- eth/tracers/firehose_concurrency.go | 24 +++++------ .../tracetest/firehose/firehose_test.go | 16 ++----- 3 files changed, 35 insertions(+), 47 deletions(-) diff --git a/eth/tracers/firehose.go b/eth/tracers/firehose.go index da2e96eb30..293ba4d08a 100644 --- a/eth/tracers/firehose.go +++ b/eth/tracers/firehose.go @@ -170,13 +170,15 @@ func (c *FirehoseConfig) ForcedBackwardCompatibility() bool { type Firehose struct { // Global state - outputBuffer *bytes.Buffer - initSent *atomic.Bool - config *FirehoseConfig - chainConfig *params.ChainConfig - hasher crypto.KeccakState // Keccak256 hasher instance shared across tracer needs (non-concurrent safe) - hasherBuf common.Hash // Keccak256 hasher result array shared across tracer needs (non-concurrent safe) - tracerID string + outputBuffer *bytes.Buffer + initSent *atomic.Bool + config *FirehoseConfig + chainConfig *params.ChainConfig + hasher crypto.KeccakState // Keccak256 hasher instance shared across tracer needs (non-concurrent safe) + hasherBuf common.Hash // Keccak256 hasher result array shared across tracer needs (non-concurrent safe) + tracerID string + concurrentFlushQueue *ConcurrentFlushQueue + concurrentFlushBufferSize int // The FirehoseTracer is used in multiple chains, some for which were produced using a legacy version // of the whole tracing infrastructure. This legacy version had many small bugs here and there that // we must "reproduce" on some chain to ensure that the FirehoseTracer produces the same output @@ -199,8 +201,6 @@ type Firehose struct { blockReorderOrdinalSnapshot uint64 blockReorderOrdinalOnce sync.Once blockIsGenesis bool - BlockFlushQueue *BlockFlushQueue - blockFlushBufferSize int // Transaction state evm *tracing.VMContext @@ -260,10 +260,10 @@ func NewFirehose(config *FirehoseConfig) *Firehose { concurrentBlockFlushing: config.ConcurrentBlockFlushing, // Block state - blockOrdinal: &Ordinal{}, - blockFinality: &FinalityStatus{}, - blockReorderOrdinal: false, - blockFlushBufferSize: 100, + blockOrdinal: &Ordinal{}, + blockFinality: &FinalityStatus{}, + blockReorderOrdinal: false, + concurrentFlushBufferSize: 100, // Transaction state transactionLogIndex: 0, @@ -375,13 +375,13 @@ func (f *Firehose) OnBlockchainInit(chainConfig *params.ChainConfig) { } if f.config.ConcurrentBlockFlushing > 0 { - log.Info("Firehose concurrent block flushing enabled, starting goroutines") - f.BlockFlushQueue = NewBlockFlushQueue( - f.config.ConcurrentBlockFlushing, - f.blockFlushBufferSize, + log.Info("Firehose concurrent block flushing enabled, starting goroutine") + f.concurrentFlushQueue = NewConcurrentFlushQueue( + f.concurrentFlushBufferSize, f.printBlockToFirehose, f.flushToFirehose, ) + f.concurrentFlushQueue.Start(f.config.ConcurrentBlockFlushing) } log.Info("Firehose tracer initialized", @@ -530,7 +530,7 @@ func (f *Firehose) OnBlockEnd(err error) { // Flush block to firehose and optionally use goroutine if f.concurrentBlockFlushing > 0 { - f.BlockFlushQueue.Enqueue(f.block, f.blockFinality) + f.concurrentFlushQueue.Enqueue(f.block, f.blockFinality) } else { f.printBlockToFirehose(f.block, f.blockFinality) } @@ -653,9 +653,9 @@ func (f *Firehose) reorderCallOrdinals(call *pbeth.Call, ordinalBase uint64) (or } func (f *Firehose) OnClose() { - if f.BlockFlushQueue != nil { + if f.concurrentFlushQueue != nil { log.Info("Firehose closing, flushing queued blocks to standard output") - f.BlockFlushQueue.Close() + f.concurrentFlushQueue.CloseChannels() } } @@ -1949,7 +1949,7 @@ func (f *Firehose) printBlockToFirehose(block *pbeth.Block, finalityStatus *Fina buf.WriteString("\n") if f.concurrentBlockFlushing > 0 { - f.BlockFlushQueue.outputQueue <- &outputJob{ + f.concurrentFlushQueue.outputQueue <- &outputJob{ blockNum: block.Number, data: buf.Bytes(), } diff --git a/eth/tracers/firehose_concurrency.go b/eth/tracers/firehose_concurrency.go index 39d5183462..6059df7bcd 100644 --- a/eth/tracers/firehose_concurrency.go +++ b/eth/tracers/firehose_concurrency.go @@ -15,7 +15,7 @@ type outputJob struct { data []byte } -type BlockFlushQueue struct { +type ConcurrentFlushQueue struct { bufferSize int startSignal chan uint64 @@ -29,12 +29,8 @@ type BlockFlushQueue struct { closeOnce sync.Once } -func NewBlockFlushQueue(concurrency int, bufferSize int, printBlockFunc func(*pbeth.Block, *FinalityStatus), outputFunc func([]byte)) *BlockFlushQueue { - if concurrency <= 0 { - panic("BlockFlushQueue requires concurrency > 0") - } - - q := &BlockFlushQueue{ +func NewConcurrentFlushQueue(bufferSize int, printBlockFunc func(*pbeth.Block, *FinalityStatus), outputFunc func([]byte)) *ConcurrentFlushQueue { + return &ConcurrentFlushQueue{ startSignal: make(chan uint64, 1), jobQueue: make(chan *blockPrintJob, bufferSize), outputQueue: make(chan *outputJob, bufferSize), @@ -42,7 +38,9 @@ func NewBlockFlushQueue(concurrency int, bufferSize int, printBlockFunc func(*pb bufferSize: bufferSize, printBlockFunc: printBlockFunc, } +} +func (q *ConcurrentFlushQueue) Start(concurrency int) { for i := 0; i < concurrency; i++ { q.jobWG.Add(1) go q.worker() @@ -50,11 +48,9 @@ func NewBlockFlushQueue(concurrency int, bufferSize int, printBlockFunc func(*pb q.outputWG.Add(1) go q.outputOrderer() - - return q } -func (q *BlockFlushQueue) Enqueue(block *pbeth.Block, finality *FinalityStatus) { +func (q *ConcurrentFlushQueue) Enqueue(block *pbeth.Block, finality *FinalityStatus) { select { case q.startSignal <- block.Number: default: @@ -66,10 +62,10 @@ func (q *BlockFlushQueue) Enqueue(block *pbeth.Block, finality *FinalityStatus) } } -// Close CloseBlockPrintQueue signals block printing goroutines to shut down and waits for them. +// CloseChannels signals goroutines to shut down and waits for them. // It blocks until all concurrent block flushing operations are completed, ensuring a clean // shutdown of the printing pipeline. -func (q *BlockFlushQueue) Close() { +func (q *ConcurrentFlushQueue) CloseChannels() { q.closeOnce.Do(func() { close(q.jobQueue) q.jobWG.Wait() @@ -79,7 +75,7 @@ func (q *BlockFlushQueue) Close() { } // Instantiates a worker that listens for jobs -func (q *BlockFlushQueue) worker() { +func (q *ConcurrentFlushQueue) worker() { defer q.jobWG.Done() for job := range q.jobQueue { q.printBlockFunc(job.block, job.finality) @@ -87,7 +83,7 @@ func (q *BlockFlushQueue) worker() { } // Channel ensuring that blocks are linearly flushed out in order -func (q *BlockFlushQueue) outputOrderer() { +func (q *ConcurrentFlushQueue) outputOrderer() { defer q.outputWG.Done() buffer := make(map[uint64][]byte) next := <-q.startSignal diff --git a/eth/tracers/internal/tracetest/firehose/firehose_test.go b/eth/tracers/internal/tracetest/firehose/firehose_test.go index fc8f65c26e..72ee89a586 100644 --- a/eth/tracers/internal/tracetest/firehose/firehose_test.go +++ b/eth/tracers/internal/tracetest/firehose/firehose_test.go @@ -55,15 +55,11 @@ func TestFirehosePrestate(t *testing.T) { ConcurrentBlockFlushing: concurrent, } - tracer, tracingHooks, onClose := newFirehoseTestTracer(t, model, config) - defer onClose() + tracer, tracingHooks, _ := newFirehoseTestTracer(t, model, config) runPrestateBlock(t, filepath.Join(folder, "prestate.json"), tracingHooks) - if tracer.BlockFlushQueue != nil { - tracer.BlockFlushQueue.Close() - } - + tracer.OnClose() genesisLine, blockLines, unknownLines := readTracerFirehoseLines(t, tracer) require.Len(t, unknownLines, 0, "Lines:\n%s", strings.Join( slicesMap(unknownLines, func(l unknownLine) string { return "- '" + string(l) + "'" }), "\n")) @@ -214,8 +210,7 @@ func testBlockTracesCorrectly(t *testing.T, genesisSpec *core.Genesis, engine co ConcurrentBlockFlushing: concurrent, } - tracer, tracingHooks, onClose := newFirehoseTestTracer(t, model, config) - defer onClose() + tracer, tracingHooks, _ := newFirehoseTestTracer(t, model, config) chain, err := core.NewBlockChain(rawdb.NewMemoryDatabase(), nil, genesisSpec, nil, engine, vm.Config{Tracer: tracingHooks}, nil) require.NoError(t, err, "failed to create tester chain") @@ -229,10 +224,7 @@ func testBlockTracesCorrectly(t *testing.T, genesisSpec *core.Genesis, engine co n, err := chain.InsertChain(blocks) require.NoError(t, err, "failed to insert chain block %d", n) - if tracer.BlockFlushQueue != nil { - tracer.BlockFlushQueue.Close() - } - + tracer.OnClose() assertBlockEquals(t, tracer, filepath.Join("testdata", goldenDir, string(model)), len(blocks)) }) }