From d68757ae3ba7bd633ba6a4221562d6ddd48f569c Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Fri, 2 Mar 2018 15:56:47 +0100 Subject: [PATCH] eth/filters: use shared code to serve filter tasks, fix tests --- accounts/abi/bind/backends/simulated.go | 18 ++----------- eth/bloombits.go | 19 ++------------ eth/filters/bench_test.go | 4 +-- eth/filters/filter.go | 22 ++++++++++++++++ eth/filters/filter_system_test.go | 35 ++++++++++++++++--------- eth/filters/filter_test.go | 4 +-- les/bloombits.go | 19 ++------------ 7 files changed, 54 insertions(+), 67 deletions(-) diff --git a/accounts/abi/bind/backends/simulated.go b/accounts/abi/bind/backends/simulated.go index 9ad23d1fc2..990cee9d77 100644 --- a/accounts/abi/bind/backends/simulated.go +++ b/accounts/abi/bind/backends/simulated.go @@ -362,23 +362,9 @@ func (b *SimulatedBackend) FilterLogs(ctx context.Context, query ethereum.Filter filter = filters.NewRangeFilter(&filterBackend{b.database, b.blockchain, filterTasks}, from, to, query.Addresses, query.Topics) // start up a filter task service - const filterTaskServiceThreads = 16 - stop := make(chan struct{}) + stop := make(chan bool) defer close(stop) - - for i := 0; i < filterTaskServiceThreads; i++ { - go func() { - for { - select { - case <-stop: - return - - case task := <-filterTasks: - task.Do() - } - } - }() - } + filters.ServeFilterTasks(filterTasks, stop) } // Run the filter and return all the logs logs, err := filter.Logs(ctx) diff --git a/eth/bloombits.go b/eth/bloombits.go index 122033fc5e..33a4c2d793 100644 --- a/eth/bloombits.go +++ b/eth/bloombits.go @@ -26,6 +26,7 @@ import ( "github.com/ethereum/go-ethereum/core/bloombits" "github.com/ethereum/go-ethereum/core/rawdb" "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/eth/filters" "github.com/ethereum/go-ethereum/ethdb" ) @@ -34,10 +35,6 @@ const ( // instance to service bloombits lookups for all running filters. bloomServiceThreads = 16 - // filterTaskServiceThreads is the number of goroutines used globally by an Ethereum - // instance to service block filter tasks for all running filters. - filterTaskServiceThreads = 16 - // bloomFilterThreads is the number of goroutines used locally per filter to // multiplex requests onto the global servicing goroutines. bloomFilterThreads = 3 @@ -82,19 +79,7 @@ func (eth *Ethereum) startBloomHandlers(sectionSize uint64) { }() } - for i := 0; i < filterTaskServiceThreads; i++ { - go func() { - for { - select { - case <-eth.shutdownChan: - return - - case task := <-eth.filterTasks: - task.Do() - } - } - }() - } + filters.ServeFilterTasks(eth.filterTasks, eth.shutdownChan) } const ( diff --git a/eth/filters/bench_test.go b/eth/filters/bench_test.go index c5f681e024..add30aae0d 100644 --- a/eth/filters/bench_test.go +++ b/eth/filters/bench_test.go @@ -130,7 +130,7 @@ func benchmarkBloomBits(b *testing.B, sectionSize uint64) { if i%20 == 0 { db.Close() db, _ = ethdb.NewLDBDatabase(benchDataDir, 128, 1024) - backend = &testBackend{mux, db, cnt, new(event.Feed), new(event.Feed), new(event.Feed), new(event.Feed)} + backend = &testBackend{mux, db, cnt, new(event.Feed), new(event.Feed), new(event.Feed), new(event.Feed), nil} } var addr common.Address addr[0] = byte(i) @@ -191,7 +191,7 @@ func BenchmarkNoBloomBits(b *testing.B) { fmt.Println("Running filter benchmarks...") start := time.Now() mux := new(event.TypeMux) - backend := &testBackend{mux, db, 0, new(event.Feed), new(event.Feed), new(event.Feed), new(event.Feed)} + backend := &testBackend{mux, db, 0, new(event.Feed), new(event.Feed), new(event.Feed), new(event.Feed), nil} filter := NewRangeFilter(backend, 0, int64(*headNum), []common.Address{{}}, nil) filter.Logs(context.Background()) d := time.Since(start) diff --git a/eth/filters/filter.go b/eth/filters/filter.go index c7234f7956..33471bff3a 100644 --- a/eth/filters/filter.go +++ b/eth/filters/filter.go @@ -66,6 +66,7 @@ type BlockFilterTask struct { done chan struct{} } +// Do is a callback function called for each block filter task by ServeFilterTasks func (t *BlockFilterTask) Do() { defer close(t.done) @@ -83,6 +84,27 @@ func (t *BlockFilterTask) Do() { } } +// filterTaskServiceThreads is the number of goroutines used globally by a filter +// backend to service block filter tasks for all running filters. +const filterTaskServiceThreads = 16 + +// ServeFilterTasks is called by the filter backend to start up filter task service threads +func ServeFilterTasks(tasks chan *BlockFilterTask, stop chan bool) { + for i := 0; i < filterTaskServiceThreads; i++ { + go func() { + for { + select { + case <-stop: + return + + case task := <-tasks: + task.Do() + } + } + }() + } +} + // Filter can be used to retrieve and filter logs. type Filter struct { backend Backend diff --git a/eth/filters/filter_system_test.go b/eth/filters/filter_system_test.go index e71080b1ab..f9652ee2c3 100644 --- a/eth/filters/filter_system_test.go +++ b/eth/filters/filter_system_test.go @@ -39,13 +39,14 @@ import ( ) type testBackend struct { - mux *event.TypeMux - db ethdb.Database - sections uint64 - txFeed *event.Feed - rmLogsFeed *event.Feed - logsFeed *event.Feed - chainFeed *event.Feed + mux *event.TypeMux + db ethdb.Database + sections uint64 + txFeed *event.Feed + rmLogsFeed *event.Feed + logsFeed *event.Feed + chainFeed *event.Feed + filterTasks chan *BlockFilterTask } func (b *testBackend) ChainDb() ethdb.Database { @@ -56,6 +57,14 @@ func (b *testBackend) EventMux() *event.TypeMux { return b.mux } +func (b *testBackend) FilterTaskChannel() chan *BlockFilterTask { + if b.filterTasks == nil { + b.filterTasks = make(chan *BlockFilterTask) + ServeFilterTasks(b.filterTasks, nil) + } + return b.filterTasks +} + func (b *testBackend) HeaderByNumber(ctx context.Context, blockNr rpc.BlockNumber) (*types.Header, error) { var ( hash common.Hash @@ -166,7 +175,7 @@ func TestBlockSubscription(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} api = NewPublicFilterAPI(backend, false) genesis = new(core.Genesis).MustCommit(db) chain, _ = core.GenerateChain(params.TestChainConfig, genesis, ethash.NewFaker(), db, 10, func(i int, gen *core.BlockGen) {}) @@ -223,7 +232,7 @@ func TestPendingTxFilter(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} api = NewPublicFilterAPI(backend, false) transactions = []*types.Transaction{ @@ -283,7 +292,7 @@ func TestLogFilterCreation(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} api = NewPublicFilterAPI(backend, false) testCases = []struct { @@ -332,7 +341,7 @@ func TestInvalidLogFilterCreation(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} api = NewPublicFilterAPI(backend, false) ) @@ -389,7 +398,7 @@ func TestLogFilter(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} api = NewPublicFilterAPI(backend, false) firstAddr = common.HexToAddress("0x1111111111111111111111111111111111111111") @@ -508,7 +517,7 @@ func TestPendingLogsSubscription(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} api = NewPublicFilterAPI(backend, false) firstAddr = common.HexToAddress("0x1111111111111111111111111111111111111111") diff --git a/eth/filters/filter_test.go b/eth/filters/filter_test.go index 396a03d611..2a51c547db 100644 --- a/eth/filters/filter_test.go +++ b/eth/filters/filter_test.go @@ -57,7 +57,7 @@ func BenchmarkFilters(b *testing.B) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} key1, _ = crypto.HexToECDSA("b71c71a67e1177ad4e901695e1b4b9ee17ae16c6668d313eac2f96dbcda3f291") addr1 = crypto.PubkeyToAddress(key1.PublicKey) addr2 = common.BytesToAddress([]byte("jeff")) @@ -116,7 +116,7 @@ func TestFilters(t *testing.T) { rmLogsFeed = new(event.Feed) logsFeed = new(event.Feed) chainFeed = new(event.Feed) - backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed} + backend = &testBackend{mux, db, 0, txFeed, rmLogsFeed, logsFeed, chainFeed, nil} key1, _ = crypto.HexToECDSA("b71c71a67e1177ad4e901695e1b4b9ee17ae16c6668d313eac2f96dbcda3f291") addr = crypto.PubkeyToAddress(key1.PublicKey) diff --git a/les/bloombits.go b/les/bloombits.go index 3b47cdc3d6..51db1fdae0 100644 --- a/les/bloombits.go +++ b/les/bloombits.go @@ -20,6 +20,7 @@ import ( "time" "github.com/ethereum/go-ethereum/common/bitutil" + "github.com/ethereum/go-ethereum/eth/filters" "github.com/ethereum/go-ethereum/light" ) @@ -28,10 +29,6 @@ const ( // instance to service bloombits lookups for all running filters. bloomServiceThreads = 16 - // filterTaskServiceThreads is the number of goroutines used globally by an Ethereum - // instance to service block filter tasks for all running filters. - filterTaskServiceThreads = 16 - // bloomFilterThreads is the number of goroutines used locally per filter to // multiplex requests onto the global servicing goroutines. bloomFilterThreads = 3 @@ -76,17 +73,5 @@ func (eth *LightEthereum) startBloomHandlers(sectionSize uint64) { }() } - for i := 0; i < filterTaskServiceThreads; i++ { - go func() { - for { - select { - case <-eth.shutdownChan: - return - - case task := <-eth.filterTasks: - task.Do() - } - } - }() - } + filters.ServeFilterTasks(eth.filterTasks, eth.shutdownChan) }