From a053ac76fd6c463256d413c714e96179454e4b11 Mon Sep 17 00:00:00 2001 From: David Zhou Date: Tue, 13 May 2025 13:19:07 -0400 Subject: [PATCH] BUG: Fixed order --- eth/tracers/firehose.go | 19 +++++++++++++++++++ eth/tracers/firehose_concurrency_test.go | 2 +- 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/eth/tracers/firehose.go b/eth/tracers/firehose.go index b3a80fa658..d09392fb04 100644 --- a/eth/tracers/firehose.go +++ b/eth/tracers/firehose.go @@ -272,7 +272,26 @@ func NewFirehose(config *FirehoseConfig) *Firehose { // Output channel to order the flushing linearly firehose.flushOutputDone.Add(1) + go func() { + defer firehose.flushOutputDone.Done() + expected := uint64(0) + buffer := make(map[uint64][]byte) + + for result := range firehose.blockOutputQueue { + buffer[result.blockNum] = result.data + + for { + data, ok := buffer[expected] + if !ok { + break + } + firehose.flushToFirehose(data) + delete(buffer, expected) + expected++ + } + } + }() } return firehose diff --git a/eth/tracers/firehose_concurrency_test.go b/eth/tracers/firehose_concurrency_test.go index 795b2d7212..f7d2dd369a 100644 --- a/eth/tracers/firehose_concurrency_test.go +++ b/eth/tracers/firehose_concurrency_test.go @@ -54,7 +54,7 @@ func TestFirehose_BlockPrintsToFirehose_SingleBlock(t *testing.T) { func TestFirehose_BlocksPrintToFirehose_MultipleBlocksInOrder(t *testing.T) { - const blockCount = 10 + const blockCount = 100 const baseBlockNum = 0 f := NewFirehose(&FirehoseConfig{