Add thrift plugin support thrift TMultiplexedProcessor. (#22)
This commit is contained in:
parent
a27cf2af71
commit
0e45ed0b13
|
|
@ -13,14 +13,14 @@
|
|||
|
||||
<!-- ==== 🔌 Remove this line WHEN AND ONLY WHEN you're adding a new plugin, follow the checklist 👇 ====
|
||||
### Add an agent plugin to support <framework name>
|
||||
- [ ] Add a test case for the new plugin, refer to [the doc](https://github.com/apache/skywalking/blob/master/docs/en/guides/Plugin-test.md)
|
||||
- [ ] Add a component id in [the component-libraries.yml](https://github.com/apache/skywalking/blob/master/oap-server/server-bootstrap/src/main/resources/component-libraries.yml)
|
||||
- [ ] Add a test case for the new plugin, refer to [the doc](https://github.com/apache/skywalking-java/blob/main/docs/en/setup/service-agent/java-agent/Plugin-test.md)
|
||||
- [ ] Add a component id in [the component-libraries.yml](https://github.com/apache/skywalking/blob/master/oap-server/server-starter/src/main/resources/component-libraries.yml)
|
||||
- [ ] Add a logo in [the UI repo](https://github.com/apache/skywalking-rocketbot-ui/tree/master/src/views/components/topology/assets)
|
||||
==== 🔌 Remove this line WHEN AND ONLY WHEN you're adding a new plugin, follow the checklist 👆 ==== -->
|
||||
|
||||
<!-- ==== 📈 Remove this line WHEN AND ONLY WHEN you're improving the performance, follow the checklist 👇 ====
|
||||
### Improve the performance of <class or module or ...>
|
||||
- [ ] Add a benchmark for the improvement, refer to [the existing ones](https://github.com/apache/skywalking/blob/master/apm-commons/apm-datacarrier/src/test/java/org/apache/skywalking/apm/commons/datacarrier/LinkedArrayBenchmark.java)
|
||||
- [ ] Add a benchmark for the improvement, refer to [the existing ones](https://github.com/apache/skywalking-java/blob/main/apm-commons/apm-datacarrier/src/test/java/org/apache/skywalking/apm/commons/datacarrier/LinkedArrayBenchmark.java)
|
||||
- [ ] The benchmark result.
|
||||
```text
|
||||
<Paste the benchmark results here>
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ Release Notes.
|
|||
* Support mTLS for gRPC channel.
|
||||
* fix the bug that plugin record wrong time elapse for lettuce plugin
|
||||
* fix the bug that the wrong db.instance value displayed on Skywalking-UI when existing multi-database-instance on same host port pair.
|
||||
* Add thrift plugin support thrift TMultiplexedProcessor.
|
||||
|
||||
#### Documentation
|
||||
|
||||
|
|
|
|||
|
|
@ -51,8 +51,10 @@ public class TBaseProcessorInterceptor implements InstanceConstructorInterceptor
|
|||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
MethodInterceptResult result) throws Throwable {
|
||||
ServerInProtocolWrapper in = (ServerInProtocolWrapper) allArguments[0];
|
||||
in.initial(new Context(processMap));
|
||||
if (allArguments[0] instanceof ServerInProtocolWrapper) {
|
||||
ServerInProtocolWrapper in = (ServerInProtocolWrapper) allArguments[0];
|
||||
in.initial(new Context(processMap));
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -61,7 +63,9 @@ public class TBaseProcessorInterceptor implements InstanceConstructorInterceptor
|
|||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
Object ret) throws Throwable {
|
||||
ContextManager.stopSpan();
|
||||
if (allArguments[0] instanceof ServerInProtocolWrapper) {
|
||||
ContextManager.stopSpan();
|
||||
}
|
||||
return ret;
|
||||
}
|
||||
|
||||
|
|
@ -71,6 +75,8 @@ public class TBaseProcessorInterceptor implements InstanceConstructorInterceptor
|
|||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
Throwable t) {
|
||||
ContextManager.activeSpan().log(t);
|
||||
if (allArguments[0] instanceof ServerInProtocolWrapper) {
|
||||
ContextManager.activeSpan().log(t);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,82 @@
|
|||
/*
|
||||
* 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.thrift;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.apache.skywalking.apm.agent.core.context.ContextManager;
|
||||
import org.apache.skywalking.apm.agent.core.logging.api.ILog;
|
||||
import org.apache.skywalking.apm.agent.core.logging.api.LogManager;
|
||||
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.EnhancedInstance;
|
||||
import org.apache.skywalking.apm.agent.core.plugin.interceptor.enhance.InstanceConstructorInterceptor;
|
||||
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.plugin.thrift.wrapper.Context;
|
||||
import org.apache.skywalking.apm.plugin.thrift.wrapper.ServerInProtocolWrapper;
|
||||
import org.apache.thrift.ProcessFunction;
|
||||
import org.apache.thrift.TBaseAsyncProcessor;
|
||||
|
||||
/**
|
||||
* To wrap the ProcessFunction for getting arguments of method.
|
||||
*
|
||||
* @see TBaseAsyncProcessor
|
||||
* @see TBaseProcessorInterceptor
|
||||
* @see TMultiplexedProcessorInterceptor
|
||||
*/
|
||||
public class TMultiplexedProcessorInterceptor implements InstanceConstructorInterceptor, InstanceMethodsAroundInterceptor {
|
||||
private Map<String, ProcessFunction> processMap = new HashMap<>();
|
||||
|
||||
private static final ILog LOGGER = LogManager.getLogger(TMultiplexedProcessorInterceptor.class);
|
||||
|
||||
@Override
|
||||
public void onConstruct(final EnhancedInstance objInst, final Object[] allArguments) throws Throwable {
|
||||
objInst.setSkyWalkingDynamicField(processMap);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void beforeMethod(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
MethodInterceptResult result) throws Throwable {
|
||||
|
||||
ServerInProtocolWrapper in = (ServerInProtocolWrapper) allArguments[0];
|
||||
in.initial(new Context(processMap));
|
||||
}
|
||||
|
||||
@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);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,80 @@
|
|||
/*
|
||||
* 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.thrift;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.apache.skywalking.apm.agent.core.logging.api.ILog;
|
||||
import org.apache.skywalking.apm.agent.core.logging.api.LogManager;
|
||||
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.thrift.ProcessFunction;
|
||||
import org.apache.thrift.TBaseAsyncProcessor;
|
||||
import org.apache.thrift.TBaseProcessor;
|
||||
import org.apache.thrift.TProcessor;
|
||||
|
||||
public class TMultiplexedProcessorRegisterDefaultInterceptor implements InstanceMethodsAroundInterceptor {
|
||||
|
||||
private static final ILog LOGGER = LogManager.getLogger(TMultiplexedProcessorRegisterDefaultInterceptor.class);
|
||||
|
||||
@Override
|
||||
public void beforeMethod(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
MethodInterceptResult result) throws Throwable {
|
||||
|
||||
Map<String, ProcessFunction> processMap = (Map<String, ProcessFunction>) objInst.getSkyWalkingDynamicField();
|
||||
TProcessor processor = (TProcessor) allArguments[0];
|
||||
processMap.putAll(getProcessMap(processor));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object afterMethod(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
Object ret) throws Throwable {
|
||||
return ret;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleMethodException(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
Throwable t) {
|
||||
}
|
||||
|
||||
private Map<String, ProcessFunction> getProcessMap(TProcessor processor) {
|
||||
Map<String, ProcessFunction> hashMap = new HashMap<>();
|
||||
if (processor instanceof TBaseProcessor) {
|
||||
Map<String, ProcessFunction> processMapView = ((TBaseProcessor) processor).getProcessMapView();
|
||||
hashMap.putAll(processMapView);
|
||||
} else if (processor instanceof TBaseAsyncProcessor) {
|
||||
Map<String, ProcessFunction> processMapView = ((TBaseProcessor) processor).getProcessMapView();
|
||||
hashMap.putAll(processMapView);
|
||||
} else {
|
||||
LOGGER.warn("Not support this processor:{}", processor.getClass().getName());
|
||||
}
|
||||
return hashMap;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,82 @@
|
|||
/*
|
||||
* 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.thrift;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import org.apache.skywalking.apm.agent.core.logging.api.ILog;
|
||||
import org.apache.skywalking.apm.agent.core.logging.api.LogManager;
|
||||
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.thrift.ProcessFunction;
|
||||
import org.apache.thrift.TBaseAsyncProcessor;
|
||||
import org.apache.thrift.TBaseProcessor;
|
||||
import org.apache.thrift.TProcessor;
|
||||
import org.apache.thrift.protocol.TMultiplexedProtocol;
|
||||
|
||||
public class TMultiplexedProcessorRegisterInterceptor implements InstanceMethodsAroundInterceptor {
|
||||
|
||||
private static final ILog LOGGER = LogManager.getLogger(TMultiplexedProcessorRegisterInterceptor.class);
|
||||
|
||||
@Override
|
||||
public void beforeMethod(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
MethodInterceptResult result) throws Throwable {
|
||||
|
||||
Map<String, ProcessFunction> processMap = (Map<String, ProcessFunction>) objInst.getSkyWalkingDynamicField();
|
||||
String serviceName = (String) allArguments[0];
|
||||
TProcessor processor = (TProcessor) allArguments[1];
|
||||
processMap.putAll(getProcessMap(serviceName, processor));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object afterMethod(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
Object ret) throws Throwable {
|
||||
return ret;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleMethodException(EnhancedInstance objInst,
|
||||
Method method,
|
||||
Object[] allArguments,
|
||||
Class<?>[] argumentsTypes,
|
||||
Throwable t) {
|
||||
}
|
||||
|
||||
private Map<String, ProcessFunction> getProcessMap(String serviceName, TProcessor processor) {
|
||||
Map<String, ProcessFunction> hashMap = new HashMap<>();
|
||||
if (processor instanceof TBaseProcessor) {
|
||||
Map<String, ProcessFunction> processMapView = ((TBaseProcessor) processor).getProcessMapView();
|
||||
processMapView.forEach((k, v) -> hashMap.put(serviceName + TMultiplexedProtocol.SEPARATOR + k, v));
|
||||
} else if (processor instanceof TBaseAsyncProcessor) {
|
||||
Map<String, ProcessFunction> processMapView = ((TBaseAsyncProcessor) processor).getProcessMapView();
|
||||
processMapView.forEach((k, v) -> hashMap.put(serviceName + TMultiplexedProtocol.SEPARATOR + k, v));
|
||||
} else {
|
||||
LOGGER.warn("Not support this processor:{}", processor.getClass().getName());
|
||||
}
|
||||
return hashMap;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,115 @@
|
|||
/*
|
||||
* 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.thrift.define;
|
||||
|
||||
import net.bytebuddy.description.method.MethodDescription;
|
||||
import net.bytebuddy.matcher.ElementMatcher;
|
||||
import net.bytebuddy.matcher.ElementMatchers;
|
||||
import org.apache.skywalking.apm.agent.core.plugin.interceptor.ConstructorInterceptPoint;
|
||||
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.NameMatch;
|
||||
|
||||
public class TMultiplexedProcessorInstrumentation extends ClassInstanceMethodsEnhancePluginDefine {
|
||||
private static final String ENHANCE_CLASS = "org.apache.thrift.TMultiplexedProcessor";
|
||||
private static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.thrift.TMultiplexedProcessorInterceptor";
|
||||
private static final String REGISTER_PROCESSOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.thrift.TMultiplexedProcessorRegisterInterceptor";
|
||||
private static final String REGISTER_DEFAULT_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.thrift.TMultiplexedProcessorRegisterDefaultInterceptor";
|
||||
|
||||
@Override
|
||||
public ConstructorInterceptPoint[] getConstructorsInterceptPoints() {
|
||||
return new ConstructorInterceptPoint[] {
|
||||
new ConstructorInterceptPoint() {
|
||||
@Override
|
||||
public ElementMatcher<MethodDescription> getConstructorMatcher() {
|
||||
return ElementMatchers.any();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getConstructorInterceptor() {
|
||||
return INTERCEPTOR_CLASS;
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public InstanceMethodsInterceptPoint[] getInstanceMethodsInterceptPoints() {
|
||||
return new InstanceMethodsInterceptPoint[] {
|
||||
new InstanceMethodsInterceptPoint() {
|
||||
|
||||
@Override
|
||||
public ElementMatcher<MethodDescription> getMethodsMatcher() {
|
||||
return ElementMatchers.named("process").and(ElementMatchers.takesArguments(2));
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getMethodsInterceptor() {
|
||||
return INTERCEPTOR_CLASS;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isOverrideArgs() {
|
||||
return false;
|
||||
}
|
||||
},
|
||||
new InstanceMethodsInterceptPoint() {
|
||||
|
||||
@Override
|
||||
public ElementMatcher<MethodDescription> getMethodsMatcher() {
|
||||
return ElementMatchers.named("registerProcessor").and(ElementMatchers.takesArguments(2));
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getMethodsInterceptor() {
|
||||
return REGISTER_PROCESSOR_INTERCEPTOR_CLASS;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isOverrideArgs() {
|
||||
return false;
|
||||
}
|
||||
},
|
||||
new InstanceMethodsInterceptPoint() {
|
||||
|
||||
@Override
|
||||
public ElementMatcher<MethodDescription> getMethodsMatcher() {
|
||||
return ElementMatchers.named("registerDefault");
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getMethodsInterceptor() {
|
||||
return REGISTER_DEFAULT_INTERCEPTOR_CLASS;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isOverrideArgs() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
protected ClassMatch enhanceClass() {
|
||||
return NameMatch.byName(ENHANCE_CLASS);
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -20,4 +20,5 @@ thrift=org.apache.skywalking.apm.plugin.thrift.define.client.TAsyncMethodCallIns
|
|||
thrift=org.apache.skywalking.apm.plugin.thrift.define.client.TServiceClientInstrumentation
|
||||
thrift=org.apache.skywalking.apm.plugin.thrift.define.TBaseProcessorInstrumentation
|
||||
thrift=org.apache.skywalking.apm.plugin.thrift.define.TBaseAsyncProcessorInstrumentation
|
||||
thrift=org.apache.skywalking.apm.plugin.thrift.define.TMultiplexedProcessorInstrumentation
|
||||
thrift=org.apache.skywalking.apm.plugin.thrift.define.transport.TSocketInstrumentation
|
||||
|
|
|
|||
Loading…
Reference in New Issue