-
Notifications
You must be signed in to change notification settings - Fork 416
/
Copy pathsubscriber.go
53 lines (45 loc) · 1.26 KB
/
subscriber.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
package metrics
import (
"github.com/ThreeDotsLabs/watermill/message"
"github.com/prometheus/client_golang/prometheus"
)
var (
subscriberLabelKeys = []string{
labelKeyHandlerName,
labelKeySubscriberName,
}
)
// SubscriberPrometheusMetricsDecorator decorates a subscriber to capture Prometheus metrics.
type SubscriberPrometheusMetricsDecorator struct {
message.Subscriber
subscriberName string
subscriberMessagesReceivedTotal *prometheus.CounterVec
closing chan struct{}
}
func (s SubscriberPrometheusMetricsDecorator) recordMetrics(msg *message.Message) {
if msg == nil {
return
}
ctx := msg.Context()
labels := labelsFromCtx(ctx, subscriberLabelKeys...)
if labels[labelKeySubscriberName] == "" {
labels[labelKeySubscriberName] = s.subscriberName
}
if labels[labelKeyHandlerName] == "" {
labels[labelKeyHandlerName] = labelValueNoHandler
}
go func() {
if subscribeAlreadyObserved(ctx) {
// decorator idempotency when applied decorator multiple times
return
}
select {
case <-msg.Acked():
labels[labelAcked] = "acked"
case <-msg.Nacked():
labels[labelAcked] = "nacked"
}
s.subscriberMessagesReceivedTotal.With(labels).Inc()
}()
msg.SetContext(setSubscribeObservedToCtx(msg.Context()))
}