-
Notifications
You must be signed in to change notification settings - Fork 3
/
opentracing.go
128 lines (104 loc) · 3.6 KB
/
opentracing.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
package opentracing
import (
"context"
"fmt"
"github.com/alexfalkowski/go-service/meta"
"github.com/alexfalkowski/go-service/time"
"github.com/alexfalkowski/go-service/transport/nsq/handler"
"github.com/alexfalkowski/go-service/transport/nsq/message"
"github.com/alexfalkowski/go-service/transport/nsq/producer"
"github.com/opentracing/opentracing-go"
"github.com/opentracing/opentracing-go/ext"
"github.com/opentracing/opentracing-go/log"
)
const (
nsqID = "nsq.id"
nsqBody = "nsq.body"
nsqTimestamp = "nsq.timestamp"
nsqAttempts = "nsq.attempts"
nsqAddress = "nsq.address"
nsqDuration = "nsq.duration_ms"
nsqStartTime = "nsq.start_time"
nsqTopic = "nsq.topic"
nsqChannel = "nsq.channel"
component = "component"
nsqComponent = "nsq"
)
// NewHandler for opentracing.
func NewHandler(topic, channel string, h handler.Handler) *Handler {
return &Handler{topic: topic, channel: channel, Handler: h}
}
// Handler for opentracing.
type Handler struct {
topic, channel string
handler.Handler
}
func (h *Handler) Handle(ctx context.Context, message *message.Message) error {
start := time.Now().UTC()
tracer := opentracing.GlobalTracer()
traceCtx, _ := tracer.Extract(opentracing.TextMap, headersTextMap(message.Headers))
operationName := fmt.Sprintf("consume %s:%s", h.topic, h.channel)
opts := []opentracing.StartSpanOption{
ext.RPCServerOption(traceCtx),
opentracing.Tag{Key: nsqStartTime, Value: start.Format(time.RFC3339)},
opentracing.Tag{Key: nsqTopic, Value: h.topic},
opentracing.Tag{Key: nsqChannel, Value: h.channel},
opentracing.Tag{Key: nsqID, Value: string(message.ID[:])},
opentracing.Tag{Key: nsqBody, Value: string(message.Body)},
opentracing.Tag{Key: nsqTimestamp, Value: message.Timestamp},
opentracing.Tag{Key: nsqAttempts, Value: message.Attempts},
opentracing.Tag{Key: nsqAddress, Value: message.NSQDAddress},
opentracing.Tag{Key: component, Value: nsqComponent},
ext.SpanKindConsumer,
}
span := tracer.StartSpan(operationName, opts...)
defer span.Finish()
ctx = opentracing.ContextWithSpan(ctx, span)
err := h.Handler.Handle(ctx, message)
for k, v := range meta.Attributes(ctx) {
span.SetTag(k, v)
}
span.SetTag(nsqDuration, time.ToMilliseconds(time.Since(start)))
if err != nil {
setError(span, err)
}
return err
}
// NewProducer for opentracing.
func NewProducer(p producer.Producer) *Producer {
return &Producer{Producer: p}
}
// Producer for opentracing.
type Producer struct {
producer.Producer
}
func (p *Producer) Publish(ctx context.Context, topic string, message *message.Message) error {
start := time.Now().UTC()
tracer := opentracing.GlobalTracer()
operationName := fmt.Sprintf("publish %s", topic)
opts := []opentracing.StartSpanOption{
opentracing.Tag{Key: nsqStartTime, Value: start.Format(time.RFC3339)},
opentracing.Tag{Key: nsqBody, Value: string(message.Body)},
opentracing.Tag{Key: nsqTopic, Value: topic},
opentracing.Tag{Key: component, Value: nsqComponent},
ext.SpanKindProducer,
}
span, ctx := opentracing.StartSpanFromContextWithTracer(ctx, tracer, operationName, opts...)
defer span.Finish()
if err := tracer.Inject(span.Context(), opentracing.TextMap, headersTextMap(message.Headers)); err != nil {
return err
}
err := p.Producer.Publish(ctx, topic, message)
for k, v := range meta.Attributes(ctx) {
span.SetTag(k, v)
}
span.SetTag(nsqDuration, time.ToMilliseconds(time.Since(start)))
if err != nil {
setError(span, err)
}
return err
}
func setError(span opentracing.Span, err error) {
ext.Error.Set(span, true)
span.LogFields(log.String("event", "error"), log.String("message", err.Error()))
}