mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 06:36:43 +00:00
DEV: Implementation of concurrency in firehose
This commit is contained in:
parent
f70862f7e4
commit
4c868888a7
2 changed files with 47 additions and 1 deletions
|
|
@ -186,6 +186,10 @@ type Firehose struct {
|
||||||
// Testing state, only used in tests and private configs
|
// Testing state, only used in tests and private configs
|
||||||
testingBuffer *bytes.Buffer
|
testingBuffer *bytes.Buffer
|
||||||
testingIgnoreGenesisBlock bool
|
testingIgnoreGenesisBlock bool
|
||||||
|
|
||||||
|
// Worker queuefor blockprinting
|
||||||
|
blockPrintQueue chan *blockPrintJob
|
||||||
|
workerWg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
const FirehoseProtocolVersion = "3.0"
|
const FirehoseProtocolVersion = "3.0"
|
||||||
|
|
@ -238,6 +242,8 @@ func NewFirehose(config *FirehoseConfig) *Firehose {
|
||||||
callStack: NewCallStack(),
|
callStack: NewCallStack(),
|
||||||
deferredCallState: NewDeferredCallState(),
|
deferredCallState: NewDeferredCallState(),
|
||||||
latestCallEnterSuicided: false,
|
latestCallEnterSuicided: false,
|
||||||
|
|
||||||
|
blockPrintQueue: make(chan *blockPrintJob, 100),
|
||||||
}
|
}
|
||||||
|
|
||||||
if config.private != nil {
|
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
|
return firehose
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -483,7 +496,13 @@ func (f *Firehose) OnBlockEnd(err error) {
|
||||||
}
|
}
|
||||||
|
|
||||||
f.ensureInBlockAndNotInTrx()
|
f.ensureInBlockAndNotInTrx()
|
||||||
f.printBlockToFirehose(f.block, f.blockFinality)
|
|
||||||
|
job := &blockPrintJob{
|
||||||
|
block: f.block,
|
||||||
|
finality: f.blockFinality,
|
||||||
|
}
|
||||||
|
f.blockPrintQueue <- job
|
||||||
|
|
||||||
} else {
|
} else {
|
||||||
// An error occurred, could have happen in transaction/call context, we must not check if in trx/call, only check in block
|
// 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)
|
f.ensureInBlock(0)
|
||||||
|
|
@ -2791,3 +2810,21 @@ func (m Memory) GetPtr(offset, size int64) []byte {
|
||||||
reminder := m[offset:]
|
reminder := m[offset:]
|
||||||
return append(reminder, make([]byte, int(size)-len(reminder))...)
|
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()
|
||||||
|
}
|
||||||
|
|
|
||||||
9
eth/tracers/firehose_concurrency_test.go
Normal file
9
eth/tracers/firehose_concurrency_test.go
Normal file
|
|
@ -0,0 +1,9 @@
|
||||||
|
package tracers
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestFirehoseBlockPrintingOrder(t *testing.T) {
|
||||||
|
|
||||||
|
}
|
||||||
Loading…
Reference in a new issue