diff --git a/eth/filters/api.go b/eth/filters/api.go index 0eb977cadf..ebb2a45a2f 100644 --- a/eth/filters/api.go +++ b/eth/filters/api.go @@ -271,13 +271,47 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc var ( rpcSub = notifier.CreateSubscription() matchedLogs = make(chan []*types.Log) - rmLogsCh = make(chan []*types.Log) ) - syncHistLogs := func(n int64) error { - // the ctx will be canceled after this function is called + filterLiveLogs := func() error { + logsSub, err := api.events.SubscribeLogs(ethereum.FilterQuery(crit), matchedLogs) + if err != nil { + return err + } + + go func() { + for { + select { + case logs := <-matchedLogs: + for _, log := range logs { + log := log + notifier.Notify(rpcSub.ID, &log) + } + case <-rpcSub.Err(): // client send an unsubscribe request + logsSub.Unsubscribe() + return + case <-notifier.Closed(): // connection dropped + logsSub.Unsubscribe() + return + } + } + }() + return nil + } + + // do live filter only + if crit.FromBlock == nil { + err := filterLiveLogs() + return rpcSub, err + } + + filterHistLogs := func(n int64) error { + // the ctx will be canceled after this callback(filter.Logs) is called // and we are run in a background goroutine, use a new context instead - cctx := context.Background() + var ( + cctx = context.Background() + rmLogsCh = make(chan []*types.Log) + ) for { // get the latest block header head := api.sys.backend.CurrentHeader().Number.Int64() @@ -331,47 +365,16 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc } } - syncLiveLogs := func() error { - logsSub, err := api.events.SubscribeLogs(ethereum.FilterQuery(crit), matchedLogs) - if err != nil { - return err + n := crit.FromBlock.Int64() + go func() { + // do historical sync first + if err := filterHistLogs(n); err != nil { + return } - - go func() { - for { - select { - case logs := <-matchedLogs: - for _, log := range logs { - log := log - notifier.Notify(rpcSub.ID, &log) - } - case <-rpcSub.Err(): // client send an unsubscribe request - logsSub.Unsubscribe() - return - case <-notifier.Closed(): // connection dropped - logsSub.Unsubscribe() - return - } - } - }() - return nil - } - - if crit.FromBlock != nil { - n := crit.FromBlock.Int64() - go func() { - // do historical sync first - if err := syncHistLogs(n); err != nil { - return - } - // subscribe from latest - syncLiveLogs() - }() - return rpcSub, nil - } - - err := syncLiveLogs() - return rpcSub, err + // then subscribe from the header + filterLiveLogs() + }() + return rpcSub, nil } // FilterCriteria represents a request to create a new filter.