diff --git a/apm-protocol/apm-network/src/main/proto/Common.proto b/apm-protocol/apm-network/src/main/proto/Common.proto new file mode 100644 index 000000000..9142490f3 --- /dev/null +++ b/apm-protocol/apm-network/src/main/proto/Common.proto @@ -0,0 +1,10 @@ +syntax = "proto3"; + +option java_multiple_files = true; +option java_package = "org.apache.skywalking.apm.network.proto"; + +enum SpanType { + Entry = 0; + Exit = 1; + Local = 2; +} \ No newline at end of file diff --git a/apm-protocol/apm-network/src/main/proto/DiscoveryService.proto b/apm-protocol/apm-network/src/main/proto/DiscoveryService.proto index f2068f027..c4c06ca5c 100644 --- a/apm-protocol/apm-network/src/main/proto/DiscoveryService.proto +++ b/apm-protocol/apm-network/src/main/proto/DiscoveryService.proto @@ -3,6 +3,7 @@ syntax = "proto3"; option java_multiple_files = true; option java_package = "org.apache.skywalking.apm.network.proto"; +import "Common.proto"; import "Downstream.proto"; service InstanceDiscoveryService { @@ -67,4 +68,5 @@ message ServiceNameMappingElement { message ServiceNameElement { string serviceName = 1; int32 applicationId = 2; + SpanType srcSpanType = 3; } diff --git a/apm-protocol/apm-network/src/main/proto/TraceSegmentService.proto b/apm-protocol/apm-network/src/main/proto/TraceSegmentService.proto index 5ac1aba91..f58fbbf3f 100644 --- a/apm-protocol/apm-network/src/main/proto/TraceSegmentService.proto +++ b/apm-protocol/apm-network/src/main/proto/TraceSegmentService.proto @@ -3,6 +3,7 @@ syntax = "proto3"; option java_multiple_files = true; option java_package = "org.apache.skywalking.apm.network.proto"; +import "Common.proto"; import "Downstream.proto"; import "KeyWithStringValue.proto"; @@ -66,12 +67,6 @@ enum RefType { CrossThread = 1; } -enum SpanType { - Entry = 0; - Exit = 1; - Local = 2; -} - enum SpanLayer { Unknown = 0; Database = 1; diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java index e1ca40ab0..33d189958 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/boot/ServiceManager.java @@ -19,8 +19,9 @@ package org.apache.skywalking.apm.agent.core.boot; -import java.util.HashMap; +import java.util.Collections; import java.util.Iterator; +import java.util.LinkedHashMap; import java.util.Map; import java.util.ServiceLoader; import org.apache.skywalking.apm.agent.core.logging.api.LogManager; @@ -36,7 +37,7 @@ public enum ServiceManager { INSTANCE; private static final ILog logger = LogManager.getLogger(ServiceManager.class); - private Map bootedServices = new HashMap(); + private Map bootedServices = Collections.emptyMap(); public void boot() { bootedServices = loadAllServices(); @@ -57,7 +58,7 @@ public enum ServiceManager { } private Map loadAllServices() { - HashMap bootedServices = new HashMap(); + Map bootedServices = new LinkedHashMap(); Iterator serviceIterator = load().iterator(); while (serviceIterator.hasNext()) { BootService bootService = serviceIterator.next(); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java index 03a1900a7..f848502b7 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/TracingContext.java @@ -296,7 +296,7 @@ public class TracingContext implements AbstractTracerContext { AbstractSpan parentSpan = peek(); final int parentSpanId = parentSpan == null ? -1 : parentSpan.getSpanId(); AbstractTracingSpan span = (AbstractTracingSpan)DictionaryManager.findOperationNameCodeSection() - .findOrPrepare4Register(segment.getApplicationId(), operationName) + .findOrPrepare4Register(segment.getApplicationId(), operationName, false, false) .doInCondition(new PossibleFound.FoundAndObtain() { @Override public Object doProcess(int operationId) { diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/StackBasedTracingSpan.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/StackBasedTracingSpan.java index 22a21fd8b..4ee8957c8 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/StackBasedTracingSpan.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/context/trace/StackBasedTracingSpan.java @@ -16,7 +16,6 @@ * */ - package org.apache.skywalking.apm.agent.core.context.trace; import org.apache.skywalking.apm.agent.core.dictionary.DictionaryManager; @@ -48,7 +47,7 @@ public abstract class StackBasedTracingSpan extends AbstractTracingSpan { if (--stackDepth == 0) { if (this.operationId == DictionaryUtil.nullValue()) { this.operationId = (Integer)DictionaryManager.findOperationNameCodeSection() - .findOrPrepare4Register(owner.getApplicationId(), operationName) + .findOrPrepare4Register(owner.getApplicationId(), operationName, this.isEntry(), this.isExit()) .doInCondition( new PossibleFound.FoundAndObtain() { @Override public Object doProcess(int value) { diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/dictionary/OperationNameDictionary.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/dictionary/OperationNameDictionary.java index 18c8b35e6..7075f2860 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/dictionary/OperationNameDictionary.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/dictionary/OperationNameDictionary.java @@ -16,7 +16,6 @@ * */ - package org.apache.skywalking.apm.agent.core.dictionary; import io.netty.util.internal.ConcurrentSet; @@ -28,6 +27,7 @@ import org.apache.skywalking.apm.network.proto.ServiceNameDiscoveryServiceGrpc; import org.apache.skywalking.apm.network.proto.ServiceNameElement; import org.apache.skywalking.apm.network.proto.ServiceNameMappingCollection; import org.apache.skywalking.apm.network.proto.ServiceNameMappingElement; +import org.apache.skywalking.apm.network.proto.SpanType; import static org.apache.skywalking.apm.agent.core.conf.Config.Dictionary.OPERATION_NAME_BUFFER_SIZE; @@ -39,19 +39,21 @@ public enum OperationNameDictionary { private Map operationNameDictionary = new ConcurrentHashMap(); private Set unRegisterOperationNames = new ConcurrentSet(); - public PossibleFound findOrPrepare4Register(int applicationId, String operationName) { - return find0(applicationId, operationName, true); + public PossibleFound findOrPrepare4Register(int applicationId, String operationName, + boolean isEntry, boolean isExit) { + return find0(applicationId, operationName, isEntry, isExit, true); } public PossibleFound findOnly(int applicationId, String operationName) { - return find0(applicationId, operationName, false); + return find0(applicationId, operationName, false, false, false); } - private PossibleFound find0(int applicationId, String operationName, boolean registerWhenNotFound) { + private PossibleFound find0(int applicationId, String operationName, + boolean isEntry, boolean isExit, boolean registerWhenNotFound) { if (operationName == null || operationName.length() == 0) { return new NotFound(); } - OperationNameKey key = new OperationNameKey(applicationId, operationName); + OperationNameKey key = new OperationNameKey(applicationId, operationName, isEntry, isExit); Integer operationId = operationNameDictionary.get(key); if (operationId != null) { return new Found(operationId); @@ -78,9 +80,12 @@ public enum OperationNameDictionary { ServiceNameMappingCollection serviceNameMappingCollection = serviceNameDiscoveryServiceBlockingStub.discovery(builder.build()); if (serviceNameMappingCollection.getElementsCount() > 0) { for (ServiceNameMappingElement serviceNameMappingElement : serviceNameMappingCollection.getElementsList()) { + ServiceNameElement element = serviceNameMappingElement.getElement(); OperationNameKey key = new OperationNameKey( - serviceNameMappingElement.getElement().getApplicationId(), - serviceNameMappingElement.getElement().getServiceName()); + element.getApplicationId(), + element.getServiceName(), + SpanType.Entry.equals(element.getSrcSpanType()), + SpanType.Exit.equals(element.getSrcSpanType())); unRegisterOperationNames.remove(key); operationNameDictionary.put(key, serviceNameMappingElement.getServiceId()); } @@ -91,10 +96,14 @@ public enum OperationNameDictionary { private class OperationNameKey { private int applicationId; private String operationName; + private boolean isEntry; + private boolean isExit; - public OperationNameKey(int applicationId, String operationName) { + public OperationNameKey(int applicationId, String operationName, boolean isEntry, boolean isExit) { this.applicationId = applicationId; this.operationName = operationName; + this.isEntry = isEntry; + this.isExit = isExit; } public int getApplicationId() { @@ -123,5 +132,13 @@ public enum OperationNameDictionary { result = 31 * result + operationName.hashCode(); return result; } + + boolean isEntry() { + return isEntry; + } + + boolean isExit() { + return isExit; + } } } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java index f7910f71d..b5d4f442e 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/apache/skywalking/apm/agent/core/remote/CollectorDiscoveryService.java @@ -44,13 +44,15 @@ public class CollectorDiscoveryService implements BootService { @Override public void boot() throws Throwable { + DiscoveryRestServiceClient discoveryRestServiceClient = new DiscoveryRestServiceClient(); + discoveryRestServiceClient.run(); future = Executors.newSingleThreadScheduledExecutor(new DefaultNamedThreadFactory("CollectorDiscoveryService")) - .scheduleAtFixedRate(new RunnableWithExceptionProtection(new DiscoveryRestServiceClient(), + .scheduleAtFixedRate(new RunnableWithExceptionProtection(discoveryRestServiceClient, new RunnableWithExceptionProtection.CallbackWhenException() { @Override public void handle(Throwable t) { logger.error("unexpected exception.", t); } - }), 0, + }), Config.Collector.DISCOVERY_CHECK_INTERVAL, Config.Collector.DISCOVERY_CHECK_INTERVAL, TimeUnit.SECONDS); } diff --git a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/CallbackInstrumentation.java b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/CallbackInstrumentation.java index e478d6204..1744ce702 100644 --- a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/CallbackInstrumentation.java +++ b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/CallbackInstrumentation.java @@ -39,7 +39,7 @@ public class CallbackInstrumentation extends ClassInstanceMethodsEnhancePluginDe public static final String ENHANCE_CLASS = "org.apache.kafka.clients.producer.Callback"; public static final String ENHANCE_METHOD = "onCompletion"; - public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v11.CallbackInterceptor"; + public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.CallbackInterceptor"; @Override protected ConstructorInterceptPoint[] getConstructorsInterceptPoints() { return new ConstructorInterceptPoint[0]; diff --git a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaConsumerInstrumentation.java b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaConsumerInstrumentation.java index e4e1445f6..878a8302c 100644 --- a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaConsumerInstrumentation.java +++ b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaConsumerInstrumentation.java @@ -48,12 +48,13 @@ import static org.apache.skywalking.apm.agent.core.plugin.match.NameMatch.byName public class KafkaConsumerInstrumentation extends ClassInstanceMethodsEnhancePluginDefine { public static final String CONSTRUCTOR_INTERCEPT_TYPE = "org.apache.kafka.clients.consumer.ConsumerConfig"; - public static final String CONSTRUCTOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v11.ConsumerConstructorInterceptor"; - public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v11.KafkaConsumerInterceptor"; + public static final String CONSTRUCTOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.ConsumerConstructorInterceptor"; + public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.KafkaConsumerInterceptor"; public static final String ENHANCE_METHOD = "pollOnce"; public static final String ENHANCE_CLASS = "org.apache.kafka.clients.consumer.KafkaConsumer"; public static final String SUBSCRIBE_METHOD = "subscribe"; public static final String SUBSCRIBE_INTERCEPT_TYPE = "org.apache.kafka.clients.consumer.ConsumerRebalanceListener"; + public static final String SUBSCRIBE_INTERCEPT_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.SubscribeMethodInterceptor"; @Override protected ConstructorInterceptPoint[] getConstructorsInterceptPoints() { return new ConstructorInterceptPoint[] { @@ -90,7 +91,7 @@ public class KafkaConsumerInstrumentation extends ClassInstanceMethodsEnhancePlu } @Override public String getMethodsInterceptor() { - return "org.apache.skywalking.apm.plugin.kafka.v11.SubscribeMethodInterceptor"; + return SUBSCRIBE_INTERCEPT_CLASS; } @Override public boolean isOverrideArgs() { diff --git a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaProducerInstrumentation.java b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaProducerInstrumentation.java index 54932a8af..f93e936c6 100644 --- a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaProducerInstrumentation.java +++ b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/KafkaProducerInstrumentation.java @@ -48,9 +48,9 @@ import static org.apache.skywalking.apm.agent.core.plugin.match.NameMatch.byName */ public class KafkaProducerInstrumentation extends ClassInstanceMethodsEnhancePluginDefine { - public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v11.KafkaProducerInterceptor"; + public static final String INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.KafkaProducerInterceptor"; public static final String ENHANCE_CLASS = "org.apache.kafka.clients.producer.KafkaProducer"; - public static final String CONSTRUCTOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v11.ProducerConstructorInterceptor"; + public static final String CONSTRUCTOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.ProducerConstructorInterceptor"; public static final String CONSTRUCTOR_INTERCEPTOR_FLAG = "org.apache.kafka.clients.producer.ProducerConfig"; @Override protected ConstructorInterceptPoint[] getConstructorsInterceptPoints() { diff --git a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/ProducerRecordInstrumentation.java b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/ProducerRecordInstrumentation.java index fe4d8b203..e2035b29f 100644 --- a/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/ProducerRecordInstrumentation.java +++ b/apm-sniffer/apm-sdk-plugin/kafka-v1-plugin/src/main/java/org/apache/skywalking/apm/plugin/kafka/v1/define/ProducerRecordInstrumentation.java @@ -39,7 +39,7 @@ import static org.apache.skywalking.apm.agent.core.plugin.match.NameMatch.byName */ public class ProducerRecordInstrumentation extends ClassInstanceMethodsEnhancePluginDefine { - public static final String CONSTRUCTOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v11.ProducerRecordConstructorInterceptor"; + public static final String CONSTRUCTOR_INTERCEPTOR_CLASS = "org.apache.skywalking.apm.plugin.kafka.v1.ProducerRecordConstructorInterceptor"; public static final String ENHANCE_CLASS = "org.apache.kafka.clients.producer.ProducerRecord"; @Override protected ConstructorInterceptPoint[] getConstructorsInterceptPoints() { diff --git a/docs/Supported-list.md b/docs/Supported-list.md index b1187e394..85dc2fcdf 100644 --- a/docs/Supported-list.md +++ b/docs/Supported-list.md @@ -27,7 +27,7 @@ * [Dubbo](https://github.com/alibaba/dubbo) 2.5.4 -> 2.6.0 * [Dubbox](https://github.com/dangdangdotcom/dubbox) 2.8.4 * [Motan](https://github.com/weibocom/motan) 0.2.x -> 1.1.0 - * [gRPC](https://github.com/grpc/grpc-java) 1.6+ + * [gRPC](https://github.com/grpc/grpc-java) 1.x * MQ * [RocketMQ](https://github.com/apache/rocketmq) 4.x * [Kafka](http://kafka.apache.org) 0.11.0.0 -> 1.0