Polish up rabbitmq-5.x plugin to fix missing broker tag on consumer side (#362)

This commit is contained in:
pg.yang 2022-10-21 21:19:57 +08:00 committed by GitHub
parent 6d5bbe8e5f
commit cf85f03ad5
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
9 changed files with 231 additions and 183 deletions

View File

@ -18,6 +18,7 @@ Release Notes.
* Polish up nats plugin to unify MQ related tags
* Correct the duration of the transaction span for Neo4J 4.x.
* Plugin-test configuration.yml dependencies support docker service command field
* Polish up rabbitmq-5.x plugin to fix missing broker tag on consumer side
#### Documentation

View File

@ -22,7 +22,7 @@ import com.rabbitmq.client.Connection;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.EnhancedInstance;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.InstanceConstructorInterceptor;
public class RabbitMQProducerAndConsumerConstructorInterceptor implements InstanceConstructorInterceptor {
public class ChannelNConstructorInterceptor implements InstanceConstructorInterceptor {
@Override
public void onConstruct(EnhancedInstance objInst, Object[] allArguments) {
Connection connection = (Connection) allArguments[0];

View File

@ -18,61 +18,29 @@
package org.apache.skywalking.apm.plugin.rabbitmq;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Envelope;
import org.apache.skywalking.apm.agent.core.context.CarrierItem;
import org.apache.skywalking.apm.agent.core.context.ContextCarrier;
import org.apache.skywalking.apm.agent.core.context.ContextManager;
import org.apache.skywalking.apm.agent.core.context.tag.Tags;
import org.apache.skywalking.apm.agent.core.context.trace.AbstractSpan;
import org.apache.skywalking.apm.agent.core.context.trace.SpanLayer;
import com.rabbitmq.client.Consumer;
import java.lang.reflect.Method;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.EnhancedInstance;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.InstanceMethodsAroundInterceptor;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.MethodInterceptResult;
import org.apache.skywalking.apm.network.trace.component.ComponentsDefine;
import java.lang.reflect.Method;
public class RabbitMQConsumerInterceptor implements InstanceMethodsAroundInterceptor {
public static final String OPERATE_NAME_PREFIX = "RabbitMQ/";
public static final String CONSUMER_OPERATE_NAME_SUFFIX = "/Consumer";
@Override
public void beforeMethod(EnhancedInstance objInst, Method method, Object[] allArguments, Class<?>[] argumentsTypes,
MethodInterceptResult result) throws Throwable {
ContextCarrier contextCarrier = new ContextCarrier();
String url = (String) objInst.getSkyWalkingDynamicField();
Envelope envelope = (Envelope) allArguments[1];
AMQP.BasicProperties properties = (AMQP.BasicProperties) allArguments[2];
AbstractSpan activeSpan = ContextManager.createEntrySpan(OPERATE_NAME_PREFIX + "Topic/" + envelope.getExchange() + "Queue/" + envelope
.getRoutingKey() + CONSUMER_OPERATE_NAME_SUFFIX, null).start(System.currentTimeMillis());
Tags.MQ_BROKER.set(activeSpan, url);
Tags.MQ_TOPIC.set(activeSpan, envelope.getExchange());
Tags.MQ_QUEUE.set(activeSpan, envelope.getRoutingKey());
activeSpan.setComponent(ComponentsDefine.RABBITMQ_CONSUMER);
SpanLayer.asMQ(activeSpan);
CarrierItem next = contextCarrier.items();
while (next.hasNext()) {
next = next.next();
if (properties.getHeaders() != null && properties.getHeaders().get(next.getHeadKey()) != null) {
next.setHeadValue(properties.getHeaders().get(next.getHeadKey()).toString());
}
}
ContextManager.extract(contextCarrier);
Consumer consumer = (Consumer) allArguments[6];
allArguments[6] = new TracerConsumer(consumer, (String) objInst.getSkyWalkingDynamicField());
}
@Override
public Object afterMethod(EnhancedInstance objInst, Method method, Object[] allArguments, Class<?>[] argumentsTypes,
Object ret) throws Throwable {
ContextManager.stopSpan();
return ret;
}
@Override
public void handleMethodException(EnhancedInstance objInst, Method method, Object[] allArguments,
Class<?>[] argumentsTypes, Throwable t) {
ContextManager.activeSpan().log(t);
}
}

View File

@ -0,0 +1,102 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.skywalking.apm.plugin.rabbitmq;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.Envelope;
import com.rabbitmq.client.ShutdownSignalException;
import java.io.IOException;
import org.apache.skywalking.apm.agent.core.context.CarrierItem;
import org.apache.skywalking.apm.agent.core.context.ContextCarrier;
import org.apache.skywalking.apm.agent.core.context.ContextManager;
import org.apache.skywalking.apm.agent.core.context.tag.Tags;
import org.apache.skywalking.apm.agent.core.context.trace.AbstractSpan;
import org.apache.skywalking.apm.agent.core.context.trace.SpanLayer;
import org.apache.skywalking.apm.network.trace.component.ComponentsDefine;
public class TracerConsumer implements Consumer {
private Consumer delegate;
private String serverUrl;
public static final String OPERATE_NAME_PREFIX = "RabbitMQ/";
public static final String CONSUMER_OPERATE_NAME_SUFFIX = "/Consumer";
public TracerConsumer(final Consumer delegate, final String serverUrl) {
this.delegate = delegate;
this.serverUrl = serverUrl;
}
@Override
public void handleConsumeOk(final String consumerTag) {
this.delegate.handleConsumeOk(consumerTag);
}
@Override
public void handleCancelOk(final String consumerTag) {
this.delegate.handleRecoverOk(consumerTag);
}
@Override
public void handleCancel(final String consumerTag) throws IOException {
this.delegate.handleCancel(consumerTag);
}
@Override
public void handleShutdownSignal(final String consumerTag, final ShutdownSignalException sig) {
this.delegate.handleShutdownSignal(consumerTag, sig);
}
@Override
public void handleRecoverOk(final String consumerTag) {
this.delegate.handleRecoverOk(consumerTag);
}
@Override
public void handleDelivery(final String consumerTag,
final Envelope envelope,
final AMQP.BasicProperties properties,
final byte[] body) throws IOException {
ContextCarrier contextCarrier = new ContextCarrier();
AbstractSpan activeSpan = ContextManager.createEntrySpan(
OPERATE_NAME_PREFIX + "Topic/" + envelope.getExchange() + "Queue/" + envelope
.getRoutingKey() + CONSUMER_OPERATE_NAME_SUFFIX, null).start(System.currentTimeMillis());
Tags.MQ_BROKER.set(activeSpan, serverUrl);
Tags.MQ_TOPIC.set(activeSpan, envelope.getExchange());
Tags.MQ_QUEUE.set(activeSpan, envelope.getRoutingKey());
activeSpan.setComponent(ComponentsDefine.RABBITMQ_CONSUMER);
SpanLayer.asMQ(activeSpan);
CarrierItem next = contextCarrier.items();
while (next.hasNext()) {
next = next.next();
if (properties.getHeaders() != null && properties.getHeaders().get(next.getHeadKey()) != null) {
next.setHeadValue(properties.getHeaders().get(next.getHeadKey()).toString());
}
}
ContextManager.extract(contextCarrier);
try {
this.delegate.handleDelivery(consumerTag, envelope, properties, body);
} catch (Exception e) {
activeSpan.log(e).errorOccurred();
} finally {
ContextManager.stopSpan(activeSpan);
}
}
}

View File

@ -21,19 +21,23 @@ package org.apache.skywalking.apm.plugin.rabbitmq.define;
import net.bytebuddy.description.method.MethodDescription;
import net.bytebuddy.matcher.ElementMatcher;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.ConstructorInterceptPoint;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.DeclaredInstanceMethodsInterceptPoint;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.InstanceMethodsInterceptPoint;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.ClassInstanceMethodsEnhancePluginDefine;
import org.apache.skywalking.apm.agent.core.plugin.match.ClassMatch;
import org.apache.skywalking.apm.agent.core.plugin.match.MultiClassNameMatch;
import static net.bytebuddy.matcher.ElementMatchers.named;
import static net.bytebuddy.matcher.ElementMatchers.takesArguments;
import static org.apache.skywalking.apm.agent.core.plugin.bytebuddy.ArgumentTypeNameMatch.takesArgumentWithType;
public class RabbitMQProducerInstrumentation extends ClassInstanceMethodsEnhancePluginDefine {
public class ChannelNInstrumentation extends ClassInstanceMethodsEnhancePluginDefine {
public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.rabbitmq.RabbitMQProducerInterceptor";
public static final String ENHANCE_CLASS_PRODUCER = "com.rabbitmq.client.impl.ChannelN";
public static final String ENHANCE_METHOD_DISPATCH = "basicPublish";
public static final String INTERCEPTOR_CONSTRUCTOR = "org.apache.skywalking.apm.plugin.rabbitmq.RabbitMQProducerAndConsumerConstructorInterceptor";
public static final String PUBLISH_ENHANCE_METHOD = "basicPublish";
public static final String INTERCEPTOR_CONSTRUCTOR = "org.apache.skywalking.apm.plugin.rabbitmq.ChannelNConstructorInterceptor";
public static final String CONSUME_ENHANCE_METHOD = "basicConsume";
public static final String CONSUME_INTERCEPTOR_CONSTRUCTOR = "org.apache.skywalking.apm.plugin.rabbitmq.RabbitMQConsumerInterceptor";
@Override
public ConstructorInterceptPoint[] getConstructorsInterceptPoints() {
@ -58,7 +62,7 @@ public class RabbitMQProducerInstrumentation extends ClassInstanceMethodsEnhance
new InstanceMethodsInterceptPoint() {
@Override
public ElementMatcher<MethodDescription> getMethodsMatcher() {
return named(ENHANCE_METHOD_DISPATCH).and(takesArgumentWithType(4, "com.rabbitmq.client.AMQP$BasicProperties"));
return named(PUBLISH_ENHANCE_METHOD).and(takesArgumentWithType(4, "com.rabbitmq.client.AMQP$BasicProperties"));
}
@Override
@ -66,6 +70,22 @@ public class RabbitMQProducerInstrumentation extends ClassInstanceMethodsEnhance
return INTERCEPTOR_CLASS;
}
@Override
public boolean isOverrideArgs() {
return true;
}
},
new DeclaredInstanceMethodsInterceptPoint() {
@Override
public ElementMatcher<MethodDescription> getMethodsMatcher() {
return named(CONSUME_ENHANCE_METHOD).and(takesArguments(7));
}
@Override
public String getMethodsInterceptor() {
return CONSUME_INTERCEPTOR_CONSTRUCTOR;
}
@Override
public boolean isOverrideArgs() {
return true;

View File

@ -1,82 +0,0 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.apm.plugin.rabbitmq.define;
import net.bytebuddy.description.method.MethodDescription;
import net.bytebuddy.matcher.ElementMatcher;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.ConstructorInterceptPoint;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.DeclaredInstanceMethodsInterceptPoint;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.InstanceMethodsInterceptPoint;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.ClassInstanceMethodsEnhancePluginDefine;
import org.apache.skywalking.apm.agent.core.plugin.match.ClassMatch;
import org.apache.skywalking.apm.agent.core.plugin.match.HierarchyMatch;
import static net.bytebuddy.matcher.ElementMatchers.named;
import static org.apache.skywalking.apm.agent.core.plugin.bytebuddy.ArgumentTypeNameMatch.takesArgumentWithType;
public class RabbitMQConsumerInstrumentation extends ClassInstanceMethodsEnhancePluginDefine {
public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.rabbitmq.RabbitMQConsumerInterceptor";
public static final String ENHANCE_CLASS_PRODUCER = "com.rabbitmq.client.Consumer";
public static final String ENHANCE_METHOD_DISPATCH = "handleDelivery";
public static final String INTERCEPTOR_CONSTRUCTOR = "org.apache.skywalking.apm.plugin.rabbitmq.RabbitMQProducerAndConsumerConstructorInterceptor";
@Override
public ConstructorInterceptPoint[] getConstructorsInterceptPoints() {
return new ConstructorInterceptPoint[] {
new ConstructorInterceptPoint() {
@Override
public ElementMatcher<MethodDescription> getConstructorMatcher() {
return takesArgumentWithType(0, "com.rabbitmq.client.impl.AMQConnection");
}
@Override
public String getConstructorInterceptor() {
return INTERCEPTOR_CONSTRUCTOR;
}
}
};
}
@Override
public InstanceMethodsInterceptPoint[] getInstanceMethodsInterceptPoints() {
return new InstanceMethodsInterceptPoint[] {
new DeclaredInstanceMethodsInterceptPoint() {
@Override
public ElementMatcher<MethodDescription> getMethodsMatcher() {
return named(ENHANCE_METHOD_DISPATCH).and(takesArgumentWithType(2, "com.rabbitmq.client.AMQP$BasicProperties"));
}
@Override
public String getMethodsInterceptor() {
return INTERCEPTOR_CLASS;
}
@Override
public boolean isOverrideArgs() {
return false;
}
}
};
}
@Override
protected ClassMatch enhanceClass() {
return HierarchyMatch.byHierarchyMatch(new String[] {ENHANCE_CLASS_PRODUCER});
}
}

View File

@ -14,5 +14,4 @@
# See the License for the specific language governing permissions and
# limitations under the License.
rabbitmq-5.x=org.apache.skywalking.apm.plugin.rabbitmq.define.RabbitMQProducerInstrumentation
rabbitmq-5.x=org.apache.skywalking.apm.plugin.rabbitmq.define.RabbitMQConsumerInstrumentation
rabbitmq-5.x=org.apache.skywalking.apm.plugin.rabbitmq.define.ChannelNInstrumentation

View File

@ -43,9 +43,9 @@ import static org.hamcrest.core.Is.is;
@RunWith(PowerMockRunner.class)
@PowerMockRunnerDelegate(TracingSegmentRunner.class)
public class RabbitMQProducerAndConsumerConstructorInterceptorTest {
public class ChannelNConstructorInterceptorTest {
private RabbitMQProducerAndConsumerConstructorInterceptor rabbitMQProducerAndConsumerConstructorInterceptor;
private ChannelNConstructorInterceptor channelNConstructorInterceptor;
private EnhancedInstance enhancedInstance = new EnhancedInstance() {
private String test;
@ -230,8 +230,8 @@ public class RabbitMQProducerAndConsumerConstructorInterceptorTest {
@Test
public void TestRabbitMQConsumerAndProducerConstructorInterceptor() {
rabbitMQProducerAndConsumerConstructorInterceptor = new RabbitMQProducerAndConsumerConstructorInterceptor();
rabbitMQProducerAndConsumerConstructorInterceptor.onConstruct(enhancedInstance, new Object[] {testConnection});
channelNConstructorInterceptor = new ChannelNConstructorInterceptor();
channelNConstructorInterceptor.onConstruct(enhancedInstance, new Object[] {testConnection});
assertThat((String) enhancedInstance.getSkyWalkingDynamicField(), is("127.0.0.1:5672"));
}
}

View File

@ -19,7 +19,13 @@
package org.apache.skywalking.apm.plugin.rabbitmq;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Consumer;
import com.rabbitmq.client.Envelope;
import com.rabbitmq.client.ShutdownSignalException;
import java.io.IOException;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import org.apache.skywalking.apm.agent.core.context.SW8CarrierItem;
import org.apache.skywalking.apm.agent.core.context.trace.TraceSegment;
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.EnhancedInstance;
@ -35,10 +41,6 @@ import org.junit.runner.RunWith;
import org.powermock.modules.junit4.PowerMockRunner;
import org.powermock.modules.junit4.PowerMockRunnerDelegate;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import static org.hamcrest.CoreMatchers.is;
@RunWith(PowerMockRunner.class)
@ -48,75 +50,113 @@ public class RabbitMQConsumerInterceptorTest {
@SegmentStoragePoint
private SegmentStorage segmentStorage;
private RabbitMQConsumerInterceptor rabbitMQConsumerInterceptor;
@Rule
public AgentServiceRule serviceRule = new AgentServiceRule();
private EnhancedInstance enhancedInstance = new EnhancedInstance() {
@Override
public Object getSkyWalkingDynamicField() {
return "127.0.0.1:5272";
}
@Override
public void setSkyWalkingDynamicField(Object value) {
}
};
private RabbitMQConsumerInterceptor rabbitMQConsumerInterceptor;
@Before
public void setUp() throws Exception {
rabbitMQConsumerInterceptor = new RabbitMQConsumerInterceptor();
}
@Test
public void TestRabbitMQConsumerInterceptor() throws Throwable {
Envelope envelope = new Envelope(1111, false, "", "rabbitmq-test");
Map<String, Object> headers = new HashMap<String, Object>();
headers.put(SW8CarrierItem.HEADER_NAME, "1-My40LjU=-MS4yLjM=-3-c2VydmljZQ==-aW5zdGFuY2U=-L2FwcA==-MTI3LjAuMC4xOjgwODA=");
AMQP.BasicProperties.Builder propsBuilder = new AMQP.BasicProperties.Builder();
Object[] arguments = new Object[] {
0,
envelope,
propsBuilder.headers(headers).build()
};
rabbitMQConsumerInterceptor.beforeMethod(enhancedInstance, null, arguments, null, null);
rabbitMQConsumerInterceptor.afterMethod(enhancedInstance, null, arguments, null, null);
List<TraceSegment> traceSegments = segmentStorage.getTraceSegments();
Assert.assertThat(traceSegments.size(), is(1));
}
@Test
public void testRabbitMQConsumerInterceptorWithNilHeaders() throws Throwable {
Envelope envelope = new Envelope(1111, false, "", "rabbitmq-test");
AMQP.BasicProperties.Builder propsBuilder = new AMQP.BasicProperties.Builder();
Object[] arguments = new Object[] {
0,
envelope,
propsBuilder.headers(null).build()
final Object[] args = {
null,
null,
null,
null,
null,
null,
getConsumer()
};
rabbitMQConsumerInterceptor.beforeMethod(enhancedInstance, null, arguments, null, null);
rabbitMQConsumerInterceptor.afterMethod(enhancedInstance, null, arguments, null, null);
rabbitMQConsumerInterceptor.beforeMethod(getEnhancedInstance(), null, args, new Class[0], null);
rabbitMQConsumerInterceptor.afterMethod(getEnhancedInstance(), null, args, new Class[0], null);
((Consumer) args[6]).handleDelivery("tag", new Envelope(1L, false, "exchange", "routerKey"),
new AMQP.BasicProperties(), new byte[0]
);
List<TraceSegment> traceSegments = segmentStorage.getTraceSegments();
Assert.assertThat(traceSegments.size(), is(1));
}
@Test
public void testRabbitMQConsumerInterceptorWithEmptyHeaders() throws Throwable {
Envelope envelope = new Envelope(1111, false, "", "rabbitmq-test");
Map<String, Object> headers = new HashMap<String, Object>();
AMQP.BasicProperties.Builder propsBuilder = new AMQP.BasicProperties.Builder();
Object[] arguments = new Object[] {
0,
envelope,
propsBuilder.headers(headers).build()
public void testRabbitMQConsumerInterceptor() throws Throwable {
final Object[] args = {
null,
null,
null,
null,
null,
null,
getConsumer()
};
rabbitMQConsumerInterceptor.beforeMethod(enhancedInstance, null, arguments, null, null);
rabbitMQConsumerInterceptor.afterMethod(enhancedInstance, null, arguments, null, null);
Map<String, Object> headers = new HashMap<>();
headers.put(
SW8CarrierItem.HEADER_NAME,
"1-My40LjU=-MS4yLjM=-3-c2VydmljZQ==-aW5zdGFuY2U=-L2FwcA==-MTI3LjAuMC4xOjgwODA="
);
AMQP.BasicProperties.Builder propsBuilder = new AMQP.BasicProperties.Builder();
propsBuilder.headers(headers);
rabbitMQConsumerInterceptor.beforeMethod(getEnhancedInstance(), null, args, new Class[0], null);
rabbitMQConsumerInterceptor.afterMethod(getEnhancedInstance(), null, args, new Class[0], null);
((Consumer) args[6]).handleDelivery("tag", new Envelope(1L, false, "exchange", "routerKey"),
propsBuilder.build(), new byte[0]
);
List<TraceSegment> traceSegments = segmentStorage.getTraceSegments();
Assert.assertThat(traceSegments.size(), is(1));
}
public EnhancedInstance getEnhancedInstance() {
return new EnhancedInstance() {
@Override
public Object getSkyWalkingDynamicField() {
return "serverAddr";
}
@Override
public void setSkyWalkingDynamicField(final Object value) {
}
};
}
public Consumer getConsumer() {
return new Consumer() {
@Override
public void handleConsumeOk(final String consumerTag) {
}
@Override
public void handleCancelOk(final String consumerTag) {
}
@Override
public void handleCancel(final String consumerTag) throws IOException {
}
@Override
public void handleShutdownSignal(final String consumerTag, final ShutdownSignalException sig) {
}
@Override
public void handleRecoverOk(final String consumerTag) {
}
@Override
public void handleDelivery(final String consumerTag,
final Envelope envelope,
final AMQP.BasicProperties properties,
final byte[] body) throws IOException {
}
};
}
}