forked from bnb-chain/bnc-tendermint
-
Notifications
You must be signed in to change notification settings - Fork 0
/
indexer_service.go
67 lines (54 loc) · 1.68 KB
/
indexer_service.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
package blockindex
import (
"context"
cmn "github.com/tendermint/tendermint/libs/common"
"github.com/tendermint/tendermint/types"
)
const (
subscriber = "BlockIndexerService"
)
// IndexerService connects event bus and block indexer together in order
// to index blocks coming from event bus.
type IndexerService struct {
cmn.BaseService
idr BlockIndexer
eventBus *types.EventBus
onIndex func(int64)
}
// NewIndexerService returns a new service instance.
func NewIndexerService(idr BlockIndexer, eventBus *types.EventBus) *IndexerService {
is := &IndexerService{idr: idr, eventBus: eventBus}
is.BaseService = *cmn.NewBaseService(nil, "BlockIndexerService", is)
return is
}
func (is *IndexerService) SetOnIndex(callback func(int64)) {
is.onIndex = callback
}
// OnStart implements cmn.Service by subscribing for blocks and indexing them by hash.
func (is *IndexerService) OnStart() error {
blockHeadersSub, err := is.eventBus.SubscribeUnbuffered(context.Background(), subscriber, types.EventQueryNewBlockHeader)
if err != nil {
return err
}
go func() {
for {
msg := <-blockHeadersSub.Out()
header := msg.Data().(types.EventDataNewBlockHeader).Header
if err := is.idr.Index(&header); err != nil {
is.Logger.Error("Failed to index block", "height", header.Height, "err", err)
} else {
is.Logger.Info("Indexed block", "height", header.Height, "hash", header.LastBlockID.Hash)
}
if is.onIndex != nil {
is.onIndex(header.Height)
}
}
}()
return nil
}
// OnStop implements cmn.Service by unsubscribing from blocks.
func (is *IndexerService) OnStop() {
if is.eventBus.IsRunning() {
_ = is.eventBus.UnsubscribeAll(context.Background(), subscriber)
}
}