Skip to content

Commit 7f44037

Browse files
committed
feat: adding support for grpc oltp
1 parent 0900dc5 commit 7f44037

6 files changed

Lines changed: 272 additions & 2 deletions

File tree

execution/grpc/client.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ func newHTTP2Client() *http.Client {
5050
// - *Client: The initialized gRPC client
5151
func NewClient(url string, opts ...connect.ClientOption) *Client {
5252
// Prepend WithGRPC to use the native gRPC protocol (required for tonic/gRPC servers)
53-
opts = append([]connect.ClientOption{connect.WithGRPC()}, opts...)
53+
opts = append([]connect.ClientOption{connect.WithInterceptors(outboundPropagationInterceptor()), connect.WithGRPC()}, opts...)
5454
return &Client{
5555
client: v1connect.NewExecutorServiceClient(
5656
newHTTP2Client(),

execution/grpc/go.mod

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,20 @@ require (
77
connectrpc.com/grpcreflect v1.3.0
88
github.com/evstack/ev-node v1.1.1
99
github.com/evstack/ev-node/core v1.0.0
10+
go.opentelemetry.io/otel v1.43.0
11+
go.opentelemetry.io/otel/sdk v1.43.0
12+
go.opentelemetry.io/otel/trace v1.43.0
1013
golang.org/x/net v0.53.0
1114
google.golang.org/protobuf v1.36.11
1215
)
1316

14-
require golang.org/x/text v0.36.0 // indirect
17+
require (
18+
github.com/cespare/xxhash/v2 v2.3.0 // indirect
19+
github.com/go-logr/logr v1.4.3 // indirect
20+
github.com/go-logr/stdr v1.2.2 // indirect
21+
github.com/google/uuid v1.6.0 // indirect
22+
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
23+
go.opentelemetry.io/otel/metric v1.43.0 // indirect
24+
golang.org/x/sys v0.43.0 // indirect
25+
golang.org/x/text v0.36.0 // indirect
26+
)

execution/grpc/go.sum

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,14 +2,35 @@ connectrpc.com/connect v1.19.2 h1:McQ83FGdzL+t60peksi0gXC7MQ/iLKgLduAnThbM0mo=
22
connectrpc.com/connect v1.19.2/go.mod h1:tN20fjdGlewnSFeZxLKb0xwIZ6ozc3OQs2hTXy4du9w=
33
connectrpc.com/grpcreflect v1.3.0 h1:Y4V+ACf8/vOb1XOc251Qun7jMB75gCUNw6llvB9csXc=
44
connectrpc.com/grpcreflect v1.3.0/go.mod h1:nfloOtCS8VUQOQ1+GTdFzVg2CJo4ZGaat8JIovCtDYs=
5+
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
6+
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
57
github.com/evstack/ev-node v1.1.1 h1:J9h5PKx177XdvNWLZCDOkWJEGRIrPzYxkCFhbGkVUm8=
68
github.com/evstack/ev-node v1.1.1/go.mod h1:/d/i+SSTDFnxffoijcrwmlt0LgfUU8D4S3HQqucwtu8=
79
github.com/evstack/ev-node/core v1.0.0 h1:s0Tx0uWHme7SJn/ZNEtee4qNM8UO6PIxXnHhPbbKTz8=
810
github.com/evstack/ev-node/core v1.0.0/go.mod h1:n2w/LhYQTPsi48m6lMj16YiIqsaQw6gxwjyJvR+B3sY=
11+
github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A=
12+
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
13+
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
14+
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
15+
github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE=
916
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
1017
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
18+
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
19+
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
20+
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
21+
go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y=
22+
go.opentelemetry.io/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I=
23+
go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0=
24+
go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM=
25+
go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY=
26+
go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg=
27+
go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg=
28+
go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A=
29+
go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0=
1130
golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA=
1231
golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs=
32+
golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI=
33+
golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
1334
golang.org/x/text v0.36.0 h1:JfKh3XmcRPqZPKevfXVpI1wXPTqbkE5f7JA92a55Yxg=
1435
golang.org/x/text v0.36.0/go.mod h1:NIdBknypM8iqVmPiuco0Dh6P5Jcdk8lJL0CUebqK164=
1536
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=

execution/grpc/handler.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
// - http.Handler: The configured HTTP handler
2626
func NewExecutorServiceHandler(executor execution.Executor, opts ...connect.HandlerOption) http.Handler {
2727
server := NewServer(executor)
28+
opts = append([]connect.HandlerOption{connect.WithInterceptors(inboundPropagationInterceptor())}, opts...)
2829

2930
mux := http.NewServeMux()
3031

execution/grpc/otel_propagation.go

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
package grpc
2+
3+
import (
4+
"context"
5+
6+
"connectrpc.com/connect"
7+
"go.opentelemetry.io/otel"
8+
"go.opentelemetry.io/otel/propagation"
9+
)
10+
11+
func inboundPropagationInterceptor() connect.UnaryInterceptorFunc {
12+
return connect.UnaryInterceptorFunc(func(next connect.UnaryFunc) connect.UnaryFunc {
13+
return func(ctx context.Context, req connect.AnyRequest) (connect.AnyResponse, error) {
14+
prop := otel.GetTextMapPropagator()
15+
ctx = prop.Extract(ctx, propagation.HeaderCarrier(req.Header()))
16+
return next(ctx, req)
17+
}
18+
})
19+
}
20+
21+
func outboundPropagationInterceptor() connect.UnaryInterceptorFunc {
22+
return connect.UnaryInterceptorFunc(func(next connect.UnaryFunc) connect.UnaryFunc {
23+
return func(ctx context.Context, req connect.AnyRequest) (connect.AnyResponse, error) {
24+
prop := otel.GetTextMapPropagator()
25+
prop.Inject(ctx, propagation.HeaderCarrier(req.Header()))
26+
return next(ctx, req)
27+
}
28+
})
29+
}
Lines changed: 207 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,207 @@
1+
package grpc
2+
3+
import (
4+
"context"
5+
"net/http/httptest"
6+
"testing"
7+
"time"
8+
9+
"connectrpc.com/connect"
10+
"go.opentelemetry.io/otel"
11+
"go.opentelemetry.io/otel/baggage"
12+
"go.opentelemetry.io/otel/propagation"
13+
sdktrace "go.opentelemetry.io/otel/sdk/trace"
14+
"go.opentelemetry.io/otel/sdk/trace/tracetest"
15+
"go.opentelemetry.io/otel/trace"
16+
17+
"github.com/evstack/ev-node/core/execution"
18+
)
19+
20+
func setupTracer(t *testing.T) (*tracetest.SpanRecorder, func()) {
21+
t.Helper()
22+
rec := tracetest.NewSpanRecorder()
23+
tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(rec))
24+
oldTP := otel.GetTracerProvider()
25+
oldProp := otel.GetTextMapPropagator()
26+
otel.SetTracerProvider(tp)
27+
otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator(propagation.TraceContext{}, propagation.Baggage{}))
28+
return rec, func() {
29+
_ = tp.Shutdown(context.Background())
30+
otel.SetTracerProvider(oldTP)
31+
otel.SetTextMapPropagator(oldProp)
32+
}
33+
}
34+
35+
func TestInboundMetadataCreatesChildSpanWithSameTraceID(t *testing.T) {
36+
rec, cleanup := setupTracer(t)
37+
defer cleanup()
38+
39+
tracer := otel.Tracer("test")
40+
parentCtx, parent := tracer.Start(context.Background(), "parent")
41+
defer parent.End()
42+
parentTraceID := parent.SpanContext().TraceID()
43+
44+
mockExec := &mockExecutor{getTxsFunc: func(ctx context.Context) ([][]byte, error) {
45+
_, span := tracer.Start(ctx, "server-child")
46+
span.End()
47+
return [][]byte{}, nil
48+
}}
49+
50+
handler := NewExecutorServiceHandler(mockExec)
51+
ts := httptest.NewServer(handler)
52+
defer ts.Close()
53+
54+
client := NewClient(ts.URL)
55+
_, err := client.GetTxs(parentCtx)
56+
if err != nil {
57+
t.Fatalf("GetTxs failed: %v", err)
58+
}
59+
60+
var found bool
61+
for _, s := range rec.Ended() {
62+
if s.Name() == "server-child" {
63+
found = true
64+
if s.SpanContext().TraceID() != parentTraceID {
65+
t.Fatalf("trace id mismatch: got %s want %s", s.SpanContext().TraceID(), parentTraceID)
66+
}
67+
}
68+
}
69+
if !found {
70+
t.Fatalf("server-child span not found")
71+
}
72+
}
73+
74+
func TestOutboundGRPCCallCarriesTraceparentMetadata(t *testing.T) {
75+
rec, cleanup := setupTracer(t)
76+
_ = rec
77+
defer cleanup()
78+
79+
tracer := otel.Tracer("test")
80+
ctx, parent := tracer.Start(context.Background(), "parent")
81+
defer parent.End()
82+
83+
gotTraceparent := ""
84+
captureHeader := connect.UnaryInterceptorFunc(func(next connect.UnaryFunc) connect.UnaryFunc {
85+
return func(ctx context.Context, req connect.AnyRequest) (connect.AnyResponse, error) {
86+
gotTraceparent = req.Header().Get("traceparent")
87+
return next(ctx, req)
88+
}
89+
})
90+
91+
mockExec := &mockExecutor{}
92+
handler := NewExecutorServiceHandler(mockExec, connect.WithInterceptors(captureHeader))
93+
ts := httptest.NewServer(handler)
94+
defer ts.Close()
95+
96+
client := NewClient(ts.URL)
97+
if _, err := client.GetTxs(ctx); err != nil {
98+
t.Fatalf("GetTxs failed: %v", err)
99+
}
100+
if gotTraceparent == "" {
101+
t.Fatalf("expected traceparent metadata to be propagated")
102+
}
103+
}
104+
105+
func TestOutboundGRPCCallCarriesPropagationHeaders(t *testing.T) {
106+
rec, cleanup := setupTracer(t)
107+
_ = rec
108+
defer cleanup()
109+
110+
tracer := otel.Tracer("test")
111+
ctx, parent := tracer.Start(context.Background(), "parent")
112+
defer parent.End()
113+
member, err := baggage.NewMember("tenant", "alpha")
114+
if err != nil {
115+
t.Fatalf("failed to create baggage member: %v", err)
116+
}
117+
bg, err := baggage.New(member)
118+
if err != nil {
119+
t.Fatalf("failed to create baggage: %v", err)
120+
}
121+
ctx = baggage.ContextWithBaggage(ctx, bg)
122+
123+
var gotTraceparent string
124+
var gotBaggage string
125+
captureHeader := connect.UnaryInterceptorFunc(func(next connect.UnaryFunc) connect.UnaryFunc {
126+
return func(ctx context.Context, req connect.AnyRequest) (connect.AnyResponse, error) {
127+
gotTraceparent = req.Header().Get("traceparent")
128+
gotBaggage = req.Header().Get("baggage")
129+
return next(ctx, req)
130+
}
131+
})
132+
133+
mockExec := &mockExecutor{}
134+
handler := NewExecutorServiceHandler(mockExec, connect.WithInterceptors(captureHeader))
135+
ts := httptest.NewServer(handler)
136+
defer ts.Close()
137+
138+
client := NewClient(ts.URL)
139+
if _, err := client.GetTxs(ctx); err != nil {
140+
t.Fatalf("GetTxs failed: %v", err)
141+
}
142+
143+
if gotTraceparent == "" {
144+
t.Fatalf("expected traceparent metadata to be propagated")
145+
}
146+
if gotBaggage == "" {
147+
t.Fatalf("expected baggage metadata to be propagated")
148+
}
149+
}
150+
151+
func TestEndToEndParentChildAcrossServerClientHop(t *testing.T) {
152+
rec, cleanup := setupTracer(t)
153+
defer cleanup()
154+
155+
tracer := otel.Tracer("test")
156+
var midSpan trace.Span
157+
158+
downstreamExec := &mockExecutor{getExecutionInfoFunc: func(ctx context.Context) (executionInfo execution.ExecutionInfo, err error) {
159+
_, span := tracer.Start(ctx, "downstream-child")
160+
span.End()
161+
return execution.ExecutionInfo{MaxGas: 1}, nil
162+
}}
163+
downstreamHandler := NewExecutorServiceHandler(downstreamExec)
164+
downstreamSrv := httptest.NewServer(downstreamHandler)
165+
defer downstreamSrv.Close()
166+
downstreamClient := NewClient(downstreamSrv.URL)
167+
168+
upstreamExec := &mockExecutor{getTxsFunc: func(ctx context.Context) ([][]byte, error) {
169+
ctx, span := tracer.Start(ctx, "upstream-mid")
170+
midSpan = span
171+
defer span.End()
172+
_, err := downstreamClient.GetExecutionInfo(ctx)
173+
if err != nil {
174+
return nil, err
175+
}
176+
return [][]byte{}, nil
177+
}}
178+
upstreamHandler := NewExecutorServiceHandler(upstreamExec)
179+
upstreamSrv := httptest.NewServer(upstreamHandler)
180+
defer upstreamSrv.Close()
181+
182+
client := NewClient(upstreamSrv.URL)
183+
rootCtx, root := tracer.Start(context.Background(), "root")
184+
defer root.End()
185+
if _, err := client.GetTxs(rootCtx); err != nil {
186+
t.Fatalf("GetTxs failed: %v", err)
187+
}
188+
189+
time.Sleep(10 * time.Millisecond)
190+
191+
rootTraceID := root.SpanContext().TraceID()
192+
if midSpan.SpanContext().TraceID() != rootTraceID {
193+
t.Fatalf("mid span trace id mismatch")
194+
}
195+
var found bool
196+
for _, s := range rec.Ended() {
197+
if s.Name() == "downstream-child" {
198+
found = true
199+
if s.SpanContext().TraceID() != rootTraceID {
200+
t.Fatalf("downstream trace id mismatch")
201+
}
202+
}
203+
}
204+
if !found {
205+
t.Fatalf("downstream-child span not found")
206+
}
207+
}

0 commit comments

Comments
 (0)