Skip to content
This repository was archived by the owner on Feb 12, 2024. It is now read-only.

Commit 2ac2f2e

Browse files
andbosonAndrey Kolomiets
authored andcommitted
added decoder error process (#38)
* added decoder error process
1 parent d3cb480 commit 2ac2f2e

5 files changed

Lines changed: 167 additions & 41 deletions

File tree

CHANGELOG.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,10 @@
11
# Changelog
22

3+
4+
## [2.3.0] - 2019-02-08
5+
### Change
6+
- fixed opentracing span creation in amqp_kit
7+
38
## [2.2.0]
49
### Add
510
- add tracing database

amqp-kit/subscriber.go

Lines changed: 1 addition & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,6 @@ import (
66

77
"github.com/go-kit/kit/endpoint"
88
"github.com/opentracing-contrib/go-amqp/amqptracer"
9-
"github.com/opentracing/opentracing-go"
109
"github.com/streadway/amqp"
1110
)
1211

@@ -92,14 +91,7 @@ func (s Subscriber) ServeDelivery(ch Channel) func(deliv *amqp.Delivery) {
9291

9392
//extract tracing headers and start root span
9493
spCtx, _ := amqptracer.Extract(deliv.Headers)
95-
sp := opentracing.StartSpan(
96-
"ConsumeMessage",
97-
opentracing.FollowsFrom(spCtx),
98-
)
99-
defer sp.Finish()
100-
101-
// Update the context with the span for the subsequent reference.
102-
ctx = opentracing.ContextWithSpan(ctx, sp)
94+
ctx = context.WithValue(ctx, amqpCtx, &spCtx)
10395

10496
request, err := s.dec(ctx, deliv)
10597

amqp-kit/subscriber_test.go

Lines changed: 140 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,140 @@
1+
package amqp_kit
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"testing"
7+
"time"
8+
9+
"github.com/go-kit/kit/endpoint"
10+
"github.com/opentracing/opentracing-go"
11+
"github.com/opentracing/opentracing-go/mocktracer"
12+
"github.com/streadway/amqp"
13+
"github.com/stretchr/testify/suite"
14+
)
15+
16+
type subsSuite struct {
17+
suite.Suite
18+
dsn string
19+
conn *amqp.Connection
20+
}
21+
22+
func (s *subsSuite) SetupSuite() {
23+
var err error
24+
s.dsn = MakeDsn(&Config{
25+
Address: "127.0.0.1:5672",
26+
User: "guest",
27+
Password: "guest"},
28+
)
29+
30+
s.conn, err = amqp.Dial(s.dsn)
31+
s.Require().NoError(err)
32+
}
33+
34+
func (s *subsSuite) TearDownSuite() {
35+
s.conn.Close()
36+
}
37+
38+
func TestSubscriberSuite(t *testing.T) {
39+
suite.Run(t, new(subsSuite))
40+
}
41+
42+
func (s *subsSuite) TestSubscriberTracing() {
43+
tracer := mocktracer.New()
44+
opentracing.SetGlobalTracer(tracer)
45+
46+
ch, err := s.conn.Channel()
47+
s.NoError(err)
48+
49+
err = Declare(ch, `exc`, `test-s`, []string{`key.request.test_2`})
50+
s.NoError(err)
51+
ctx := context.Background()
52+
dec1 := make(chan []byte)
53+
subs := []SubscribeInfo{
54+
{
55+
Q: `test-s`,
56+
Key: `key.request.test_2`,
57+
E: endpoint.Chain(
58+
TraceEndpoint(tracer, `test_endpoint`),
59+
)(func(ctx context.Context, request interface{}) (response interface{}, err error) {
60+
dec1 <- request.([]byte)
61+
62+
if string(request.([]byte)) == "endpoint_error" {
63+
return nil, fmt.Errorf("endpoint_error_resp")
64+
}
65+
66+
return nil, nil
67+
}),
68+
Dec: func(i context.Context, delivery *amqp.Delivery) (request interface{}, err error) {
69+
s.Equal(delivery.RoutingKey, `key.request.test_2`)
70+
71+
if delivery.CorrelationId == "errDecodeID" {
72+
return nil, fmt.Errorf("decode error")
73+
}
74+
75+
if delivery.CorrelationId == "errEndpointID" {
76+
return []byte("endpoint_error"), nil
77+
}
78+
return delivery.Body, nil
79+
},
80+
Enc: EncodeJSONResponse,
81+
O: []SubscriberOption{SubscriberAfter(SetAckAfterEndpoint(true))},
82+
},
83+
}
84+
85+
ser := NewServer(subs, s.conn)
86+
err = ser.Serve()
87+
s.NoError(err)
88+
89+
ch, err = s.conn.Channel()
90+
s.NoError(err)
91+
pub := NewPublisher(ch)
92+
93+
// success
94+
err = pub.PublishWithTracing(ctx, "exc", "key.request.test_2", `cor_1`, []byte(`{"f1":"b1"}`))
95+
s.NoError(err)
96+
97+
select {
98+
case d := <-dec1:
99+
s.Equal(d, []byte(`{"f1":"b1"}`))
100+
case <-time.After(5 * time.Second):
101+
}
102+
103+
finishedSpans := opentracing.GlobalTracer().(*mocktracer.MockTracer).FinishedSpans()
104+
s.Len(finishedSpans, 2)
105+
s.Equal(finishedSpans[0].SpanContext.TraceID, finishedSpans[1].SpanContext.TraceID)
106+
107+
opentracing.GlobalTracer().(*mocktracer.MockTracer).Reset()
108+
109+
// decode error
110+
err = pub.PublishWithTracing(ctx, "exc", "key.request.test_2", `errDecodeID`, []byte(`{"f1":"b1"}`))
111+
s.NoError(err)
112+
113+
select {
114+
case d := <-dec1:
115+
s.Equal(d, []byte(`{"f1":"b1"}`))
116+
case <-time.After(5 * time.Second):
117+
}
118+
119+
finishedSpans = opentracing.GlobalTracer().(*mocktracer.MockTracer).FinishedSpans()
120+
s.Len(finishedSpans, 1)
121+
122+
// endpoint error
123+
opentracing.GlobalTracer().(*mocktracer.MockTracer).Reset()
124+
125+
err = pub.PublishWithTracing(ctx, "exc", "key.request.test_2", `errEndpointID`, []byte(`{"f1":"b1"}`))
126+
s.NoError(err)
127+
128+
select {
129+
case d := <-dec1:
130+
s.Equal(d, []byte(`endpoint_error`))
131+
case <-time.After(5 * time.Second):
132+
}
133+
134+
finishedSpans = opentracing.GlobalTracer().(*mocktracer.MockTracer).FinishedSpans()
135+
s.Len(finishedSpans, 2)
136+
s.Equal(finishedSpans[1].Tag(tagError), "endpoint_error_resp")
137+
138+
err = ser.Stop()
139+
s.Require().NoError(err)
140+
}

amqp-kit/tracing.go

Lines changed: 19 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -4,47 +4,36 @@ import (
44
"context"
55

66
"github.com/go-kit/kit/endpoint"
7-
"github.com/opentracing-contrib/go-amqp/amqptracer"
87
"github.com/opentracing/opentracing-go"
98
otext "github.com/opentracing/opentracing-go/ext"
10-
"github.com/streadway/amqp"
119
)
1210

13-
func DecodeWithTrace(next DecodeRequestFunc, operationName string) DecodeRequestFunc {
14-
return func(ctx context.Context, r *amqp.Delivery) (i interface{}, err error) {
15-
//extract tracing headers
16-
spCtx, _ := amqptracer.Extract(r.Headers)
17-
sp := opentracing.StartSpan(
18-
operationName,
19-
opentracing.FollowsFrom(spCtx),
20-
)
21-
defer sp.Finish()
22-
23-
// Update the context with the span for the subsequent reference.
24-
otext.SpanKindRPCServer.Set(sp)
25-
ctx = opentracing.ContextWithSpan(ctx, sp)
26-
27-
return next(ctx, r)
28-
}
29-
}
11+
const tagError = "error"
12+
13+
type amqpSpanCtx string
3014

31-
// set operation name for parent span, started in subscriber.go
32-
// if span not found, start new span with given tracer
15+
var amqpCtx amqpSpanCtx
16+
17+
// Get context value spanContext and start Span with given operationName.
18+
// Set an error as tag if raised.
3319
func TraceEndpoint(tracer opentracing.Tracer, operationName string) endpoint.Middleware {
3420
return func(next endpoint.Endpoint) endpoint.Endpoint {
3521
return func(ctx context.Context, request interface{}) (interface{}, error) {
36-
parentSpan := opentracing.SpanFromContext(ctx)
37-
if parentSpan != nil {
38-
parentSpan.SetOperationName(operationName)
39-
} else {
40-
parentSpan = tracer.StartSpan(operationName)
41-
defer parentSpan.Finish()
22+
var sp opentracing.Span
23+
if spCtx, ok := ctx.Value(amqpCtx).(*opentracing.SpanContext); ok {
24+
sp = tracer.StartSpan(operationName, opentracing.FollowsFrom(*spCtx))
25+
defer sp.Finish()
26+
27+
otext.SpanKindRPCServer.Set(sp)
28+
ctx = opentracing.ContextWithSpan(ctx, sp)
4229
}
4330

44-
otext.SpanKindRPCServer.Set(parentSpan)
45-
ctx = opentracing.ContextWithSpan(ctx, parentSpan)
31+
i, err := next(ctx, request)
32+
if err != nil && sp != nil {
33+
sp.SetTag(tagError, err.Error())
34+
}
4635

47-
return next(ctx, request)
36+
return i, err
4837
}
4938
}
5039
}

go.sum

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,8 +28,8 @@ github.com/RoaringBitmap/roaring v0.4.16/go.mod h1:8khRDP4HmeXns4xIj9oGrKSz7XTQi
2828
github.com/SAP/go-hdb v0.13.1/go.mod h1:etBT+FAi1t5k3K3tf5vQTnosgYmhDkRi8jEnQqCnxF0=
2929
github.com/SAP/go-hdb v0.13.2 h1:19+oZb2RGKCJpSZjFN/dcdOuSACSWGJlinHwgVyvUnM=
3030
github.com/SAP/go-hdb v0.13.2/go.mod h1:etBT+FAi1t5k3K3tf5vQTnosgYmhDkRi8jEnQqCnxF0=
31-
github.com/SermoDigital/jose v0.9.1 h1:atYaHPD3lPICcbK1owly3aPm0iaJGSGPi0WD4vLznv8=
32-
github.com/SermoDigital/jose v0.9.1/go.mod h1:ARgCUhI1MHQH+ONky/PAtmVHQrP5JlGY0F3poXOp/fA=
31+
github.com/SermoDigital/jose v0.9.2-0.20161205224733-f6df55f235c2 h1:koK7z0nSsRiRiBWwa+E714Puh+DO+ZRdIyAXiXzL+lg=
32+
github.com/SermoDigital/jose v0.9.2-0.20161205224733-f6df55f235c2/go.mod h1:ARgCUhI1MHQH+ONky/PAtmVHQrP5JlGY0F3poXOp/fA=
3333
github.com/VividCortex/gohistogram v1.0.0 h1:6+hBz+qvs0JOrrNhhmR7lFxo5sINxBCGXrdtl/UvroE=
3434
github.com/VividCortex/gohistogram v1.0.0/go.mod h1:Pf5mBqqDxYaXu3hDrrU+w6nw50o/4+TcAqDqk/vUH7g=
3535
github.com/alcortesm/tgz v0.0.0-20161220082320-9c5fe88206d7/go.mod h1:6zEj6s6u/ghQa61ZWa/C2Aw3RkjiTBOix7dkqa1VLIs=

0 commit comments

Comments
 (0)