Skip to content

Commit

Permalink
fix(amqplib): use extracted context for message consuming (#1354)
Browse files Browse the repository at this point in the history
* fix(amqplib): use extracted context for message consuming
- fixes baggage propagation

* style(amqplib): fix linting

---------

Co-authored-by: Gerhard Stöbich <[email protected]>
Co-authored-by: Amir Blum <[email protected]>
  • Loading branch information
3 people authored Feb 7, 2023
1 parent 799b10b commit ad92673
Show file tree
Hide file tree
Showing 2 changed files with 46 additions and 2 deletions.
2 changes: 1 addition & 1 deletion plugins/node/instrumentation-amqplib/src/amqplib.ts
Original file line number Diff line number Diff line change
Expand Up @@ -451,7 +451,7 @@ export class AmqplibInstrumentation extends InstrumentationBase<typeof amqp> {
msg[MESSAGE_STORED_SPAN] = span;
}

context.with(trace.setSpan(context.active(), span), () => {
context.with(trace.setSpan(parentContext, span), () => {
onMessage.call(this, msg);
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,20 +28,35 @@ import {
MessagingDestinationKindValues,
SemanticAttributes,
} from '@opentelemetry/semantic-conventions';
import { context, SpanKind } from '@opentelemetry/api';
import { Baggage, context, propagation, SpanKind } from '@opentelemetry/api';
import { asyncConfirmSend, asyncConsume, shouldTest } from './utils';
import {
censoredUrl,
rabbitMqUrl,
TEST_RABBITMQ_HOST,
TEST_RABBITMQ_PORT,
} from './config';
import {
CompositePropagator,
W3CBaggagePropagator,
W3CTraceContextPropagator,
} from '@opentelemetry/core';

const msgPayload = 'payload from test';
const queueName = 'queue-name-from-unittest';

describe('amqplib instrumentation callback model', () => {
let conn: amqpCallback.Connection;
before(() => {
propagation.setGlobalPropagator(
new CompositePropagator({
propagators: [
new W3CBaggagePropagator(),
new W3CTraceContextPropagator(),
],
})
);
});
before(function (done) {
if (!shouldTest) {
this.skip();
Expand Down Expand Up @@ -186,6 +201,35 @@ describe('amqplib instrumentation callback model', () => {
});
});

it('baggage is available while consuming', done => {
const baggageContext = propagation.setBaggage(
context.active(),
propagation.createBaggage({
key1: { value: 'value1' },
})
);
context.with(baggageContext, () => {
channel.sendToQueue(queueName, Buffer.from(msgPayload));
let extractedBaggage: Baggage | undefined;
asyncConsume(
channel,
queueName,
[
msg => {
extractedBaggage = propagation.getActiveBaggage();
},
],
{
noAck: true,
}
).then(() => {
expect(extractedBaggage).toBeDefined();
expect(extractedBaggage!.getEntry('key1')).toBeDefined();
done();
});
});
});

it('end span with ack sync', done => {
channel.sendToQueue(queueName, Buffer.from(msgPayload));

Expand Down

0 comments on commit ad92673

Please sign in to comment.