eth/filters: use shared code to serve filter tasks, fix tests

This commit is contained in:
Zsolt Felfoldi 2018-03-02 15:56:47 +01:00
parent fc6fa6cb0c
commit d68757ae3b
7 changed files with 54 additions and 67 deletions

View file

@ -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) filter = filters.NewRangeFilter(&filterBackend{b.database, b.blockchain, filterTasks}, from, to, query.Addresses, query.Topics)
// start up a filter task service // start up a filter task service
const filterTaskServiceThreads = 16 stop := make(chan bool)
stop := make(chan struct{})
defer close(stop) defer close(stop)
filters.ServeFilterTasks(filterTasks, stop)
for i := 0; i < filterTaskServiceThreads; i++ {
go func() {
for {
select {
case <-stop:
return
case task := <-filterTasks:
task.Do()
}
}
}()
}
} }
// Run the filter and return all the logs // Run the filter and return all the logs
logs, err := filter.Logs(ctx) logs, err := filter.Logs(ctx)

View file

@ -26,6 +26,7 @@ import (
"github.com/ethereum/go-ethereum/core/bloombits" "github.com/ethereum/go-ethereum/core/bloombits"
"github.com/ethereum/go-ethereum/core/rawdb" "github.com/ethereum/go-ethereum/core/rawdb"
"github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/eth/filters"
"github.com/ethereum/go-ethereum/ethdb" "github.com/ethereum/go-ethereum/ethdb"
) )
@ -34,10 +35,6 @@ const (
// instance to service bloombits lookups for all running filters. // instance to service bloombits lookups for all running filters.
bloomServiceThreads = 16 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 // bloomFilterThreads is the number of goroutines used locally per filter to
// multiplex requests onto the global servicing goroutines. // multiplex requests onto the global servicing goroutines.
bloomFilterThreads = 3 bloomFilterThreads = 3
@ -82,19 +79,7 @@ func (eth *Ethereum) startBloomHandlers(sectionSize uint64) {
}() }()
} }
for i := 0; i < filterTaskServiceThreads; i++ { filters.ServeFilterTasks(eth.filterTasks, eth.shutdownChan)
go func() {
for {
select {
case <-eth.shutdownChan:
return
case task := <-eth.filterTasks:
task.Do()
}
}
}()
}
} }
const ( const (

View file

@ -130,7 +130,7 @@ func benchmarkBloomBits(b *testing.B, sectionSize uint64) {
if i%20 == 0 { if i%20 == 0 {
db.Close() db.Close()
db, _ = ethdb.NewLDBDatabase(benchDataDir, 128, 1024) 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 var addr common.Address
addr[0] = byte(i) addr[0] = byte(i)
@ -191,7 +191,7 @@ func BenchmarkNoBloomBits(b *testing.B) {
fmt.Println("Running filter benchmarks...") fmt.Println("Running filter benchmarks...")
start := time.Now() start := time.Now()
mux := new(event.TypeMux) 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 := NewRangeFilter(backend, 0, int64(*headNum), []common.Address{{}}, nil)
filter.Logs(context.Background()) filter.Logs(context.Background())
d := time.Since(start) d := time.Since(start)

View file

@ -66,6 +66,7 @@ type BlockFilterTask struct {
done chan struct{} done chan struct{}
} }
// Do is a callback function called for each block filter task by ServeFilterTasks
func (t *BlockFilterTask) Do() { func (t *BlockFilterTask) Do() {
defer close(t.done) 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. // Filter can be used to retrieve and filter logs.
type Filter struct { type Filter struct {
backend Backend backend Backend

View file

@ -46,6 +46,7 @@ type testBackend struct {
rmLogsFeed *event.Feed rmLogsFeed *event.Feed
logsFeed *event.Feed logsFeed *event.Feed
chainFeed *event.Feed chainFeed *event.Feed
filterTasks chan *BlockFilterTask
} }
func (b *testBackend) ChainDb() ethdb.Database { func (b *testBackend) ChainDb() ethdb.Database {
@ -56,6 +57,14 @@ func (b *testBackend) EventMux() *event.TypeMux {
return b.mux 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) { func (b *testBackend) HeaderByNumber(ctx context.Context, blockNr rpc.BlockNumber) (*types.Header, error) {
var ( var (
hash common.Hash hash common.Hash
@ -166,7 +175,7 @@ func TestBlockSubscription(t *testing.T) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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) api = NewPublicFilterAPI(backend, false)
genesis = new(core.Genesis).MustCommit(db) genesis = new(core.Genesis).MustCommit(db)
chain, _ = core.GenerateChain(params.TestChainConfig, genesis, ethash.NewFaker(), db, 10, func(i int, gen *core.BlockGen) {}) 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) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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) api = NewPublicFilterAPI(backend, false)
transactions = []*types.Transaction{ transactions = []*types.Transaction{
@ -283,7 +292,7 @@ func TestLogFilterCreation(t *testing.T) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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) api = NewPublicFilterAPI(backend, false)
testCases = []struct { testCases = []struct {
@ -332,7 +341,7 @@ func TestInvalidLogFilterCreation(t *testing.T) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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) api = NewPublicFilterAPI(backend, false)
) )
@ -389,7 +398,7 @@ func TestLogFilter(t *testing.T) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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) api = NewPublicFilterAPI(backend, false)
firstAddr = common.HexToAddress("0x1111111111111111111111111111111111111111") firstAddr = common.HexToAddress("0x1111111111111111111111111111111111111111")
@ -508,7 +517,7 @@ func TestPendingLogsSubscription(t *testing.T) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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) api = NewPublicFilterAPI(backend, false)
firstAddr = common.HexToAddress("0x1111111111111111111111111111111111111111") firstAddr = common.HexToAddress("0x1111111111111111111111111111111111111111")

View file

@ -57,7 +57,7 @@ func BenchmarkFilters(b *testing.B) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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") key1, _ = crypto.HexToECDSA("b71c71a67e1177ad4e901695e1b4b9ee17ae16c6668d313eac2f96dbcda3f291")
addr1 = crypto.PubkeyToAddress(key1.PublicKey) addr1 = crypto.PubkeyToAddress(key1.PublicKey)
addr2 = common.BytesToAddress([]byte("jeff")) addr2 = common.BytesToAddress([]byte("jeff"))
@ -116,7 +116,7 @@ func TestFilters(t *testing.T) {
rmLogsFeed = new(event.Feed) rmLogsFeed = new(event.Feed)
logsFeed = new(event.Feed) logsFeed = new(event.Feed)
chainFeed = 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") key1, _ = crypto.HexToECDSA("b71c71a67e1177ad4e901695e1b4b9ee17ae16c6668d313eac2f96dbcda3f291")
addr = crypto.PubkeyToAddress(key1.PublicKey) addr = crypto.PubkeyToAddress(key1.PublicKey)

View file

@ -20,6 +20,7 @@ import (
"time" "time"
"github.com/ethereum/go-ethereum/common/bitutil" "github.com/ethereum/go-ethereum/common/bitutil"
"github.com/ethereum/go-ethereum/eth/filters"
"github.com/ethereum/go-ethereum/light" "github.com/ethereum/go-ethereum/light"
) )
@ -28,10 +29,6 @@ const (
// instance to service bloombits lookups for all running filters. // instance to service bloombits lookups for all running filters.
bloomServiceThreads = 16 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 // bloomFilterThreads is the number of goroutines used locally per filter to
// multiplex requests onto the global servicing goroutines. // multiplex requests onto the global servicing goroutines.
bloomFilterThreads = 3 bloomFilterThreads = 3
@ -76,17 +73,5 @@ func (eth *LightEthereum) startBloomHandlers(sectionSize uint64) {
}() }()
} }
for i := 0; i < filterTaskServiceThreads; i++ { filters.ServeFilterTasks(eth.filterTasks, eth.shutdownChan)
go func() {
for {
select {
case <-eth.shutdownChan:
return
case task := <-eth.filterTasks:
task.Do()
}
}
}()
}
} }