-
Notifications
You must be signed in to change notification settings - Fork 0
/
notifier.go
150 lines (126 loc) · 4.24 KB
/
notifier.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
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
package bundle
import (
"encoding/json"
"github.com/warp-contracts/syncer/src/utils/config"
"github.com/warp-contracts/syncer/src/utils/model"
"github.com/warp-contracts/syncer/src/utils/monitoring"
"github.com/warp-contracts/syncer/src/utils/streamer"
"github.com/warp-contracts/syncer/src/utils/task"
"gorm.io/gorm"
)
// Gets a live stream of unbundled intearctions, parses them and puts them on the output channel
type Notifier struct {
*task.Task
db *gorm.DB
streamer *streamer.Streamer
monitor monitoring.Monitor
// Data about the interactions that need to be bundled
output chan *model.BundleItem
}
func NewNotifier(config *config.Config) (self *Notifier) {
self = new(Notifier)
if config.Bundler.NotifierDisabled {
self.Task = task.NewTask(config, "bundler-notifier")
return
}
self.streamer = streamer.NewStreamer(config, "bundler-notifier").
WithNotificationChannelName("bundle_items_pending").
WithCapacity(10)
self.Task = task.NewTask(config, "notifier").
// Live source of interactions that need to be bundled
WithSubtask(self.streamer.Task).
// Interactions that somehow wasn't sent through the notification channel. Probably because of a restart.
WithSubtaskFunc(self.run).
// Workers unmarshal big JSON messages and optionally fetch data from the database if the messages wuldn't fit in the notification channel
WithWorkerPool(config.Bundler.NotifierWorkerPoolSize, config.Bundler.NotifierWorkerQueueSize)
return
}
func (self *Notifier) WithDB(db *gorm.DB) *Notifier {
self.db = db
return self
}
func (self *Notifier) WithMonitor(monitor monitoring.Monitor) *Notifier {
self.monitor = monitor
return self
}
func (self *Notifier) WithOutputChannel(bundleItems chan *model.BundleItem) *Notifier {
self.output = bundleItems
return self
}
func (self *Notifier) run() error {
for {
select {
case <-self.StopChannel:
self.Log.Debug("Stop passing interactions from notification")
return nil
case msg, ok := <-self.streamer.Output:
if !ok {
self.Log.Info("Notification streamer channel closed")
return nil
}
self.SubmitToWorker(func() {
var notification model.BundleItemNotification
err := json.Unmarshal([]byte(msg), ¬ification)
if err != nil {
self.Log.WithError(err).Error("Failed to unmarshal notification")
return
}
bundleItem := model.BundleItem{
InteractionID: notification.InteractionID,
}
if notification.Transaction != nil || notification.DataItem != nil {
if notification.Transaction != nil {
bundleItem.Transaction = *notification.Transaction
}
if notification.Tags != nil {
bundleItem.Tags = *notification.Tags
}
if notification.DataItem != nil {
err = bundleItem.DataItem.Scan(*notification.DataItem)
if err != nil {
self.Log.WithError(err).Error("Failed to scan data item from notification")
return
}
}
} else {
// Transaction was too big to fit into the notification channel
// Only id is there, we need to fetch the rest of the data from the database
err = self.db.WithContext(self.Ctx).
Model(&model.BundleItem{}).
Select("transaction", "tags", "data_item").
Where("interaction_id = ?", notification.InteractionID).
Scan(&bundleItem).
Error
if err != nil {
// Action will be retried automatically, no need to do it here
self.Log.WithError(err).Error("Failed to get bundle item")
self.monitor.GetReport().Bundler.Errors.AdditionalFetchError.Inc()
return
}
self.monitor.GetReport().Bundler.State.AdditionalFetches.Inc()
}
select {
case <-self.StopChannel:
return
case self.output <- &bundleItem:
}
// Update metrics
self.monitor.GetReport().Bundler.State.BundlesFromNotifications.Inc()
// This might be the workload that unpauses the streamer
if self.GetWorkerQueueFillFactor() < 0.1 {
err := self.streamer.Resume()
if err != nil {
self.Log.WithError(err).Error("Failed to resume streamer")
}
}
})
// Pause streamer if the queue is too full or resume it
if self.GetWorkerQueueFillFactor() > 0.9 {
err := self.streamer.Pause()
if err != nil {
self.Log.WithError(err).Error("Failed to pause streamer")
}
}
}
}
}