From 7520990fa6fdc3c96b5a2cb74fcde7c1920f6953 Mon Sep 17 00:00:00 2001 From: David Zhou Date: Fri, 9 May 2025 16:20:57 -0400 Subject: [PATCH] DEV: Implemented onClose() --- eth/tracers/firehose.go | 26 +++++++++++++----------- eth/tracers/firehose_concurrency_test.go | 7 ++++--- 2 files changed, 18 insertions(+), 15 deletions(-) diff --git a/eth/tracers/firehose.go b/eth/tracers/firehose.go index 655f5852ee..6ea929bc47 100644 --- a/eth/tracers/firehose.go +++ b/eth/tracers/firehose.go @@ -78,6 +78,7 @@ func NewTracingHooksFromFirehose(tracer *Firehose) *tracing.Hooks { OnBlockStart: tracer.OnBlockStart, OnBlockEnd: tracer.OnBlockEnd, OnSkippedBlock: tracer.OnSkippedBlock, + OnClose: tracer.OnClose, // TODO: OnClose OnTxStart: tracer.OnTxStart, @@ -498,13 +499,13 @@ func (f *Firehose) OnBlockEnd(err error) { f.fixOrdinalsForEndOfBlockChanges() } - // f.printBlockToFirehose(f.block, f.blockFinality) - f.ensureInBlockAndNotInTrx() - job := &blockPrintJob{ - block: f.block, - finality: f.blockFinality, - } - f.blockPrintQueue <- job + f.printBlockToFirehose(f.block, f.blockFinality) + //f.ensureInBlockAndNotInTrx() + //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 @@ -623,6 +624,12 @@ func (f *Firehose) reorderCallOrdinals(call *pbeth.Call, ordinalBase uint64) (or return call.EndOrdinal } +func (f *Firehose) OnClose() { + log.Info("Firehose closing: waiting for worker goroutines to finish and shutting down channels") + close(f.blockPrintQueue) + f.workerWg.Wait() +} + func (f *Firehose) OnSystemCallStart() { firehoseInfo("system call start") f.ensureInBlockAndNotInTrx() @@ -2826,8 +2833,3 @@ func (f *Firehose) blockPrintWorker() { 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 index d1619da7f8..d69291e858 100644 --- a/eth/tracers/firehose_concurrency_test.go +++ b/eth/tracers/firehose_concurrency_test.go @@ -37,7 +37,7 @@ func TestFirehose_BlockPrintsToFirehose_SingleBlock(t *testing.T) { f.OnBlockEnd(nil) } - f.Shutdown() + f.OnClose() output := f.InternalTestingBuffer().String() @@ -50,11 +50,12 @@ func TestFirehose_BlockPrintsToFirehose_SingleBlock(t *testing.T) { fields := strings.SplitN(line, " ", 4) if len(fields) >= 3 { + require.Equal(t, "FIRE", fields[0]) + require.Equal(t, "BLOCK", fields[1]) outNumber = append(outNumber, fields[2]) } } - require.Contains(t, output, "FIRE BLOCK", "expected FIRE BLOCK output not found") require.Equal(t, []string{"123", "124", "125"}, outNumber) } @@ -89,7 +90,7 @@ func TestFirehose_BlocksPrintToFirehose_MultipleBlocksInOrder(t *testing.T) { f.OnBlockEnd(nil) } - f.Shutdown() + f.OnClose() output := f.InternalTestingBuffer().String() extractedBlocks := extractBlocksFromOutput(t, output)