mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-23 13:16:42 +00:00
eth/filters: wrap Notify, for unit test
Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
parent
c93423eb60
commit
964e801961
1 changed files with 13 additions and 7 deletions
|
|
@ -261,19 +261,25 @@ func (api *FilterAPI) NewHeads(ctx context.Context) (*rpc.Subscription, error) {
|
||||||
return rpcSub, nil
|
return rpcSub, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type notify interface {
|
||||||
|
Notify(id rpc.ID, data interface{}) error
|
||||||
|
Closed() <-chan interface{}
|
||||||
|
}
|
||||||
|
|
||||||
// Logs creates a subscription that fires for all new log that match the given filter criteria.
|
// Logs creates a subscription that fires for all new log that match the given filter criteria.
|
||||||
func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subscription, error) {
|
func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subscription, error) {
|
||||||
notifier, supported := rpc.NotifierFromContext(ctx)
|
notifier, supported := rpc.NotifierFromContext(ctx)
|
||||||
if !supported {
|
if !supported {
|
||||||
return &rpc.Subscription{}, rpc.ErrNotificationsUnsupported
|
return &rpc.Subscription{}, rpc.ErrNotificationsUnsupported
|
||||||
}
|
}
|
||||||
|
rpcSub := notifier.CreateSubscription()
|
||||||
|
err := api.logs(ctx, notifier, rpcSub, crit)
|
||||||
|
return rpcSub, err
|
||||||
|
}
|
||||||
|
|
||||||
var (
|
func (api *FilterAPI) logs(ctx context.Context, notifier notify, rpcSub *rpc.Subscription, crit FilterCriteria) error {
|
||||||
rpcSub = notifier.CreateSubscription()
|
|
||||||
matchedLogs = make(chan []*types.Log)
|
|
||||||
)
|
|
||||||
|
|
||||||
filterLiveLogs := func() error {
|
filterLiveLogs := func() error {
|
||||||
|
matchedLogs := make(chan []*types.Log)
|
||||||
logsSub, err := api.events.SubscribeLogs(ethereum.FilterQuery(crit), matchedLogs)
|
logsSub, err := api.events.SubscribeLogs(ethereum.FilterQuery(crit), matchedLogs)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
|
|
@ -306,7 +312,7 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc
|
||||||
// do live filter only
|
// do live filter only
|
||||||
if isLiveFilter {
|
if isLiveFilter {
|
||||||
err := filterLiveLogs()
|
err := filterLiveLogs()
|
||||||
return rpcSub, err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
filterHistLogs := func(n int64) (int64, error) {
|
filterHistLogs := func(n int64) (int64, error) {
|
||||||
|
|
@ -398,7 +404,7 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc
|
||||||
// then subscribe from the header
|
// then subscribe from the header
|
||||||
filterLiveLogs()
|
filterLiveLogs()
|
||||||
}()
|
}()
|
||||||
return rpcSub, nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// FilterCriteria represents a request to create a new filter.
|
// FilterCriteria represents a request to create a new filter.
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue