Skip to content

Commit 5260332

Browse files
committed
feat: migrate eventing-natss to otel
Signed-off-by: Calum Murray <cmurray@redhat.com>
1 parent 93bae99 commit 5260332

6 files changed

Lines changed: 189 additions & 89 deletions

File tree

pkg/channel/jetstream/dispatcher/dispatcher.go

Lines changed: 20 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -25,13 +25,13 @@ import (
2525
"sync"
2626

2727
"github.com/nats-io/nats.go"
28+
"go.opentelemetry.io/otel"
29+
"go.opentelemetry.io/otel/metric"
2830

2931
cejs "github.com/cloudevents/sdk-go/protocol/nats_jetstream/v2"
3032
ce "github.com/cloudevents/sdk-go/v2"
3133
"github.com/cloudevents/sdk-go/v2/binding"
3234
"github.com/cloudevents/sdk-go/v2/event"
33-
"github.com/google/uuid"
34-
"go.opencensus.io/trace"
3535
"go.uber.org/zap"
3636
"k8s.io/apimachinery/pkg/types"
3737
"k8s.io/apimachinery/pkg/util/sets"
@@ -43,7 +43,6 @@ import (
4343
eventingchannels "knative.dev/eventing/pkg/channel"
4444
"knative.dev/eventing/pkg/eventingtls"
4545
"knative.dev/eventing/pkg/kncloudevents"
46-
"knative.dev/pkg/kmeta"
4746
"knative.dev/pkg/logging"
4847
)
4948

@@ -55,7 +54,6 @@ import (
5554
type Dispatcher struct {
5655
receiver *eventingchannels.EventReceiver
5756
dispatcher *kncloudevents.Dispatcher
58-
reporter eventingchannels.StatsReporter
5957

6058
js nats.JetStreamContext
6159

@@ -74,12 +72,9 @@ type Dispatcher struct {
7472
func NewDispatcher(ctx context.Context, args NatsDispatcherArgs) (*Dispatcher, error) {
7573
logger := logging.FromContext(ctx)
7674

77-
reporter := eventingchannels.NewStatsReporter(args.ContainerName, kmeta.ChildName(args.PodName, uuid.New().String()))
78-
7975
oidcTokenProvider := auth.NewOIDCTokenProvider(ctx)
8076
d := &Dispatcher{
8177
dispatcher: kncloudevents.NewDispatcher(eventingtls.ClientConfig{}, oidcTokenProvider),
82-
reporter: reporter,
8378

8479
js: args.JetStream,
8580

@@ -94,7 +89,6 @@ func NewDispatcher(ctx context.Context, args NatsDispatcherArgs) (*Dispatcher, e
9489
receiverFunc, err := eventingchannels.NewEventReceiver(
9590
d.messageReceiver,
9691
logger.Desugar(),
97-
reporter,
9892
eventingchannels.ResolveChannelFromHostHeader(d.getChannelReferenceFromHost),
9993
)
10094
if err != nil {
@@ -274,7 +268,6 @@ func (d *Dispatcher) subscribe(ctx context.Context, config ChannelConfig, sub Su
274268
pushConsumer := &PushConsumer{
275269
sub: sub,
276270
dispatcher: d.dispatcher,
277-
reporter: d.reporter,
278271
channelNamespace: config.Namespace,
279272
logger: logger,
280273
ctx: ctx,
@@ -287,14 +280,30 @@ func (d *Dispatcher) subscribe(ctx context.Context, config ChannelConfig, sub Su
287280
return SubscriberStatusTypeError, err
288281
}
289282
pushConsumer.jsSub = jsSub
283+
284+
mp := otel.GetMeterProvider()
285+
tp := otel.GetTracerProvider()
286+
pushConsumer.tracer = tp.Tracer(scopeName)
287+
288+
meter := mp.Meter(scopeName)
289+
pushConsumer.dispatchDuration, err = meter.Float64Histogram(
290+
"kn.eventing.dispatch.duration",
291+
metric.WithDescription("The duration to dispatch the event"),
292+
metric.WithUnit("s"),
293+
metric.WithExplicitBucketBoundaries(latencyBounds...),
294+
)
295+
if err != nil {
296+
logger.Errorw("failed to set up dispatch duration metric", zap.Error(err))
297+
return SubscriberStatusTypeError, err
298+
}
290299
consumer = pushConsumer
291300
} else {
292301
natsSub, err := d.js.PullSubscribe(".>", info.Config.Durable, nats.Bind(info.Stream, info.Name), nats.ManualAck())
293302
if err != nil {
294303
logger.Errorw("failed to pull subscribe to jetstream", zap.Error(err))
295304
return SubscriberStatusTypeError, err
296305
}
297-
consumer, err = NewPullConsumer(ctx, natsSub, sub, d.dispatcher, d.reporter, &config)
306+
consumer, err = NewPullConsumer(ctx, natsSub, sub, d.dispatcher, &config)
298307
if err != nil {
299308
logger.Errorw("failed to create pull consumer", zap.Error(err))
300309
return SubscriberStatusTypeError, err
@@ -383,7 +392,7 @@ func (d *Dispatcher) messageReceiver(ctx context.Context, ch eventingchannels.Ch
383392
eventID := commonce.IDExtractorTransformer("")
384393

385394
transformers := append([]binding.Transformer{&eventID},
386-
tracing.SerializeTraceTransformers(trace.FromContext(ctx).SpanContext())...,
395+
tracing.SerializeTraceTransformers(ctx)...,
387396
)
388397

389398
ctx = ce.WithEncodingStructured(ctx)

pkg/channel/jetstream/dispatcher/message_dispatcher_test.go

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,10 @@ import (
2929
"time"
3030

3131
"github.com/nats-io/nats.go"
32+
"go.opentelemetry.io/otel"
33+
"go.opentelemetry.io/otel/sdk/metric"
34+
"go.opentelemetry.io/otel/sdk/trace"
35+
"go.opentelemetry.io/otel/sdk/trace/tracetest"
3236

3337
"github.com/nats-io/nats-server/v2/server"
3438
natsserver "github.com/nats-io/nats-server/v2/test"
@@ -37,6 +41,7 @@ import (
3741
"knative.dev/eventing/pkg/eventingtls"
3842
"knative.dev/pkg/apis"
3943
duckv1 "knative.dev/pkg/apis/duck/v1"
44+
"knative.dev/pkg/observability/tracing"
4045

4146
v1 "knative.dev/eventing/pkg/apis/duck/v1"
4247
"knative.dev/eventing/pkg/kncloudevents"
@@ -450,6 +455,14 @@ func TestDispatchMessage(t *testing.T) {
450455

451456
for n, tc := range testCases {
452457
t.Run(n, func(t *testing.T) {
458+
reader := metric.NewManualReader()
459+
mp := metric.NewMeterProvider(metric.WithReader(reader))
460+
otel.SetMeterProvider(mp)
461+
462+
exporter := tracetest.NewInMemoryExporter()
463+
tp := trace.NewTracerProvider(trace.WithSyncer(exporter))
464+
otel.SetTracerProvider(tp)
465+
otel.SetTextMapPropagator(tracing.DefaultTextMapPropagator())
453466
destHandler := &fakeHandler{
454467
t: t,
455468
response: tc.fakeResponse,

pkg/channel/jetstream/dispatcher/pull_consumer.go

Lines changed: 63 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -29,14 +29,19 @@ import (
2929
cejs "github.com/cloudevents/sdk-go/protocol/nats_jetstream/v2"
3030
"github.com/cloudevents/sdk-go/v2/binding"
3131

32-
"go.opencensus.io/trace"
32+
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
33+
"go.opentelemetry.io/otel"
34+
"go.opentelemetry.io/otel/metric"
35+
"go.opentelemetry.io/otel/trace"
36+
"k8s.io/apimachinery/pkg/types"
37+
3338
"knative.dev/eventing-natss/pkg/tracing"
3439

3540
"github.com/nats-io/nats.go"
3641
"go.uber.org/zap"
37-
eventingchannels "knative.dev/eventing/pkg/channel"
38-
"knative.dev/eventing/pkg/channel/fanout"
3942
"knative.dev/eventing/pkg/kncloudevents"
43+
"knative.dev/eventing/pkg/observability"
44+
eventingtracing "knative.dev/eventing/pkg/tracing"
4045
"knative.dev/pkg/logging"
4146
)
4247

@@ -64,8 +69,10 @@ func init() {
6469

6570
type PullConsumer struct {
6671
dispatcher *kncloudevents.Dispatcher
67-
reporter eventingchannels.StatsReporter
6872
channelNamespace string
73+
channelName string
74+
dispatchDuration metric.Float64Histogram
75+
tracer trace.Tracer
6976

7077
natsConsumer *nats.Subscription
7178
natsConsumerInfo *nats.ConsumerInfo
@@ -81,7 +88,7 @@ type PullConsumer struct {
8188
closed chan struct{}
8289
}
8390

84-
func NewPullConsumer(ctx context.Context, consumer *nats.Subscription, subscription Subscription, dispatcher *kncloudevents.Dispatcher, reporter eventingchannels.StatsReporter, channelConfig *ChannelConfig) (*PullConsumer, error) {
91+
func NewPullConsumer(ctx context.Context, consumer *nats.Subscription, subscription Subscription, dispatcher *kncloudevents.Dispatcher, channelConfig *ChannelConfig) (*PullConsumer, error) {
8592
consumerInfo, err := consumer.ConsumerInfo()
8693
if err != nil {
8794
return nil, fmt.Errorf("failed to get consumer info: %w", err)
@@ -90,10 +97,14 @@ func NewPullConsumer(ctx context.Context, consumer *nats.Subscription, subscript
9097
logger := logging.FromContext(ctx)
9198
updatePullSubscriptionConfig(channelConfig, &subscription)
9299

93-
return &PullConsumer{
100+
tp := otel.GetTracerProvider()
101+
mp := otel.GetMeterProvider()
102+
103+
c := &PullConsumer{
94104
dispatcher: dispatcher,
95-
reporter: reporter,
96105
channelNamespace: channelConfig.Namespace,
106+
channelName: channelConfig.Name,
107+
tracer: tp.Tracer(scopeName),
97108

98109
natsConsumer: consumer,
99110
natsConsumerInfo: consumerInfo,
@@ -104,7 +115,20 @@ func NewPullConsumer(ctx context.Context, consumer *nats.Subscription, subscript
104115

105116
closing: make(chan struct{}),
106117
closed: make(chan struct{}),
107-
}, nil
118+
}
119+
120+
meter := mp.Meter(scopeName)
121+
c.dispatchDuration, err = meter.Float64Histogram(
122+
"kn.eventing.dispatch.duration",
123+
metric.WithDescription("The duration to dispatch the event"),
124+
metric.WithUnit("s"),
125+
metric.WithExplicitBucketBoundaries(latencyBounds...),
126+
)
127+
if err != nil {
128+
return nil, err
129+
}
130+
131+
return c, nil
108132
}
109133

110134
func (c *PullConsumer) ConsumerType() ConsumerType {
@@ -234,16 +258,30 @@ func (c *PullConsumer) handleMessage(ctx context.Context, msg *nats.Msg) (err er
234258
event := tracing.ConvertNatsMsgToEvent(c.logger.Desugar(), msg)
235259
additionalHeaders := tracing.ConvertEventToHttpHeader(event)
236260

237-
sc, ok := tracing.ParseSpanContext(event)
238-
var span *trace.Span
239-
if !ok {
240-
c.logger.Warn("Cannot parse the spancontext, creating a new span")
241-
ctx, span = trace.StartSpan(ctx, jsmChannel+"-"+string(c.sub.UID))
242-
} else {
243-
ctx, span = trace.StartSpanWithRemoteParent(ctx, jsmChannel+"-"+string(c.sub.UID), sc)
244-
}
245-
246-
defer span.End()
261+
ctx = observability.WithChannelLabels(ctx, types.NamespacedName{Name: c.channelName, Namespace: c.channelNamespace})
262+
ctx = observability.WithMessagingLabels(
263+
ctx,
264+
eventingtracing.SubscriptionMessagingDestination(
265+
types.NamespacedName{
266+
Name: c.sub.Name,
267+
Namespace: c.sub.Namespace,
268+
},
269+
),
270+
"send",
271+
)
272+
ctx = observability.WithMinimalEventLabels(ctx, event)
273+
274+
ctx = tracing.ParseSpanContext(ctx, event)
275+
ctx, span := c.tracer.Start(ctx, jsmChannel+"-"+string(c.sub.UID))
276+
defer func() {
277+
if span.IsRecording() {
278+
// add full event labels here so that they only populate the span, and not any metrics
279+
ctx = observability.WithEventLabels(ctx, event)
280+
labeler, _ := otelhttp.LabelerFromContext(ctx)
281+
span.SetAttributes(labeler.Get()...)
282+
}
283+
span.End()
284+
}()
247285

248286
te := TypeExtractorTransformer("")
249287

@@ -261,10 +299,13 @@ func (c *PullConsumer) handleMessage(ctx context.Context, msg *nats.Msg) (err er
261299
WithHeader(additionalHeaders),
262300
)
263301

264-
_ = fanout.ParseDispatchResultAndReportMetrics(fanout.NewDispatchResult(err, dispatchExecutionInfo), c.reporter, eventingchannels.ReportArgs{
265-
Ns: c.channelNamespace,
266-
EventType: string(te),
267-
})
302+
labeler, _ := otelhttp.LabelerFromContext(ctx)
303+
304+
c.dispatchDuration.Record(
305+
ctx,
306+
dispatchExecutionInfo.Duration.Seconds(),
307+
metric.WithAttributes(labeler.Get()...),
308+
)
268309

269310
if err != nil {
270311
logger.Errorw("failed to forward message to downstream subscriber",

pkg/channel/jetstream/dispatcher/push_consumer.go

Lines changed: 42 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -21,32 +21,39 @@ import (
2121
"errors"
2222
"sync"
2323

24+
"go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp"
25+
"go.opentelemetry.io/otel/metric"
26+
"go.opentelemetry.io/otel/trace"
27+
"k8s.io/apimachinery/pkg/types"
2428
"knative.dev/eventing/pkg/kncloudevents"
2529

2630
cejs "github.com/cloudevents/sdk-go/protocol/nats_jetstream/v2"
2731
"github.com/cloudevents/sdk-go/v2/binding"
2832
"github.com/nats-io/nats.go"
29-
"go.opencensus.io/trace"
3033
"go.uber.org/zap"
3134
"knative.dev/eventing-natss/pkg/tracing"
32-
eventingchannels "knative.dev/eventing/pkg/channel"
33-
"knative.dev/eventing/pkg/channel/fanout"
35+
"knative.dev/eventing/pkg/observability"
36+
eventingtracing "knative.dev/eventing/pkg/tracing"
3437
"knative.dev/pkg/logging"
3538
)
3639

3740
const (
3841
jsmChannel = "jsm-channel"
42+
scopeName = "knative.dev/eventing-natss/pkg/channel/jetstream/dispatcher"
3943
)
4044

4145
var (
4246
ErrConsumerClosed = errors.New("dispatcher consumer closed")
47+
latencyBounds = []float64{0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1, 2.5, 5, 7.5, 10}
4348
)
4449

4550
type PushConsumer struct {
4651
sub Subscription
4752
dispatcher *kncloudevents.Dispatcher
48-
reporter eventingchannels.StatsReporter
4953
channelNamespace string
54+
channelName string
55+
dispatchDuration metric.Float64Histogram
56+
tracer trace.Tracer
5057

5158
jsSub *nats.Subscription
5259
natsConsumerInfo *nats.ConsumerInfo
@@ -116,16 +123,30 @@ func (c *PushConsumer) doHandle(ctx context.Context, msg *nats.Msg) {
116123
event := tracing.ConvertNatsMsgToEvent(c.logger.Desugar(), msg)
117124
additionalHeaders := tracing.ConvertEventToHttpHeader(event)
118125

119-
sc, ok := tracing.ParseSpanContext(event)
120-
var span *trace.Span
121-
if !ok {
122-
logger.Warn("Cannot parse the spancontext, creating a new span")
123-
ctx, span = trace.StartSpan(ctx, jsmChannel+"-"+string(c.sub.UID))
124-
} else {
125-
ctx, span = trace.StartSpanWithRemoteParent(ctx, jsmChannel+"-"+string(c.sub.UID), sc)
126-
}
127-
128-
defer span.End()
126+
ctx = observability.WithChannelLabels(ctx, types.NamespacedName{Name: c.channelName, Namespace: c.channelNamespace})
127+
ctx = observability.WithMessagingLabels(
128+
ctx,
129+
eventingtracing.SubscriptionMessagingDestination(
130+
types.NamespacedName{
131+
Name: c.sub.Name,
132+
Namespace: c.sub.Namespace,
133+
},
134+
),
135+
"send",
136+
)
137+
ctx = observability.WithMinimalEventLabels(ctx, event)
138+
139+
ctx = tracing.ParseSpanContext(ctx, event)
140+
ctx, span := c.tracer.Start(ctx, jsmChannel+"-"+string(c.sub.UID))
141+
defer func() {
142+
if span.IsRecording() {
143+
// add full event labels here so that they only populate the span, and not any metrics
144+
ctx = observability.WithEventLabels(ctx, event)
145+
labeler, _ := otelhttp.LabelerFromContext(ctx)
146+
span.SetAttributes(labeler.Get()...)
147+
}
148+
span.End()
149+
}()
129150

130151
te := TypeExtractorTransformer("")
131152

@@ -143,10 +164,13 @@ func (c *PushConsumer) doHandle(ctx context.Context, msg *nats.Msg) {
143164
WithHeader(additionalHeaders),
144165
)
145166

146-
_ = fanout.ParseDispatchResultAndReportMetrics(fanout.NewDispatchResult(err, dispatchExecutionInfo), c.reporter, eventingchannels.ReportArgs{
147-
Ns: c.channelNamespace,
148-
EventType: string(te),
149-
})
167+
labeler, _ := otelhttp.LabelerFromContext(ctx)
168+
169+
c.dispatchDuration.Record(
170+
ctx,
171+
dispatchExecutionInfo.Duration.Seconds(),
172+
metric.WithAttributes(labeler.Get()...),
173+
)
150174

151175
if err != nil {
152176
logger.Errorw("failed to forward message to downstream subscriber",

0 commit comments

Comments
 (0)