mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-23 13:16:42 +00:00
eth/filters: do hist sync
Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
parent
92c980f355
commit
dc6a5b1b3d
1 changed files with 39 additions and 13 deletions
|
|
@ -271,34 +271,60 @@ func (api *FilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subsc
|
||||||
var (
|
var (
|
||||||
rpcSub = notifier.CreateSubscription()
|
rpcSub = notifier.CreateSubscription()
|
||||||
matchedLogs = make(chan []*types.Log)
|
matchedLogs = make(chan []*types.Log)
|
||||||
|
rmLogsCh = make(chan []*types.Log)
|
||||||
)
|
)
|
||||||
|
|
||||||
syncHistLogs := func(n int64) error {
|
syncHistLogs := func(n int64) error {
|
||||||
for {
|
for {
|
||||||
// Get latest head block num
|
// get the latest block header
|
||||||
head := api.sys.backend.CurrentHeader().Number.Int64()
|
head := api.sys.backend.CurrentHeader().Number.Int64()
|
||||||
if n >= head {
|
if n >= head {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
// Do historical sync from n to head
|
|
||||||
|
// do historical sync from n to head
|
||||||
f := api.sys.NewRangeFilter(n, head, crit.Addresses, crit.Topics)
|
f := api.sys.NewRangeFilter(n, head, crit.Addresses, crit.Topics)
|
||||||
logChan, errChan := f.rangeLogsAsync(ctx)
|
logChan, errChan := f.rangeLogsAsync(ctx)
|
||||||
select {
|
|
||||||
case <-notifier.Closed():
|
// subscribe rmLogs
|
||||||
return errors.New("connection dropped")
|
query := ethereum.FilterQuery{FromBlock: big.NewInt(n), ToBlock: big.NewInt(head), Addresses: crit.Addresses, Topics: crit.Topics}
|
||||||
case log := <-logChan:
|
rmLogsSub, err := api.events.SubscribeLogs(query, rmLogsCh)
|
||||||
notifier.Notify(rpcSub.ID, &log)
|
if err != nil {
|
||||||
case err := <-errChan:
|
return err
|
||||||
if err != nil {
|
}
|
||||||
return err
|
|
||||||
|
var rmLogs []*types.Log
|
||||||
|
FORLOOP:
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-notifier.Closed():
|
||||||
|
rmLogsSub.Unsubscribe()
|
||||||
|
return errors.New("connection dropped")
|
||||||
|
case log := <-logChan:
|
||||||
|
notifier.Notify(rpcSub.ID, &log)
|
||||||
|
case logs := <-rmLogsCh:
|
||||||
|
for _, log := range logs {
|
||||||
|
if log.Removed {
|
||||||
|
rmLogs = append(rmLogs, log)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
case err := <-errChan:
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// range filter is done, let's also stop the rmLogs subscribe
|
||||||
|
rmLogsSub.Unsubscribe()
|
||||||
|
break FORLOOP
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TODO: Check if any logs from n to head have been reorged
|
|
||||||
// send the reorged logs
|
// send the reorged logs
|
||||||
|
for _, log := range rmLogs {
|
||||||
|
notifier.Notify(rpcSub.ID, &log)
|
||||||
|
}
|
||||||
|
|
||||||
// Head is stable, move n to current head
|
// move forward to the next batch
|
||||||
n = head
|
n = head + 1
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue