-
Notifications
You must be signed in to change notification settings - Fork 0
/
aggregator.go
135 lines (110 loc) · 2.75 KB
/
aggregator.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
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
package aggregator
import (
"context"
"encoding/json"
"os"
"os/signal"
"sync"
"syscall"
"time"
"github.com/dezswap/cosmwasm-etl/aggregator/repo"
"github.com/dezswap/cosmwasm-etl/configs"
"github.com/dezswap/cosmwasm-etl/pkg/db/parser"
"github.com/dezswap/cosmwasm-etl/pkg/logging"
"github.com/sirupsen/logrus"
)
type Aggregator interface {
Run() error
}
type aggregatorImpl struct {
chainId string
startTs time.Time
cleanDups bool
srcDbConn parser.ReadRepository
destDbConn repo.Repo
tasks []task
logger logging.Logger
}
var (
_ Aggregator = &aggregatorImpl{}
errChan chan error
)
func New(c configs.Config, logger logging.Logger) Aggregator {
repo.Logger = logger
if c.Log.Level == logrus.DebugLevel {
a, err := json.Marshal(c.Aggregator)
if err != nil {
panic(err)
}
logger.Debug(string(a))
}
srcRepo := parser.NewReadRepo(c.Aggregator.ChainId, c.Aggregator.SrcDb)
destRepo := repo.New(c.Aggregator.ChainId, c.Aggregator.DestDb)
// init tasks
rt, err := newRouterTask(c.Aggregator, logger)
if err != nil {
panic(err)
}
lht := newLpHistoryTask(c.Aggregator, srcRepo, destRepo, logger)
pt, err := newPriceTask(c.Aggregator, destRepo, logger, []task{lht})
if err != nil {
panic(err)
}
tasks := []task{
rt, lht, pt,
newPairStatsRecentUpdateTask(c.Aggregator, srcRepo, destRepo, logger, []task{pt}),
newPairStatsUpdateTask(c.Aggregator, srcRepo, destRepo, logger, []task{pt}),
}
return &aggregatorImpl{
chainId: c.Aggregator.ChainId,
startTs: c.Aggregator.StartTs,
cleanDups: c.Aggregator.CleanDups,
srcDbConn: srcRepo,
destDbConn: destRepo,
tasks: tasks,
logger: logger,
}
}
func (a aggregatorImpl) Run() error {
a.logger.Info("Aggregator has been started.")
defer a.srcDbConn.Close()
defer a.destDbConn.Close()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
sigs := make(chan os.Signal, 1)
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
go func() {
sig := <-sigs
a.logger.Infof("Signal has been caught: %s", sig.String())
cancel()
}()
errChan = make(chan error)
go func() {
err := <-errChan
a.logger.Errorf("Error has been caught: %s", err.Error())
cancel()
}()
a.runTasks(ctx)
a.logger.Info("cosmwasm-etc aggregator has been stopped.")
return nil
}
func (a aggregatorImpl) runTasks(ctx context.Context) {
wg := sync.WaitGroup{}
if a.cleanDups {
if err := a.destDbConn.DeleteDuplicates(a.startTs); err != nil {
errChan <- err
return
}
a.logger.Infof("Stats data since %s has been deleted for new update.", a.startTs.String())
}
for _, t := range a.tasks {
wg.Add(1)
go func(tk task) {
defer wg.Done()
if err := tk.Schedule(ctx, a.startTs); err != nil {
errChan <- err
}
}(t)
}
wg.Wait()
}