From 3d40ec0feb56bd017516b5013a1ca404a3c5784e Mon Sep 17 00:00:00 2001 From: David Zhou Date: Mon, 12 May 2025 11:11:41 -0400 Subject: [PATCH] TEST: Testing concurrent and not concurrent --- eth/tracers/firehose.go | 4 +- .../tracetest/firehose/firehose_test.go | 64 ++++++++++--------- .../tracetest/firehose/helpers_test.go | 6 +- 3 files changed, 39 insertions(+), 35 deletions(-) diff --git a/eth/tracers/firehose.go b/eth/tracers/firehose.go index 6905da0c9f..19aeea5a19 100644 --- a/eth/tracers/firehose.go +++ b/eth/tracers/firehose.go @@ -632,7 +632,7 @@ func (f *Firehose) reorderCallOrdinals(call *pbeth.Call, ordinalBase uint64) (or func (f *Firehose) OnClose() { log.Info("Firehose closing...") if f.concurrentBlockFlushing { - log.Info("waiting for worker goroutines to finish and shutting down channels") + log.Info("Closing channel: flushing the remaining blocks to firehose") f.CloseBlockPrintQueue() } } @@ -1860,6 +1860,7 @@ func (f *Firehose) printBlockToFirehose(block *pbeth.Block, finalityStatus *Fina panic(fmt.Errorf("failed to marshal block: %w", err)) } + // TODO: If multiple goroutine, this becomes shared resource f.outputBuffer.Reset() previousHash := block.PreviousID() @@ -2843,7 +2844,6 @@ func (f *Firehose) blockPrintWorker() { func (f *Firehose) CloseBlockPrintQueue() { if f.concurrentBlockFlushing { f.closeOnce.Do(func() { - log.Info("Closing channel: flushing the remaining blocks to firehose") close(f.blockPrintQueue) f.flushDone.Wait() }) diff --git a/eth/tracers/internal/tracetest/firehose/firehose_test.go b/eth/tracers/internal/tracetest/firehose/firehose_test.go index 36600317b9..e004105a06 100644 --- a/eth/tracers/internal/tracetest/firehose/firehose_test.go +++ b/eth/tracers/internal/tracetest/firehose/firehose_test.go @@ -38,26 +38,28 @@ func TestFirehosePrestate(t *testing.T) { "./testdata/TestFirehosePrestate/extra_account_creations", } - for _, folder := range testFolders { - name := filepath.Base(folder) + // TODO: have goroutine config on AND off + for _, concurrent := range []bool{true, false} { + for _, folder := range testFolders { + name := filepath.Base(folder) - for _, model := range tracingModels { - t.Run(string(model)+"/"+name, func(t *testing.T) { - tracer, tracingHooks, onClose := newFirehoseTestTracer(t, model) - defer onClose() + for _, model := range tracingModels { + t.Run(string(model)+"/"+name, func(t *testing.T) { + tracer, tracingHooks, onClose := newFirehoseTestTracer(t, model, concurrent) + defer onClose() - runPrestateBlock(t, filepath.Join(folder, "prestate.json"), tracingHooks) + runPrestateBlock(t, filepath.Join(folder, "prestate.json"), tracingHooks) - tracer.CloseBlockPrintQueue() + tracer.CloseBlockPrintQueue() - 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")) - require.NotNil(t, genesisLine) - blockLines.assertOnlyBlockEquals(t, filepath.Join(folder, string(model)), 1) - }) + 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")) + require.NotNil(t, genesisLine) + blockLines.assertOnlyBlockEquals(t, filepath.Join(folder, string(model)), 1) + }) + } } } - } func TestFirehose_EIP7702(t *testing.T) { // Copied from ./core/blockchain_test.go#L4180 (TestEIP7702) @@ -187,26 +189,28 @@ func TestFirehose_SystemCalls(t *testing.T) { func testBlockTracesCorrectly(t *testing.T, genesisSpec *core.Genesis, engine consensus.Engine, blocks []*types.Block, goldenDir string) { t.Helper() - for _, model := range tracingModels { - t.Run(string(model), func(t *testing.T) { - tracer, tracingHooks, onClose := newFirehoseTestTracer(t, model) - defer onClose() + for _, concurrent := range []bool{true, false} { + for _, model := range tracingModels { + t.Run(string(model), func(t *testing.T) { + tracer, tracingHooks, onClose := newFirehoseTestTracer(t, model, concurrent) + defer onClose() - chain, err := core.NewBlockChain(rawdb.NewMemoryDatabase(), nil, genesisSpec, nil, engine, vm.Config{Tracer: tracingHooks}, nil) - require.NoError(t, err, "failed to create tester chain") + chain, err := core.NewBlockChain(rawdb.NewMemoryDatabase(), nil, genesisSpec, nil, engine, vm.Config{Tracer: tracingHooks}, nil) + require.NoError(t, err, "failed to create tester chain") - chain.SetBlockValidatorAndProcessorForTesting( - ignoreValidateStateValidator{core.NewBlockValidator(genesisSpec.Config, chain)}, - core.NewStateProcessor(genesisSpec.Config, chain.HeaderChain()), - ) + chain.SetBlockValidatorAndProcessorForTesting( + ignoreValidateStateValidator{core.NewBlockValidator(genesisSpec.Config, chain)}, + core.NewStateProcessor(genesisSpec.Config, chain.HeaderChain()), + ) - defer chain.Stop() - n, err := chain.InsertChain(blocks) - require.NoError(t, err, "failed to insert chain block %d", n) + defer chain.Stop() + n, err := chain.InsertChain(blocks) + require.NoError(t, err, "failed to insert chain block %d", n) - tracer.CloseBlockPrintQueue() + tracer.CloseBlockPrintQueue() - assertBlockEquals(t, tracer, filepath.Join("testdata", goldenDir, string(model)), len(blocks)) - }) + assertBlockEquals(t, tracer, filepath.Join("testdata", goldenDir, string(model)), len(blocks)) + }) + } } } diff --git a/eth/tracers/internal/tracetest/firehose/helpers_test.go b/eth/tracers/internal/tracetest/firehose/helpers_test.go index b7263dc9ae..0541d37748 100644 --- a/eth/tracers/internal/tracetest/firehose/helpers_test.go +++ b/eth/tracers/internal/tracetest/firehose/helpers_test.go @@ -29,17 +29,17 @@ type firehoseInitLine struct { type firehoseBlockLines []firehoseBlockLine -func newFirehoseTestTracer(t *testing.T, model tracingModel) (*tracers.Firehose, *tracing.Hooks, func()) { +func newFirehoseTestTracer(t *testing.T, model tracingModel, concurrentBlockFlushing bool) (*tracers.Firehose, *tracing.Hooks, func()) { t.Helper() tracer, err := tracers.NewFirehoseFromRawJSON([]byte(fmt.Sprintf(`{ - "concurrentBlockFlushing": true, + "concurrentBlockFlushing": %t, "_private": { "flushToTestBuffer": true, "ignoreGenesisBlock": true, "forcedBackwardCompatibility": %t } - }`, model == tracingModelFirehose2_3))) + }`, concurrentBlockFlushing, model == tracingModelFirehose2_3))) require.NoError(t, err) hooks := tracers.NewTracingHooksFromFirehose(tracer)