fix rocketmq message header properties garbled characters issue (#54)

This commit is contained in:
cnScarb 2021-10-27 09:16:27 +08:00 committed by GitHub
parent 3f93954c5c
commit 30fb81a6de
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
3 changed files with 56 additions and 4 deletions

View File

@ -39,6 +39,7 @@ Release Notes.
* Fix instrumentation v2 API doesn't work for constructor instrumentation. * Fix instrumentation v2 API doesn't work for constructor instrumentation.
* Add plugin to support okhttp 2.x * Add plugin to support okhttp 2.x
* Optimize okhttp 3.x 4.x plugin to get span time cost precisely * Optimize okhttp 3.x 4.x plugin to get span time cost precisely
* Adapt message header properties of RocketMQ 4.9.x
#### Documentation #### Documentation

View File

@ -68,10 +68,13 @@ public class MessageSendInterceptor implements InstanceMethodsAroundInterceptor
while (next.hasNext()) { while (next.hasNext()) {
next = next.next(); next = next.next();
if (!StringUtil.isEmpty(next.getHeadValue())) { if (!StringUtil.isEmpty(next.getHeadValue())) {
if (properties.length() > 0 && properties.charAt(properties.length() - 1) != PROPERTY_SEPARATOR) {
// adapt for RocketMQ 4.9.x or later
properties.append(PROPERTY_SEPARATOR);
}
properties.append(next.getHeadKey()); properties.append(next.getHeadKey());
properties.append(NAME_VALUE_SEPARATOR); properties.append(NAME_VALUE_SEPARATOR);
properties.append(next.getHeadValue()); properties.append(next.getHeadValue());
properties.append(PROPERTY_SEPARATOR);
} }
} }
requestHeader.setProperties(properties.toString()); requestHeader.setProperties(properties.toString());

View File

@ -19,9 +19,13 @@
package org.apache.skywalking.apm.plugin.rocketMQ.v4; package org.apache.skywalking.apm.plugin.rocketMQ.v4;
import java.util.List; import java.util.List;
import java.util.Map;
import org.apache.rocketmq.client.impl.CommunicationMode; import org.apache.rocketmq.client.impl.CommunicationMode;
import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
import org.apache.skywalking.apm.agent.core.context.SW8ExtensionCarrierItem;
import org.apache.skywalking.apm.agent.test.tools.SegmentStorage; import org.apache.skywalking.apm.agent.test.tools.SegmentStorage;
import org.junit.Before; import org.junit.Before;
import org.junit.Rule; import org.junit.Rule;
@ -44,10 +48,17 @@ import org.apache.skywalking.apm.network.trace.component.ComponentsDefine;
import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.CoreMatchers.is;
import static org.junit.Assert.assertThat; import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import static org.mockito.Matchers.anyString; import static org.mockito.Matchers.anyString;
import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static org.powermock.api.mockito.PowerMockito.when; import static org.powermock.api.mockito.PowerMockito.when;
import static org.apache.rocketmq.common.message.MessageDecoder.NAME_VALUE_SEPARATOR;
import static org.apache.rocketmq.common.message.MessageDecoder.PROPERTY_SEPARATOR;
@RunWith(PowerMockRunner.class) @RunWith(PowerMockRunner.class)
@PowerMockRunnerDelegate(TracingSegmentRunner.class) @PowerMockRunnerDelegate(TracingSegmentRunner.class)
public class MessageSendInterceptorTest { public class MessageSendInterceptorTest {
@ -108,8 +119,8 @@ public class MessageSendInterceptorTest {
CommunicationMode.ASYNC, CommunicationMode.ASYNC,
null null
}; };
when(messageRequestHeader.getProperties()).thenReturn("");
when(message.getTags()).thenReturn("TagA"); when(message.getTags()).thenReturn("TagA");
stubMessageRequestHeader("TAGS" + NAME_VALUE_SEPARATOR + "TagA" + PROPERTY_SEPARATOR);
} }
@Test @Test
@ -117,6 +128,13 @@ public class MessageSendInterceptorTest {
messageSendInterceptor.beforeMethod(enhancedInstance, null, arguments, null, null); messageSendInterceptor.beforeMethod(enhancedInstance, null, arguments, null, null);
messageSendInterceptor.afterMethod(enhancedInstance, null, arguments, null, null); messageSendInterceptor.afterMethod(enhancedInstance, null, arguments, null, null);
Map<String, String> tags = MessageDecoder.string2messageProperties(
((SendMessageRequestHeader) arguments[3]).getProperties());
// check original header of TAGS
assertThat(tags.get("TAGS"), is("TagA"));
// check skywalking header
assertTrue(tags.containsKey(SW8ExtensionCarrierItem.HEADER_NAME));
assertThat(segmentStorage.getTraceSegments().size(), is(1)); assertThat(segmentStorage.getTraceSegments().size(), is(1));
TraceSegment traceSegment = segmentStorage.getTraceSegments().get(0); TraceSegment traceSegment = segmentStorage.getTraceSegments().get(0);
List<AbstractTracingSpan> spans = SegmentHelper.getSpans(traceSegment); List<AbstractTracingSpan> spans = SegmentHelper.getSpans(traceSegment);
@ -127,15 +145,27 @@ public class MessageSendInterceptorTest {
SpanAssert.assertLayer(mqSpan, SpanLayer.MQ); SpanAssert.assertLayer(mqSpan, SpanLayer.MQ);
SpanAssert.assertComponent(mqSpan, ComponentsDefine.ROCKET_MQ_PRODUCER); SpanAssert.assertComponent(mqSpan, ComponentsDefine.ROCKET_MQ_PRODUCER);
SpanAssert.assertTag(mqSpan, 0, "127.0.0.1"); SpanAssert.assertTag(mqSpan, 0, "127.0.0.1");
verify(messageRequestHeader).setProperties(anyString());
verify(callBack).setSkyWalkingDynamicField(Matchers.any()); verify(callBack).setSkyWalkingDynamicField(Matchers.any());
} }
@Test
public void testSendMessageNew() throws Throwable {
stubMessageRequestHeader("TAGS" + NAME_VALUE_SEPARATOR + "TagA");
testSendMessage();
}
@Test @Test
public void testSendMessageWithoutCallBack() throws Throwable { public void testSendMessageWithoutCallBack() throws Throwable {
messageSendInterceptor.beforeMethod(enhancedInstance, null, argumentsWithoutCallback, null, null); messageSendInterceptor.beforeMethod(enhancedInstance, null, argumentsWithoutCallback, null, null);
messageSendInterceptor.afterMethod(enhancedInstance, null, argumentsWithoutCallback, null, null); messageSendInterceptor.afterMethod(enhancedInstance, null, argumentsWithoutCallback, null, null);
Map<String, String> tags = MessageDecoder.string2messageProperties(
((SendMessageRequestHeader) argumentsWithoutCallback[3]).getProperties());
// check original header of TAGS
assertThat(tags.get("TAGS"), is("TagA"));
// check skywalking header
assertTrue(tags.containsKey(SW8ExtensionCarrierItem.HEADER_NAME));
assertThat(segmentStorage.getTraceSegments().size(), is(1)); assertThat(segmentStorage.getTraceSegments().size(), is(1));
TraceSegment traceSegment = segmentStorage.getTraceSegments().get(0); TraceSegment traceSegment = segmentStorage.getTraceSegments().get(0);
List<AbstractTracingSpan> spans = SegmentHelper.getSpans(traceSegment); List<AbstractTracingSpan> spans = SegmentHelper.getSpans(traceSegment);
@ -146,7 +176,25 @@ public class MessageSendInterceptorTest {
SpanAssert.assertLayer(mqSpan, SpanLayer.MQ); SpanAssert.assertLayer(mqSpan, SpanLayer.MQ);
SpanAssert.assertComponent(mqSpan, ComponentsDefine.ROCKET_MQ_PRODUCER); SpanAssert.assertComponent(mqSpan, ComponentsDefine.ROCKET_MQ_PRODUCER);
SpanAssert.assertTag(mqSpan, 0, "127.0.0.1"); SpanAssert.assertTag(mqSpan, 0, "127.0.0.1");
verify(messageRequestHeader).setProperties(anyString());
} }
@Test
public void testSendMessageWithoutCallBackNew() throws Throwable {
stubMessageRequestHeader("TAGS" + NAME_VALUE_SEPARATOR + "TagA");
testSendMessageWithoutCallBack();
}
private void stubMessageRequestHeader(String properties) {
messageRequestHeader = mock(SendMessageRequestHeader.class, RETURNS_DEEP_STUBS);
doAnswer(invocation -> {
String val = (String) invocation.getArguments()[0];
when(messageRequestHeader.getProperties()).thenReturn(val);
return null;
}).when(messageRequestHeader).setProperties(anyString());
when(messageRequestHeader.getProperties()).thenCallRealMethod();
messageRequestHeader.setProperties(properties);
arguments[3] = messageRequestHeader;
argumentsWithoutCallback[3] = messageRequestHeader;
}
} }