From 526ca333b6b2b07228ec8ae4396741b91849baaa Mon Sep 17 00:00:00 2001 From: David Zhou Date: Fri, 9 May 2025 09:34:43 -0400 Subject: [PATCH] DEV: Implementation of concurrency in firehose --- eth/tracers/firehose.go | 39 +++++++++++++++++++++++- eth/tracers/firehose_concurrency_test.go | 9 ++++++ 2 files changed, 47 insertions(+), 1 deletion(-) create mode 100644 eth/tracers/firehose_concurrency_test.go diff --git a/eth/tracers/firehose.go b/eth/tracers/firehose.go index 9d30c30f21..7ca1b01e12 100644 --- a/eth/tracers/firehose.go +++ b/eth/tracers/firehose.go @@ -186,6 +186,10 @@ type Firehose struct { // Testing state, only used in tests and private configs testingBuffer *bytes.Buffer testingIgnoreGenesisBlock bool + + // Worker queuefor blockprinting + blockPrintQueue chan *blockPrintJob + workerWg sync.WaitGroup } const FirehoseProtocolVersion = "3.0" @@ -238,6 +242,8 @@ func NewFirehose(config *FirehoseConfig) *Firehose { callStack: NewCallStack(), deferredCallState: NewDeferredCallState(), latestCallEnterSuicided: false, + + blockPrintQueue: make(chan *blockPrintJob, 100), } if config.private != nil { @@ -247,6 +253,13 @@ func NewFirehose(config *FirehoseConfig) *Firehose { } } + // Worker goroutines + numWorkers := 1 + firehose.workerWg.Add(numWorkers) + for i := 0; i < numWorkers; i++ { + go firehose.blockPrintWorker() + } + return firehose } @@ -483,7 +496,13 @@ func (f *Firehose) OnBlockEnd(err error) { } f.ensureInBlockAndNotInTrx() - f.printBlockToFirehose(f.block, f.blockFinality) + + job := &blockPrintJob{ + block: f.block, + finality: f.blockFinality, + } + f.blockPrintQueue <- job + } else { // An error occurred, could have happen in transaction/call context, we must not check if in trx/call, only check in block f.ensureInBlock(0) @@ -2791,3 +2810,21 @@ func (m Memory) GetPtr(offset, size int64) []byte { reminder := m[offset:] return append(reminder, make([]byte, int(size)-len(reminder))...) } + +type blockPrintJob struct { + block *pbeth.Block + finality *FinalityStatus +} + +func (f *Firehose) blockPrintWorker() { + defer f.workerWg.Done() + + for job := range f.blockPrintQueue { + f.printBlockToFirehose(job.block, job.finality) + } +} + +func (f *Firehose) Shutdown() { + close(f.blockPrintQueue) + f.workerWg.Wait() +} diff --git a/eth/tracers/firehose_concurrency_test.go b/eth/tracers/firehose_concurrency_test.go new file mode 100644 index 0000000000..ace16035f4 --- /dev/null +++ b/eth/tracers/firehose_concurrency_test.go @@ -0,0 +1,9 @@ +package tracers + +import ( + "testing" +) + +func TestFirehoseBlockPrintingOrder(t *testing.T) { + +}