|
|
|
@ -159,6 +159,8 @@ func (api *FilterAPI) NewPendingTransactions(ctx context.Context, fullTx *bool) |
|
|
|
|
go func() { |
|
|
|
|
txs := make(chan []*types.Transaction, 128) |
|
|
|
|
pendingTxSub := api.events.SubscribePendingTxs(txs) |
|
|
|
|
defer pendingTxSub.Unsubscribe() |
|
|
|
|
|
|
|
|
|
chainConfig := api.sys.backend.ChainConfig() |
|
|
|
|
|
|
|
|
|
for { |
|
|
|
@ -176,10 +178,8 @@ func (api *FilterAPI) NewPendingTransactions(ctx context.Context, fullTx *bool) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
case <-rpcSub.Err(): |
|
|
|
|
pendingTxSub.Unsubscribe() |
|
|
|
|
return |
|
|
|
|
case <-notifier.Closed(): |
|
|
|
|
pendingTxSub.Unsubscribe() |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
@ -233,16 +233,15 @@ func (api *FilterAPI) NewHeads(ctx context.Context) (*rpc.Subscription, error) { |
|
|
|
|
go func() { |
|
|
|
|
headers := make(chan *types.Header) |
|
|
|
|
headersSub := api.events.SubscribeNewHeads(headers) |
|
|
|
|
defer headersSub.Unsubscribe() |
|
|
|
|
|
|
|
|
|
for { |
|
|
|
|
select { |
|
|
|
|
case h := <-headers: |
|
|
|
|
notifier.Notify(rpcSub.ID, h) |
|
|
|
|
case <-rpcSub.Err(): |
|
|
|
|
headersSub.Unsubscribe() |
|
|
|
|
return |
|
|
|
|
case <-notifier.Closed(): |
|
|
|
|
headersSub.Unsubscribe() |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
@ -267,6 +266,7 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc |
|
|
|
|
if err != nil { |
|
|
|
|
return nil, err |
|
|
|
|
} |
|
|
|
|
defer logsSub.Unsubscribe() |
|
|
|
|
|
|
|
|
|
go func() { |
|
|
|
|
for { |
|
|
|
@ -277,10 +277,8 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc |
|
|
|
|
notifier.Notify(rpcSub.ID, &log) |
|
|
|
|
} |
|
|
|
|
case <-rpcSub.Err(): // client send an unsubscribe request
|
|
|
|
|
logsSub.Unsubscribe() |
|
|
|
|
return |
|
|
|
|
case <-notifier.Closed(): // connection dropped
|
|
|
|
|
logsSub.Unsubscribe() |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|