diff --git a/dist-material/release-docs/LICENSE b/dist-material/release-docs/LICENSE index 91f887319..6175b968d 100755 --- a/dist-material/release-docs/LICENSE +++ b/dist-material/release-docs/LICENSE @@ -267,7 +267,8 @@ The text of each license is the standard Apache 2.0 license. Apache: commons-collections 3.2.2: https://github.com/apache/commons-collections, Apache 2.0 Apache: commons-configuration 1.8: https://github.com/apache/commons-configuration, Apache 2.0 Apache: commons-io 2.4: https://github.com/apache/commons-io, Apache 2.0 - Apache: commons-compress 1.18: https://github.com/apache/commons-compress, Apache 2.0 + Apache: commons-compress 1.19: https://github.com/apache/commons-compress, Apache 2.0 + Apache: commons-collections4 4.1: https://mvnrepository.com/artifact/org.apache.commons/commons-collections4, Apache 2.0 Apache: tomcat 8.5.27: https://github.com/apache/tomcat/tree/trunk, Apache 2.0 Apache: freemarker 2.3.28: https://github.com/apache/freemarker, Apache 2.0 netty 5.5.0: https://github.com/netty/netty/blob/4.1/LICENSE.txt, Apache 2.0 @@ -304,7 +305,7 @@ The text of each license is the standard Apache 2.0 license. HikariCP 3.1.0: https://github.com/brettwooldridge/HikariCP, Apache 2.0 zipkin 2.9.1: https://github.com/openzipkin/zipkin, Apache 2.0 sharding-jdbc-core 2.0.3: https://github.com/sharding-sphere/sharding-sphere, Apache 2.0 - kubernetes-client 4.0.0: https://github.com/kubernetes-client/java, Apache 2.0 + kubernetes-client 8.0.0: https://github.com/kubernetes-client/java, Apache 2.0 proto files from istio/istio: https://github.com/istio/istio Apache 2.0 proto files from istio/api: https://github.com/istio/api Apache 2.0 consul-client 1.2.6: https://github.com/rickfast/consul-client, Apache 2.0 @@ -327,6 +328,12 @@ The text of each license is the standard Apache 2.0 license. moshi 1.5.0: https://github.com/square/moshi, Apache 2.0 logging-interceptor 3.13.1: https://github.com/square/okhttp/tree/master/okhttp-logging-interceptor, Apache 2.0 msgpack-core 0.8.16: https://github.com/msgpack/msgpack-java, Apache 2.0 + sundr-codegen 0.2.10: https://mvnrepository.com/artifact/io.sundr/sundr-codegen, Apache 2.0 + sundr-core 0.2.10: https://mvnrepository.com/artifact/io.sundr/sundr-core, Apache 2.0 + swagger-annotations 1.5.22: https://mvnrepository.com/artifact/io.swagger.core.v3/swagger-annotations, Apache 2.0 + resourcecify-annotations 0.21.0: https://mvnrepository.com/artifact/io.sundr/resourcecify-annotations, Apache 2.0 + jose4j 0.7.0: https://mvnrepository.com/artifact/org.bitbucket.b_c/jose4j, Apache 2.0 + converter-moshi 2.5.0: https://mvnrepository.com/artifact/com.squareup.retrofit2/converter-moshi, Apache 2.0 vavr 0.10.3: https://github.com/vavr-io/vavr, Apache 2.0 ======================================================================== @@ -341,8 +348,9 @@ The text of each license is also included at licenses/LICENSE-[project].txt. GraphQL java 6.0: https://github.com/graphql-java/graphql-java , MIT GraphQL Java Tools 4.3.0: https://github.com/graphql-java/graphql-java-tools , MIT jopt-simple 5.0.2: https://github.com/jopt-simple/jopt-simple , MIT - bcpkix-jdk15on 1.55: http://www.bouncycastle.org/licence.html , MIT - bcprov-jdk15on 1.55: http://www.bouncycastle.org/licence.html , MIT + bcpkix-jdk15on 1.61: http://www.bouncycastle.org/licence.html , MIT + bcprov-jdk15on 1.61: http://www.bouncycastle.org/licence.html , MIT + bcprov-ext-jdk15on 1.61: http://www.bouncycastle.org/licence.html , MIT minimal-json 0.9.5: https://github.com/ralfstx/minimal-json, MIT checker-qual 2.8.1: https://github.com/typetools/checker-framework, MIT influxdb-java 2.15: https://github.com/influxdata/influxdb-java, MIT @@ -429,6 +437,7 @@ popper.js 1.14.7: https://github.com/FezVrasta/popper.js MIT vue-datepicker-local 1.0.19: https://github.com/weifeiyue/vue-datepicker-local MIT vue-js-modal 1.3.31: https://github.com/euvl/vue-js-modal MIT lodash 4.17.15: https://github.com/lodash/lodash MIT +gson-fire 1.8.3: https://mvnrepository.com/artifact/io.gsonfire/gson-fire MIT ======================================== Apache 2.0 licenses diff --git a/oap-server/pom.xml b/oap-server/pom.xml index 3c6faf9f0..f7fbf6a1b 100755 --- a/oap-server/pom.xml +++ b/oap-server/pom.xml @@ -65,7 +65,7 @@ 2.6 6.3.2 2.10.5 - 4.0.0 + 8.0.0 3.1.0 2.9.1 2.6.2 diff --git a/oap-server/server-bootstrap/src/main/resources/application.yml b/oap-server/server-bootstrap/src/main/resources/application.yml index 3761bef03..cf91f6124 100755 --- a/oap-server/server-bootstrap/src/main/resources/application.yml +++ b/oap-server/server-bootstrap/src/main/resources/application.yml @@ -29,7 +29,6 @@ cluster: schema: ${SW_ZK_SCHEMA:digest} # only support digest schema expression: ${SW_ZK_EXPRESSION:skywalking:skywalking} kubernetes: - watchTimeoutSeconds: ${SW_CLUSTER_K8S_WATCH_TIMEOUT:60} namespace: ${SW_CLUSTER_K8S_NAMESPACE:default} labelSelector: ${SW_CLUSTER_K8S_LABEL:app=collector,release=skywalking} uidEnvName: ${SW_CLUSTER_K8S_UID:SKYWALKING_COLLECTOR_UID} diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesConfig.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesConfig.java index 986f3c185..739bb6978 100644 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesConfig.java +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesConfig.java @@ -18,46 +18,17 @@ package org.apache.skywalking.oap.server.cluster.plugin.kubernetes; +import lombok.Getter; +import lombok.Setter; import org.apache.skywalking.oap.server.library.module.ModuleConfig; /** * The configuration of the module of cluster.kubernetes */ +@Getter +@Setter public class ClusterModuleKubernetesConfig extends ModuleConfig { - private int watchTimeoutSeconds; private String namespace; private String labelSelector; private String uidEnvName; - - public int getWatchTimeoutSeconds() { - return watchTimeoutSeconds; - } - - public void setWatchTimeoutSeconds(int watchTimeoutSeconds) { - this.watchTimeoutSeconds = watchTimeoutSeconds; - } - - public String getNamespace() { - return namespace; - } - - public void setNamespace(String namespace) { - this.namespace = namespace; - } - - public String getLabelSelector() { - return labelSelector; - } - - public void setLabelSelector(String labelSelector) { - this.labelSelector = labelSelector; - } - - public String getUidEnvName() { - return uidEnvName; - } - - public void setUidEnvName(String uidEnvName) { - this.uidEnvName = uidEnvName; - } } diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProvider.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProvider.java index bb49ccd8b..0c8807f95 100644 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProvider.java +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProvider.java @@ -18,8 +18,7 @@ package org.apache.skywalking.oap.server.cluster.plugin.kubernetes; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.dependencies.NamespacedPodListWatch; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.dependencies.UidEnvSupplier; +import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.cluster.ClusterModule; import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery; import org.apache.skywalking.oap.server.core.cluster.ClusterRegister; @@ -58,24 +57,23 @@ public class ClusterModuleKubernetesProvider extends ModuleProvider { @Override public void prepare() throws ServiceNotProvidedException { - coordinator = new KubernetesCoordinator(getManager(), new NamespacedPodListWatch(config.getNamespace(), config.getLabelSelector(), config - .getWatchTimeoutSeconds()), new UidEnvSupplier(config.getUidEnvName())); + + coordinator = new KubernetesCoordinator(getManager(), config); this.registerServiceImplementation(ClusterRegister.class, coordinator); this.registerServiceImplementation(ClusterNodesQuery.class, coordinator); } @Override public void start() { - + NamespacedPodListInformer.INFORMER.init(config); } @Override public void notifyAfterCompleted() { - coordinator.start(); } @Override public String[] requiredModules() { - return new String[0]; + return new String[] {CoreModule.NAME}; } } diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/Event.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/Event.java deleted file mode 100644 index 103c826a6..000000000 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/Event.java +++ /dev/null @@ -1,46 +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.oap.server.cluster.plugin.kubernetes; - -/** - * The event of watch. - */ -public class Event { - private final String type; - private final String uid; - private final String host; - - public Event(final String type, final String uid, final String host) { - this.type = type; - this.uid = uid; - this.host = host; - } - - String getType() { - return type; - } - - String getUid() { - return uid; - } - - String getHost() { - return host; - } -} diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinator.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinator.java index 9f3a8779b..7e2e425fb 100644 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinator.java +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinator.java @@ -18,21 +18,13 @@ package org.apache.skywalking.oap.server.cluster.plugin.kubernetes; -import com.google.common.util.concurrent.FutureCallback; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningExecutorService; -import com.google.common.util.concurrent.MoreExecutors; -import com.google.common.util.concurrent.ThreadFactoryBuilder; -import java.util.ArrayList; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1Pod; +import io.kubernetes.client.openapi.models.V1PodStatus; +import java.util.Collections; import java.util.List; -import java.util.Map; -import java.util.concurrent.Callable; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.function.Supplier; -import javax.annotation.Nullable; +import java.util.stream.Collectors; +import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery; import org.apache.skywalking.oap.server.core.cluster.ClusterRegister; @@ -42,108 +34,64 @@ import org.apache.skywalking.oap.server.core.config.ConfigService; import org.apache.skywalking.oap.server.core.remote.client.Address; import org.apache.skywalking.oap.server.library.module.ModuleDefineHolder; import org.apache.skywalking.oap.server.telemetry.api.TelemetryRelatedContext; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Read collector pod info from api-server of kubernetes, then using all containerIp list to construct the list of * {@link RemoteInstance}. */ +@Slf4j public class KubernetesCoordinator implements ClusterRegister, ClusterNodesQuery { - private static final Logger logger = LoggerFactory.getLogger(KubernetesCoordinator.class); - private final ModuleDefineHolder manager; - private final String uid; - - private final Map cache = new ConcurrentHashMap<>(); - - private final ReusableWatch watch; - private volatile int port = -1; - KubernetesCoordinator(ModuleDefineHolder manager, final ReusableWatch watch, - final Supplier uidSupplier) { + private final String uid; + + public KubernetesCoordinator(final ModuleDefineHolder manager, + final ClusterModuleKubernetesConfig config) { + this.uid = new UidEnvSupplier(config.getUidEnvName()).get(); this.manager = manager; - this.watch = watch; - this.uid = uidSupplier.get(); - TelemetryRelatedContext.INSTANCE.setId(uid); - } - - public void start() { - ExecutorService executorService = Executors.newSingleThreadExecutor(new ThreadFactoryBuilder().setDaemon(true) - .setNameFormat("Kubernetes-ApiServer-%s") - .build()); - submitTask(MoreExecutors.listeningDecorator(executorService), executorService); - } - - @Override - public void registerRemote(RemoteInstance remoteInstance) throws ServiceRegisterException { - this.port = remoteInstance.getAddress().getPort(); - } - - private void submitTask(final ListeningExecutorService service, final ExecutorService executorService) { - watch.initOrReset(); - - ListenableFuture watchFuture = service.submit(newWatch()); - Futures.addCallback(watchFuture, new FutureCallback() { - @Override - public void onSuccess(@Nullable Object ignored) { - submitTask(service, executorService); - } - - @Override - public void onFailure(@Nullable Throwable throwable) { - logger.debug("Generate remote nodes error", throwable); - submitTask(service, executorService); - } - }, executorService); - } - - private Callable newWatch() { - return () -> { - generateRemoteNodes(); - return null; - }; - } - - private void generateRemoteNodes() { - for (Event event : watch) { - if (event == null) { - break; - } - logger.debug("Received event {} {}-{}", event.getType(), event.getUid(), event.getHost()); - switch (event.getType()) { - case "ADDED": - case "MODIFIED": - cache.put(event.getUid(), new RemoteInstance(new Address(event.getHost(), port, event.getUid() - .equals(this.uid)))); - break; - case "DELETED": - cache.remove(event.getUid()); - break; - default: - throw new RuntimeException(String.format("Unknown event %s", event.getType())); - } - } } @Override public List queryRemoteNodes() { - final List list = new ArrayList<>(); - cache.values().forEach(instance -> { - Address address = instance.getAddress(); - if (port == -1) { - logger.debug("Query kubernetes remote, port hasn't init, try to init"); - ConfigService service = manager.find(CoreModule.NAME).provider().getService(ConfigService.class); - port = service.getGRPCPort(); - logger.debug("Query kubernetes remote, port is set at {}", port); - } - list.add(new RemoteInstance(new Address(address.getHost(), port, address.isSelf()))); - }); - logger.debug("Query kubernetes remote nodes: {}", list); - return list; + List pods = NamespacedPodListInformer.INFORMER.listPods().orElseGet(this::selfPod); + + if (log.isDebugEnabled()) { + List uidList = pods + .stream() + .map(item -> item.getMetadata().getUid()) + .collect(Collectors.toList()); + log.debug("[kubernetes cluster pods uid list]:{}", uidList.toString()); + } + + if (port == -1) { + port = manager.find(CoreModule.NAME).provider().getService(ConfigService.class).getGRPCPort(); + } + + return pods.stream() + .map(pod -> new RemoteInstance( + new Address(pod.getStatus().getPodIP(), port, pod.getMetadata().getUid().equals(uid)))) + .collect(Collectors.toList()); + + } + + @Override + public void registerRemote(final RemoteInstance remoteInstance) throws ServiceRegisterException { + this.port = remoteInstance.getAddress().getPort(); + TelemetryRelatedContext.INSTANCE.setId(remoteInstance.toString()); + } + + private List selfPod() { + + V1Pod v1Pod = new V1Pod(); + v1Pod.setMetadata(new V1ObjectMeta()); + v1Pod.setStatus(new V1PodStatus()); + v1Pod.getMetadata().setUid(uid); + v1Pod.getStatus().setPodIP("127.0.0.1"); + return Collections.singletonList(v1Pod); + } } diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/NamespacedPodListInformer.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/NamespacedPodListInformer.java new file mode 100644 index 000000000..0371a3b52 --- /dev/null +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/NamespacedPodListInformer.java @@ -0,0 +1,106 @@ +/* + * 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.oap.server.cluster.plugin.kubernetes; + +import io.kubernetes.client.informer.SharedIndexInformer; +import io.kubernetes.client.informer.SharedInformerFactory; +import io.kubernetes.client.informer.cache.Lister; +import io.kubernetes.client.openapi.ApiClient; +import io.kubernetes.client.openapi.apis.CoreV1Api; +import io.kubernetes.client.openapi.models.V1Pod; +import io.kubernetes.client.openapi.models.V1PodList; +import io.kubernetes.client.util.Config; +import java.io.IOException; +import java.util.List; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public enum NamespacedPodListInformer { + + /** + * contains remote collector instances + */ + INFORMER; + + private Lister podLister; + + private SharedInformerFactory factory; + + private final ExecutorService executorService = Executors.newSingleThreadExecutor(r -> { + Thread thread = new Thread(r, "SKYWALKING_KUBERNETES_CLUSTER_INFORMER"); + thread.setDaemon(true); + return thread; + }); + + { + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + if (Objects.nonNull(factory)) { + factory.stopAllRegisteredInformers(); + } + })); + } + + public synchronized void init(ClusterModuleKubernetesConfig podConfig) { + + try { + doStartPodInformer(podConfig); + } catch (IOException e) { + log.error("cannot connect with api server in kubernetes", e); + } + } + + private void doStartPodInformer(ClusterModuleKubernetesConfig podConfig) throws IOException { + + ApiClient apiClient = Config.defaultClient(); + apiClient.setHttpClient(apiClient.getHttpClient().newBuilder().readTimeout(0, TimeUnit.SECONDS).build()); + CoreV1Api coreV1Api = new CoreV1Api(apiClient); + factory = new SharedInformerFactory(executorService); + + SharedIndexInformer podSharedIndexInformer = factory.sharedIndexInformerFor( + params -> coreV1Api.listNamespacedPodCall( + podConfig.getNamespace(), null, null, null, null, + podConfig.getLabelSelector(), Integer.MAX_VALUE, params.resourceVersion, params.timeoutSeconds, + params.watch, null + ), + V1Pod.class, V1PodList.class + ); + + factory.startAllRegisteredInformers(); + podLister = new Lister<>(podSharedIndexInformer.getIndexer()); + } + + public Optional> listPods() { + + return Optional.ofNullable(podLister.list().size() != 0 + ? podLister.list() + .stream() + .filter( + item -> "Running".equalsIgnoreCase(item.getStatus().getPhase())) + .collect(Collectors.toList()) + : null); + + } + +} diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ReusableWatch.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ReusableWatch.java deleted file mode 100644 index 67a4cfaf9..000000000 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ReusableWatch.java +++ /dev/null @@ -1,32 +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.oap.server.cluster.plugin.kubernetes; - -/** - * This watch can init or reset internal state. - * - * @param event of watch - */ -public interface ReusableWatch extends Iterable { - - /** - * Reset internal state. - */ - void initOrReset(); -} diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/dependencies/UidEnvSupplier.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/UidEnvSupplier.java similarity index 98% rename from oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/dependencies/UidEnvSupplier.java rename to oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/UidEnvSupplier.java index d444c835f..95a903f1e 100644 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/dependencies/UidEnvSupplier.java +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/UidEnvSupplier.java @@ -16,7 +16,7 @@ * */ -package org.apache.skywalking.oap.server.cluster.plugin.kubernetes.dependencies; +package org.apache.skywalking.oap.server.cluster.plugin.kubernetes; import java.util.function.Supplier; diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/dependencies/NamespacedPodListWatch.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/dependencies/NamespacedPodListWatch.java deleted file mode 100644 index 51d1125d8..000000000 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/dependencies/NamespacedPodListWatch.java +++ /dev/null @@ -1,114 +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.oap.server.cluster.plugin.kubernetes.dependencies; - -import com.google.common.reflect.TypeToken; -import io.kubernetes.client.ApiClient; -import io.kubernetes.client.ApiException; -import io.kubernetes.client.Configuration; -import io.kubernetes.client.apis.CoreV1Api; -import io.kubernetes.client.models.V1Pod; -import io.kubernetes.client.util.Config; -import io.kubernetes.client.util.Watch; -import java.io.IOException; -import java.util.Iterator; -import java.util.Objects; -import java.util.concurrent.TimeUnit; -import java.util.function.Supplier; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.Event; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.ReusableWatch; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * Watch the api {@literal https://v1-9.docs.kubernetes.io/docs/reference/generated/kubernetes-api/v1.9/#watch-64}. - */ -public class NamespacedPodListWatch implements ReusableWatch { - - private static final Logger logger = LoggerFactory.getLogger(NamespacedPodListWatch.class); - - private final String namespace; - - private final String labelSelector; - - private final int watchTimeoutSeconds; - - private Watch watch; - - public NamespacedPodListWatch(final String namespace, final String labelSelector, final int watchTimeoutSeconds) { - this.namespace = namespace; - this.labelSelector = labelSelector; - this.watchTimeoutSeconds = watchTimeoutSeconds; - } - - @Override - public void initOrReset() { - ApiClient client; - try { - client = Config.defaultClient(); - } catch (IOException e) { - throw new RuntimeException(e.getMessage(), e); - } - client.getHttpClient().setReadTimeout(watchTimeoutSeconds, TimeUnit.SECONDS); - Configuration.setDefaultApiClient(client); - CoreV1Api api = new CoreV1Api(); - try { - watch = Watch.createWatch(client, api.listNamespacedPodCall(namespace, null, null, null, null, labelSelector, Integer.MAX_VALUE, null, null, Boolean.TRUE, null, null), new TypeToken>() { - }.getType()); - } catch (final ApiException e) { - logger.error("code:{} header:{} body:{}", e.getCode(), e.getResponseHeaders(), e.getResponseBody()); - throw new RuntimeException(e.getMessage(), e); - } - } - - @Override - public Iterator iterator() { - final Iterator> watchItr = watch.iterator(); - return new Iterator() { - @Override - public boolean hasNext() { - return wrap(watchItr::hasNext, false); - } - - @Override - public Event next() { - return wrap(() -> { - final Watch.Response response = watchItr.next(); - return new Event(response.type, response.object.getMetadata().getUid(), response.object.getStatus() - .getPodIP()); - }, null); - } - - private R wrap(final Supplier action, final R defaultValue) { - Objects.requireNonNull(action); - try { - return action.get(); - } catch (final Throwable t) { - logger.trace("Throwable", t); - try { - watch.close(); - } catch (IOException e) { - logger.error("Close watch error", e); - } - } - return defaultValue; - } - }; - } -} diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProviderTest.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProviderTest.java index 1ba4bed77..79f1029b6 100644 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProviderTest.java +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/ClusterModuleKubernetesProviderTest.java @@ -18,48 +18,53 @@ package org.apache.skywalking.oap.server.cluster.plugin.kubernetes; +import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.cluster.ClusterModule; -import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery; -import org.apache.skywalking.oap.server.core.cluster.ClusterRegister; -import org.apache.skywalking.oap.server.library.module.ServiceNotProvidedException; -import org.junit.Before; +import org.apache.skywalking.oap.server.library.module.ModuleConfig; import org.junit.Test; +import org.junit.runner.RunWith; +import org.powermock.core.classloader.annotations.PowerMockIgnore; +import org.powermock.modules.junit4.PowerMockRunner; -import static org.hamcrest.core.Is.is; -import static org.junit.Assert.assertSame; -import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; +@RunWith(PowerMockRunner.class) +@PowerMockIgnore("javax.management.*") public class ClusterModuleKubernetesProviderTest { - private ClusterModuleKubernetesProvider provider; + private ClusterModuleKubernetesProvider provider = new ClusterModuleKubernetesProvider(); - @Before - public void setUp() { - provider = new ClusterModuleKubernetesProvider(); + @Test + public void name() { + assertEquals("kubernetes", provider.name()); } @Test - public void assertName() { - assertThat(provider.name(), is("kubernetes")); + public void module() { + assertEquals(ClusterModule.class, provider.module()); } @Test - public void assertModule() { - assertTrue(provider.module().isAssignableFrom(ClusterModule.class)); + public void createConfigBeanIfAbsent() { + ModuleConfig moduleConfig = provider.createConfigBeanIfAbsent(); + assertTrue(moduleConfig instanceof ClusterModuleKubernetesConfig); } @Test - public void assertCreateConfigBeanIfAbsent() { - assertTrue(ClusterModuleKubernetesConfig.class.isInstance(provider.createConfigBeanIfAbsent())); - } - - @Test - public void assertPrepare() throws ServiceNotProvidedException { + public void prepare() throws Exception { provider.prepare(); - ClusterRegister register = provider.getService(ClusterRegister.class); - ClusterNodesQuery query = provider.getService(ClusterNodesQuery.class); - assertSame(register, query); - assertTrue(KubernetesCoordinator.class.isInstance(register)); + } + + @Test + public void notifyAfterCompleted() { + provider.notifyAfterCompleted(); + } + + @Test + public void requiredModules() { + String[] modules = provider.requiredModules(); + assertArrayEquals(new String[] {CoreModule.NAME}, modules); } } \ No newline at end of file diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinatorTest.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinatorTest.java index 81a87f42d..48a831983 100644 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinatorTest.java +++ b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/KubernetesCoordinatorTest.java @@ -18,150 +18,104 @@ package org.apache.skywalking.oap.server.cluster.plugin.kubernetes; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.fixture.PlainWatch; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1Pod; +import io.kubernetes.client.openapi.models.V1PodStatus; +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; +import java.util.stream.Collectors; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.CoreModuleConfig; import org.apache.skywalking.oap.server.core.cluster.RemoteInstance; import org.apache.skywalking.oap.server.core.config.ConfigService; import org.apache.skywalking.oap.server.core.remote.client.Address; -import org.apache.skywalking.oap.server.library.module.ModuleDefineHolder; import org.apache.skywalking.oap.server.testing.module.ModuleDefineTesting; import org.apache.skywalking.oap.server.testing.module.ModuleManagerTesting; +import org.junit.Assert; +import org.junit.Before; import org.junit.Test; -import org.mockito.Mockito; +import org.junit.runner.RunWith; +import org.powermock.api.mockito.PowerMockito; +import org.powermock.api.support.membermodification.MemberModifier; +import org.powermock.core.classloader.annotations.PowerMockIgnore; +import org.powermock.core.classloader.annotations.PrepareForTest; +import org.powermock.modules.junit4.PowerMockRunner; +import org.powermock.reflect.Whitebox; -import static org.hamcrest.core.Is.is; -import static org.junit.Assert.assertThat; -import static org.mockito.Mockito.when; +import static org.powermock.api.mockito.PowerMockito.when; +@RunWith(PowerMockRunner.class) +@PowerMockIgnore("javax.management.*") +@PrepareForTest({NamespacedPodListInformer.class}) public class KubernetesCoordinatorTest { private KubernetesCoordinator coordinator; - @Test - public void assertAdded() throws InterruptedException { - PlainWatch watch = PlainWatch.create(2, "ADDED", "1", "10.0.0.1", "ADDED", "2", "10.0.0.2"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - coordinator.registerRemote(new RemoteInstance(new Address("0.0.0.0", 8454, true))); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(2)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.1")); + public static final String LOCAL_HOST = "127.0.0.1"; + public static final Integer GRPC_PORT = 8454; + public static final Integer SELF_UID = 12345; + + private Address selfAddress; + private NamespacedPodListInformer informer; + + @Before + public void prepare() throws IllegalAccessException { + coordinator = new KubernetesCoordinator(getManager(), new ClusterModuleKubernetesConfig()); + MemberModifier.field(KubernetesCoordinator.class, "uid").set(coordinator, String.valueOf(SELF_UID)); + selfAddress = new Address(LOCAL_HOST, GRPC_PORT, true); + informer = PowerMockito.mock(NamespacedPodListInformer.class); + Whitebox.setInternalState(NamespacedPodListInformer.class, "INFORMER", informer); + } @Test - public void assertModified() throws InterruptedException { - PlainWatch watch = PlainWatch.create(3, "ADDED", "1", "10.0.0.1", "ADDED", "2", "10.0.0.2", "MODIFIED", "1", "10.0.0.3"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - coordinator.registerRemote(new RemoteInstance(new Address("0.0.0.0", 8454, true))); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(2)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.3")); + public void queryRemoteNodesWhenInformerNotwork() throws Exception { + PowerMockito.doReturn(Optional.empty()).when(NamespacedPodListInformer.INFORMER).listPods(); + List remoteInstances = Whitebox.invokeMethod(coordinator, "queryRemoteNodes"); + Assert.assertEquals(1, remoteInstances.size()); + Assert.assertEquals(selfAddress, remoteInstances.get(0).getAddress()); + } @Test - public void assertDeleted() throws InterruptedException { - PlainWatch watch = PlainWatch.create(3, "ADDED", "1", "10.0.0.1", "ADDED", "2", "10.0.0.2", "DELETED", "2", "10.0.0.2"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - coordinator.registerRemote(new RemoteInstance(new Address("0.0.0.0", 8454, true))); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(1)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.1")); + public void queryRemoteNodesWhenInformerWork() throws Exception { + PowerMockito.doReturn(Optional.of(mockPodList())).when(NamespacedPodListInformer.INFORMER).listPods(); + List remoteInstances = Whitebox.invokeMethod(coordinator, "queryRemoteNodes"); + Assert.assertEquals(5, remoteInstances.size()); + List self = remoteInstances.stream() + .filter(item -> item.getAddress().isSelf()) + .collect(Collectors.toList()); + List others = remoteInstances.stream() + .filter(item -> !item.getAddress().isSelf()) + .collect(Collectors.toList()); + + Assert.assertEquals(1, self.size()); + Assert.assertEquals(4, others.size()); + } - @Test - public void assertError() throws InterruptedException { - PlainWatch watch = PlainWatch.create(3, "ADDED", "1", "10.0.0.1", "ERROR", "X", "10.0.0.2", "ADDED", "2", "10.0.0.2"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - coordinator.registerRemote(new RemoteInstance(new Address("0.0.0.0", 8454, true))); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(2)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.1")); - } - - @Test - public void assertModifiedInReceiverRole() throws InterruptedException { - PlainWatch watch = PlainWatch.create(3, "ADDED", "1", "10.0.0.1", "ADDED", "2", "10.0.0.2", "MODIFIED", "1", "10.0.0.3"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(2)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.3")); - } - - @Test - public void assertDeletedInReceiverRole() throws InterruptedException { - PlainWatch watch = PlainWatch.create(3, "ADDED", "1", "10.0.0.1", "ADDED", "2", "10.0.0.2", "DELETED", "2", "10.0.0.2"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(1)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.1")); - } - - @Test - public void assertErrorInReceiverRole() throws InterruptedException { - PlainWatch watch = PlainWatch.create(3, "ADDED", "1", "10.0.0.1", "ERROR", "X", "10.0.0.2", "ADDED", "2", "10.0.0.2"); - coordinator = new KubernetesCoordinator(getManager(), watch, () -> "1"); - coordinator.start(); - watch.await(); - assertThat(coordinator.queryRemoteNodes().size(), is(2)); - assertThat(coordinator.queryRemoteNodes() - .stream() - .filter(instance -> instance.getAddress().isSelf()) - .findFirst() - .get() - .getAddress() - .getHost(), is("10.0.0.1")); - } - - public ModuleDefineHolder getManager() { + private ModuleManagerTesting getManager() { ModuleManagerTesting moduleManagerTesting = new ModuleManagerTesting(); ModuleDefineTesting coreModuleDefine = new ModuleDefineTesting(); moduleManagerTesting.put(CoreModule.NAME, coreModuleDefine); - CoreModuleConfig config = Mockito.mock(CoreModuleConfig.class); - when(config.getGRPCHost()).thenReturn("127.0.0.1"); - when(config.getGRPCPort()).thenReturn(8454); + CoreModuleConfig config = PowerMockito.mock(CoreModuleConfig.class); + when(config.getGRPCHost()).thenReturn(LOCAL_HOST); + when(config.getGRPCPort()).thenReturn(GRPC_PORT); coreModuleDefine.provider().registerServiceImplementation(ConfigService.class, new ConfigService(config)); return moduleManagerTesting; } -} \ No newline at end of file + + private List mockPodList() { + List pods = new ArrayList<>(); + for (int i = 0; i < 5; i++) { + V1Pod v1Pod = new V1Pod(); + v1Pod.setMetadata(new V1ObjectMeta()); + v1Pod.setStatus(new V1PodStatus()); + v1Pod.getMetadata().setUid(String.valueOf(SELF_UID + i)); + v1Pod.getStatus().setPodIP(LOCAL_HOST); + pods.add(v1Pod); + } + return pods; + } +} diff --git a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/fixture/PlainWatch.java b/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/fixture/PlainWatch.java deleted file mode 100644 index 137a79319..000000000 --- a/oap-server/server-cluster-plugin/cluster-kubernetes-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/kubernetes/fixture/PlainWatch.java +++ /dev/null @@ -1,89 +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.oap.server.cluster.plugin.kubernetes.fixture; - -import java.util.ArrayList; -import java.util.Iterator; -import java.util.List; -import java.util.NoSuchElementException; -import java.util.concurrent.CountDownLatch; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.Event; -import org.apache.skywalking.oap.server.cluster.plugin.kubernetes.ReusableWatch; - -public class PlainWatch implements ReusableWatch { - - public static PlainWatch create(final int size, final String... args) { - List events = new ArrayList<>(args.length / 3); - for (int i = 0; i < args.length; i++) { - events.add(new Event(args[i++], args[i++], args[i])); - } - return new PlainWatch(events, size); - } - - private final List events; - - private final int size; - - private final CountDownLatch latch = new CountDownLatch(1); - - private Iterator iterator; - - private int count; - - private PlainWatch(final List events, final int size) { - this.events = events; - this.size = size; - } - - @Override - public void initOrReset() { - final Iterator internal = events.subList(count, events.size()).iterator(); - iterator = new Iterator() { - public boolean hasNext() { - boolean result = count < size && internal.hasNext(); - if (!result) { - latch.countDown(); - } - return result; - } - - public Event next() { - if (!this.hasNext()) { - throw new NoSuchElementException(); - } else { - ++count; - return internal.next(); - } - } - - public void remove() { - internal.remove(); - } - }; - } - - @Override - public Iterator iterator() { - return iterator; - } - - public void await() throws InterruptedException { - latch.await(); - } -} diff --git a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResource.java b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResource.java index e65399080..39b076f2f 100644 --- a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResource.java +++ b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResource.java @@ -19,14 +19,13 @@ package org.apache.skywalking.oap.server.receiver.envoy.als; -import io.kubernetes.client.models.V1ObjectMeta; -import io.kubernetes.client.models.V1OwnerReference; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1OwnerReference; +import java.util.Optional; import lombok.AccessLevel; import lombok.Getter; import lombok.RequiredArgsConstructor; -import java.util.Optional; - @RequiredArgsConstructor class DependencyResource { @Getter(AccessLevel.PACKAGE) diff --git a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/Fetcher.java b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/Fetcher.java index 7efcc8a83..1d40d7cfd 100644 --- a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/Fetcher.java +++ b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/Fetcher.java @@ -19,14 +19,13 @@ package org.apache.skywalking.oap.server.receiver.envoy.als; -import io.kubernetes.client.ApiException; -import io.kubernetes.client.models.V1ObjectMeta; -import io.kubernetes.client.models.V1OwnerReference; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - +import io.kubernetes.client.openapi.ApiException; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1OwnerReference; import java.util.Optional; import java.util.function.Function; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; interface Fetcher extends Function> { diff --git a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/K8sALSServiceMeshHTTPAnalysis.java b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/K8sALSServiceMeshHTTPAnalysis.java index caeca40d0..b1d8d9835 100644 --- a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/K8sALSServiceMeshHTTPAnalysis.java +++ b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/envoy/als/K8sALSServiceMeshHTTPAnalysis.java @@ -29,14 +29,13 @@ import io.envoyproxy.envoy.data.accesslog.v2.HTTPAccessLogEntry; import io.envoyproxy.envoy.data.accesslog.v2.HTTPRequestProperties; import io.envoyproxy.envoy.data.accesslog.v2.HTTPResponseProperties; import io.envoyproxy.envoy.service.accesslog.v2.StreamAccessLogsMessage; -import io.kubernetes.client.ApiClient; -import io.kubernetes.client.Configuration; -import io.kubernetes.client.apis.CoreV1Api; -import io.kubernetes.client.apis.ExtensionsV1beta1Api; -import io.kubernetes.client.models.V1ObjectMeta; -import io.kubernetes.client.models.V1OwnerReference; -import io.kubernetes.client.models.V1Pod; -import io.kubernetes.client.models.V1PodList; +import io.kubernetes.client.openapi.ApiClient; +import io.kubernetes.client.openapi.apis.CoreV1Api; +import io.kubernetes.client.openapi.apis.ExtensionsV1beta1Api; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1OwnerReference; +import io.kubernetes.client.openapi.models.V1Pod; +import io.kubernetes.client.openapi.models.V1PodList; import io.kubernetes.client.util.Config; import java.time.Instant; import java.util.ArrayList; @@ -72,10 +71,11 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { @Getter(AccessLevel.PROTECTED) private final AtomicReference> ipServiceMap = new AtomicReference<>(); - private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1, new ThreadFactoryBuilder() - .setNameFormat("load-pod-%d") - .setDaemon(true) - .build()); + private final ScheduledExecutorService executorService = Executors.newScheduledThreadPool( + 1, new ThreadFactoryBuilder() + .setNameFormat("load-pod-%d") + .setDaemon(true) + .build()); @Override public String name() { @@ -95,9 +95,8 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { private void loadPodInfo() { try { ApiClient client = Config.defaultClient(); - client.getHttpClient().setReadTimeout(20, TimeUnit.SECONDS); - Configuration.setDefaultApiClient(client); - CoreV1Api api = new CoreV1Api(); + CoreV1Api api = new CoreV1Api(client); + V1PodList list = api.listPodForAllNamespaces(null, null, null, null, null, null, null, null, null); Map ipMap = new HashMap<>(list.getItems().size()); long startTime = System.nanoTime(); @@ -109,9 +108,12 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { continue; } if (item.getStatus().getPodIP().equals(item.getStatus().getHostIP())) { - logger.debug("Pod {}.{} is removed because hostIP and podIP are identical ", item.getMetadata() - .getName(), item.getMetadata() - .getNamespace()); + logger.debug( + "Pod {}.{} is removed because hostIP and podIP are identical ", item.getMetadata() + .getName(), + item.getMetadata() + .getNamespace() + ); continue; } ipMap.put(item.getStatus().getPodIP(), createServiceMetaInfo(item.getMetadata())); @@ -126,8 +128,11 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { private ServiceMetaInfo createServiceMetaInfo(final V1ObjectMeta podMeta) { ExtensionsV1beta1Api extensionsApi = new ExtensionsV1beta1Api(); DependencyResource dr = new DependencyResource(podMeta); - DependencyResource meta = dr.getOwnerResource("ReplicaSet", ownerReference -> extensionsApi.readNamespacedReplicaSet(ownerReference - .getName(), podMeta.getNamespace(), "", true, true).getMetadata()); + DependencyResource meta = dr.getOwnerResource( + "ReplicaSet", ownerReference -> extensionsApi.readNamespacedReplicaSet( + ownerReference + .getName(), podMeta.getNamespace(), "", true, true) + .getMetadata()); ServiceMetaInfo result = new ServiceMetaInfo(); if (meta.getMetadata().getOwnerReferences() != null && meta.getMetadata().getOwnerReferences().size() > 0) { V1OwnerReference owner = meta.getMetadata().getOwnerReferences().get(0); @@ -197,47 +202,57 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { boolean status = responseCode >= 200 && responseCode < 400; Address downstreamRemoteAddress = properties.getDownstreamRemoteAddress(); - ServiceMetaInfo downstreamService = find(downstreamRemoteAddress.getSocketAddress() - .getAddress(), downstreamRemoteAddress.getSocketAddress() - .getPortValue()); + ServiceMetaInfo downstreamService = find( + downstreamRemoteAddress.getSocketAddress() + .getAddress(), downstreamRemoteAddress.getSocketAddress() + .getPortValue()); Address downstreamLocalAddress = properties.getDownstreamLocalAddress(); - ServiceMetaInfo localService = find(downstreamLocalAddress.getSocketAddress() - .getAddress(), downstreamLocalAddress.getSocketAddress() - .getPortValue()); + ServiceMetaInfo localService = find( + downstreamLocalAddress.getSocketAddress() + .getAddress(), downstreamLocalAddress.getSocketAddress() + .getPortValue()); if (cluster.startsWith("inbound|")) { // Server side if (downstreamService.equals(ServiceMetaInfo.UNKNOWN)) { // Ingress -> sidecar(server side) // Mesh telemetry without source, the relation would be generated. ServiceMeshMetric.Builder metric = ServiceMeshMetric.newBuilder() - .setStartTime(startTime) - .setEndTime(startTime + duration) - .setDestServiceName(localService.getServiceName()) - .setDestServiceInstance(localService.getServiceInstanceName()) - .setEndpoint(endpoint) - .setLatency((int) duration) - .setResponseCode(Math.toIntExact(responseCode)) - .setStatus(status) - .setProtocol(protocol) - .setDetectPoint(DetectPoint.server); + .setStartTime(startTime) + .setEndTime(startTime + duration) + .setDestServiceName( + localService.getServiceName()) + .setDestServiceInstance( + localService.getServiceInstanceName()) + .setEndpoint(endpoint) + .setLatency((int) duration) + .setResponseCode( + Math.toIntExact(responseCode)) + .setStatus(status) + .setProtocol(protocol) + .setDetectPoint(DetectPoint.server); logger.debug("Transformed ingress->sidecar inbound mesh metric {}", metric); forward(metric); } else { // sidecar -> sidecar(server side) ServiceMeshMetric.Builder metric = ServiceMeshMetric.newBuilder() - .setStartTime(startTime) - .setEndTime(startTime + duration) - .setSourceServiceName(downstreamService.getServiceName()) - .setSourceServiceInstance(downstreamService.getServiceInstanceName()) - .setDestServiceName(localService.getServiceName()) - .setDestServiceInstance(localService.getServiceInstanceName()) - .setEndpoint(endpoint) - .setLatency((int) duration) - .setResponseCode(Math.toIntExact(responseCode)) - .setStatus(status) - .setProtocol(protocol) - .setDetectPoint(DetectPoint.server); + .setStartTime(startTime) + .setEndTime(startTime + duration) + .setSourceServiceName( + downstreamService.getServiceName()) + .setSourceServiceInstance( + downstreamService.getServiceInstanceName()) + .setDestServiceName( + localService.getServiceName()) + .setDestServiceInstance( + localService.getServiceInstanceName()) + .setEndpoint(endpoint) + .setLatency((int) duration) + .setResponseCode( + Math.toIntExact(responseCode)) + .setStatus(status) + .setProtocol(protocol) + .setDetectPoint(DetectPoint.server); logger.debug("Transformed sidecar->sidecar(server side) inbound mesh metric {}", metric); forward(metric); @@ -245,23 +260,28 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { } else if (cluster.startsWith("outbound|")) { // sidecar(client side) -> sidecar Address upstreamRemoteAddress = properties.getUpstreamRemoteAddress(); - ServiceMetaInfo destService = find(upstreamRemoteAddress.getSocketAddress() - .getAddress(), upstreamRemoteAddress.getSocketAddress() - .getPortValue()); + ServiceMetaInfo destService = find( + upstreamRemoteAddress.getSocketAddress() + .getAddress(), upstreamRemoteAddress.getSocketAddress() + .getPortValue()); ServiceMeshMetric.Builder metric = ServiceMeshMetric.newBuilder() - .setStartTime(startTime) - .setEndTime(startTime + duration) - .setSourceServiceName(downstreamService.getServiceName()) - .setSourceServiceInstance(downstreamService.getServiceInstanceName()) - .setDestServiceName(destService.getServiceName()) - .setDestServiceInstance(destService.getServiceInstanceName()) - .setEndpoint(endpoint) - .setLatency((int) duration) - .setResponseCode(Math.toIntExact(responseCode)) - .setStatus(status) - .setProtocol(protocol) - .setDetectPoint(DetectPoint.client); + .setStartTime(startTime) + .setEndTime(startTime + duration) + .setSourceServiceName( + downstreamService.getServiceName()) + .setSourceServiceInstance( + downstreamService.getServiceInstanceName()) + .setDestServiceName( + destService.getServiceName()) + .setDestServiceInstance( + destService.getServiceInstanceName()) + .setEndpoint(endpoint) + .setLatency((int) duration) + .setResponseCode(Math.toIntExact(responseCode)) + .setStatus(status) + .setProtocol(protocol) + .setDetectPoint(DetectPoint.client); logger.debug("Transformed sidecar->sidecar(server side) inbound mesh metric {}", metric); forward(metric); @@ -280,12 +300,14 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { Address upstreamRemoteAddress = properties.getUpstreamRemoteAddress(); if (downstreamLocalAddress != null && downstreamRemoteAddress != null && upstreamRemoteAddress != null) { SocketAddress downstreamRemoteAddressSocketAddress = downstreamRemoteAddress.getSocketAddress(); - ServiceMetaInfo outside = find(downstreamRemoteAddressSocketAddress.getAddress(), downstreamRemoteAddressSocketAddress - .getPortValue()); + ServiceMetaInfo outside = find( + downstreamRemoteAddressSocketAddress.getAddress(), downstreamRemoteAddressSocketAddress + .getPortValue()); SocketAddress downstreamLocalAddressSocketAddress = downstreamLocalAddress.getSocketAddress(); - ServiceMetaInfo ingress = find(downstreamLocalAddressSocketAddress.getAddress(), downstreamLocalAddressSocketAddress - .getPortValue()); + ServiceMetaInfo ingress = find( + downstreamLocalAddressSocketAddress.getAddress(), downstreamLocalAddressSocketAddress + .getPortValue()); long startTime = formatAsLong(properties.getStartTime()); long duration = formatAsLong(properties.getTimeToLastDownstreamTxByte()); @@ -310,42 +332,51 @@ public class K8sALSServiceMeshHTTPAnalysis implements ALSHTTPAnalysis { boolean status = responseCode >= 200 && responseCode < 400; ServiceMeshMetric.Builder metric = ServiceMeshMetric.newBuilder() - .setStartTime(startTime) - .setEndTime(startTime + duration) - .setSourceServiceName(outside.getServiceName()) - .setSourceServiceInstance(outside.getServiceInstanceName()) - .setDestServiceName(ingress.getServiceName()) - .setDestServiceInstance(ingress.getServiceInstanceName()) - .setEndpoint(endpoint) - .setLatency((int) duration) - .setResponseCode(Math.toIntExact(responseCode)) - .setStatus(status) - .setProtocol(protocol) - .setDetectPoint(DetectPoint.server); + .setStartTime(startTime) + .setEndTime(startTime + duration) + .setSourceServiceName(outside.getServiceName()) + .setSourceServiceInstance( + outside.getServiceInstanceName()) + .setDestServiceName(ingress.getServiceName()) + .setDestServiceInstance( + ingress.getServiceInstanceName()) + .setEndpoint(endpoint) + .setLatency((int) duration) + .setResponseCode(Math.toIntExact(responseCode)) + .setStatus(status) + .setProtocol(protocol) + .setDetectPoint(DetectPoint.server); logger.debug("Transformed ingress inbound mesh metric {}", metric); forward(metric); SocketAddress upstreamRemoteAddressSocketAddress = upstreamRemoteAddress.getSocketAddress(); - ServiceMetaInfo targetService = find(upstreamRemoteAddressSocketAddress.getAddress(), upstreamRemoteAddressSocketAddress - .getPortValue()); + ServiceMetaInfo targetService = find( + upstreamRemoteAddressSocketAddress.getAddress(), upstreamRemoteAddressSocketAddress + .getPortValue()); long outboundStartTime = startTime + formatAsLong(properties.getTimeToFirstUpstreamTxByte()); long outboundEndTime = startTime + formatAsLong(properties.getTimeToLastUpstreamRxByte()); ServiceMeshMetric.Builder outboundMetric = ServiceMeshMetric.newBuilder() - .setStartTime(outboundStartTime) - .setEndTime(outboundEndTime) - .setSourceServiceName(ingress.getServiceName()) - .setSourceServiceInstance(ingress.getServiceInstanceName()) - .setDestServiceName(targetService.getServiceName()) - .setDestServiceInstance(targetService.getServiceInstanceName()) - .setEndpoint(endpoint) - .setLatency((int) (outboundEndTime - outboundStartTime)) - .setResponseCode(Math.toIntExact(responseCode)) - .setStatus(status) - .setProtocol(protocol) - .setDetectPoint(DetectPoint.client); + .setStartTime(outboundStartTime) + .setEndTime(outboundEndTime) + .setSourceServiceName( + ingress.getServiceName()) + .setSourceServiceInstance( + ingress.getServiceInstanceName()) + .setDestServiceName( + targetService.getServiceName()) + .setDestServiceInstance( + targetService.getServiceInstanceName()) + .setEndpoint(endpoint) + .setLatency( + (int) (outboundEndTime - outboundStartTime)) + .setResponseCode( + Math.toIntExact(responseCode)) + .setStatus(status) + .setProtocol(protocol) + .setDetectPoint(DetectPoint.client); logger.debug("Transformed ingress outbound mesh metric {}", outboundMetric); forward(outboundMetric); diff --git a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResourceTest.java b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResourceTest.java index dc3bc8358..9a4f551ca 100644 --- a/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResourceTest.java +++ b/oap-server/server-receiver-plugin/envoy-metrics-receiver-plugin/src/test/java/org/apache/skywalking/oap/server/receiver/envoy/als/DependencyResourceTest.java @@ -19,16 +19,15 @@ package org.apache.skywalking.oap.server.receiver.envoy.als; -import io.kubernetes.client.ApiException; -import io.kubernetes.client.models.V1ObjectMeta; -import io.kubernetes.client.models.V1OwnerReference; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.junit.runners.Parameterized; - +import io.kubernetes.client.openapi.ApiException; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.openapi.models.V1OwnerReference; import java.util.Arrays; import java.util.Collection; import java.util.Collections; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.Parameterized; import static org.hamcrest.core.Is.is; import static org.junit.Assert.assertThat; diff --git a/tools/dependencies/known-oap-backend-dependencies-es7.txt b/tools/dependencies/known-oap-backend-dependencies-es7.txt index 42df098f1..8cf054b79 100755 --- a/tools/dependencies/known-oap-backend-dependencies-es7.txt +++ b/tools/dependencies/known-oap-backend-dependencies-es7.txt @@ -5,24 +5,26 @@ antlr4-runtime-4.7.1.jar aopalliance-1.0.jar apollo-client-1.4.0.jar apollo-core-1.4.0.jar -bcpkix-jdk15on-1.59.jar -bcprov-ext-jdk15on-1.59.jar -bcprov-jdk15on-1.59.jar -builder-annotations-0.9.2.jar +bcpkix-jdk15on-1.61.jar +bcprov-ext-jdk15on-1.61.jar +bcprov-jdk15on-1.61.jar +builder-annotations-0.21.0.jar checker-qual-2.8.1.jar -client-java-4.0.0.jar -client-java-api-4.0.0.jar -client-java-proto-4.0.0.jar +client-java-8.0.0.jar +client-java-api-8.0.0.jar +client-java-proto-8.0.0.jar commons-codec-1.11.jar -commons-compress-1.18.jar +commons-collections4-4.1.jar +commons-compress-1.19.jar commons-dbcp-1.4.jar commons-io-2.6.jar commons-lang3-3.7.jar commons-pool-1.5.4.jar commons-text-1.4.jar +compiler-0.9.3.jar consul-client-1.2.6.jar converter-jackson-2.3.0.jar -compiler-0.9.3.jar +converter-moshi-2.5.0.jar curator-client-4.0.1.jar curator-framework-4.0.1.jar curator-recipes-4.0.1.jar @@ -49,6 +51,7 @@ grpc-protobuf-1.26.0.jar grpc-protobuf-lite-1.26.0.jar grpc-stub-1.26.0.jar gson-2.8.6.jar +gson-fire-1.8.3.jar guava-28.1-jre.jar guice-4.1.0.jar h2-1.4.196.jar @@ -59,6 +62,7 @@ httpasyncclient-4.1.4.jar httpclient-4.5.7.jar httpcore-4.4.11.jar httpcore-nio-4.4.11.jar +influxdb-java-2.15.jar jackson-annotations-2.9.5.jar jackson-core-2.9.5.jar jackson-core-asl-1.9.13.jar @@ -87,6 +91,7 @@ jna-4.5.1.jar joda-convert-1.2.jar joda-time-2.10.5.jar jopt-simple-4.6.jar +jose4j-0.7.0.jar json-flattener-0.6.0.jar jsr305-3.0.2.jar kotlin-reflect-1.1.1.jar @@ -98,7 +103,7 @@ log4j-api-2.9.0.jar log4j-core-2.9.0.jar log4j-over-slf4j-1.7.25.jar log4j-slf4j-impl-2.9.0.jar -logging-interceptor-2.7.5.jar +logging-interceptor-3.13.1.jar lucene-analyzers-common-8.0.0.jar lucene-backward-codecs-8.0.0.jar lucene-core-8.0.0.jar @@ -115,6 +120,8 @@ lucene-spatial-extras-8.0.0.jar lucene-spatial3d-8.0.0.jar lucene-suggest-8.0.0.jar minimal-json-0.9.5.jar +moshi-1.5.0.jar +msgpack-core-0.8.16.jar netty-3.10.5.Final.jar netty-buffer-4.1.42.Final.jar netty-codec-4.1.42.Final.jar @@ -129,9 +136,7 @@ netty-resolver-4.1.42.Final.jar netty-resolver-dns-4.1.42.Final.jar netty-tcnative-boringssl-static-2.0.26.Final.jar netty-transport-4.1.42.Final.jar -okhttp-2.7.5.jar okhttp-3.9.0.jar -okhttp-ws-2.7.5.jar okio-1.13.0.jar opencensus-api-0.24.0.jar opencensus-contrib-grpc-metrics-0.24.0.jar @@ -143,7 +148,7 @@ protobuf-java-util-3.11.4.jar rank-eval-client-7.0.0.jar reactive-streams-1.0.2.jar reflectasm-1.11.7.jar -resourcecify-annotations-0.9.2.jar +resourcecify-annotations-0.21.0.jar retrofit-2.3.0.jar simpleclient-0.6.0.jar simpleclient_common-0.6.0.jar @@ -151,15 +156,10 @@ simpleclient_hotspot-0.6.0.jar simpleclient_httpserver-0.6.0.jar slf4j-api-1.7.25.jar snakeyaml-1.18.jar -sundr-codegen-0.9.2.jar -sundr-core-0.9.2.jar -swagger-annotations-1.5.12.jar +sundr-codegen-0.21.0.jar +sundr-core-0.21.0.jar +swagger-annotations-1.5.22.jar t-digest-3.2.jar -zookeeper-3.4.10.jar -converter-moshi-2.5.0.jar -influxdb-java-2.15.jar -logging-interceptor-3.13.1.jar -moshi-1.5.0.jar -msgpack-core-0.8.16.jar vavr-0.10.3.jar vavr-match-0.10.3.jar +zookeeper-3.4.10.jar diff --git a/tools/dependencies/known-oap-backend-dependencies.txt b/tools/dependencies/known-oap-backend-dependencies.txt index 42620cb15..9ea4b732d 100755 --- a/tools/dependencies/known-oap-backend-dependencies.txt +++ b/tools/dependencies/known-oap-backend-dependencies.txt @@ -1,3 +1,5 @@ +HdrHistogram-2.1.9.jar +HikariCP-3.1.0.jar aggs-matrix-stats-client-6.3.2.jar animal-sniffer-annotations-1.18.jar annotations-13.0.jar @@ -5,17 +7,18 @@ antlr4-runtime-4.7.1.jar aopalliance-1.0.jar apollo-client-1.4.0.jar apollo-core-1.4.0.jar -bcpkix-jdk15on-1.59.jar -bcprov-ext-jdk15on-1.59.jar -bcprov-jdk15on-1.59.jar -builder-annotations-0.9.2.jar +bcpkix-jdk15on-1.61.jar +bcprov-ext-jdk15on-1.61.jar +bcprov-jdk15on-1.61.jar +builder-annotations-0.21.0.jar caffeine-2.6.2.jar checker-qual-2.8.1.jar -client-java-4.0.0.jar -client-java-api-4.0.0.jar -client-java-proto-4.0.0.jar +client-java-8.0.0.jar +client-java-api-8.0.0.jar +client-java-proto-8.0.0.jar commons-codec-1.11.jar -commons-compress-1.18.jar +commons-collections4-4.1.jar +commons-compress-1.19.jar commons-dbcp-1.4.jar commons-io-2.6.jar commons-lang3-3.7.jar @@ -23,6 +26,7 @@ commons-pool-1.5.4.jar commons-text-1.4.jar consul-client-1.2.6.jar converter-jackson-2.3.0.jar +converter-moshi-2.5.0.jar curator-client-4.0.1.jar curator-framework-4.0.1.jar curator-recipes-4.0.1.jar @@ -48,16 +52,16 @@ grpc-protobuf-1.26.0.jar grpc-protobuf-lite-1.26.0.jar grpc-stub-1.26.0.jar gson-2.8.6.jar +gson-fire-1.8.3.jar guava-28.1-jre.jar guice-4.1.0.jar h2-1.4.196.jar -HdrHistogram-2.1.9.jar -HikariCP-3.1.0.jar hppc-0.7.1.jar httpasyncclient-4.1.2.jar httpclient-4.5.2.jar httpcore-4.4.5.jar httpcore-nio-4.4.5.jar +influxdb-java-2.15.jar jackson-annotations-2.9.5.jar jackson-core-2.9.5.jar jackson-core-asl-1.9.13.jar @@ -86,6 +90,7 @@ jna-4.5.1.jar joda-convert-1.2.jar joda-time-2.10.5.jar jopt-simple-4.6.jar +jose4j-0.7.0.jar json-flattener-0.6.0.jar jsr305-3.0.2.jar kotlin-reflect-1.1.1.jar @@ -96,7 +101,7 @@ log4j-api-2.9.0.jar log4j-core-2.9.0.jar log4j-over-slf4j-1.7.25.jar log4j-slf4j-impl-2.9.0.jar -logging-interceptor-2.7.5.jar +logging-interceptor-3.13.1.jar lucene-analyzers-common-7.3.1.jar lucene-backward-codecs-7.3.1.jar lucene-core-7.3.1.jar @@ -113,6 +118,8 @@ lucene-spatial-extras-7.3.1.jar lucene-spatial3d-7.3.1.jar lucene-suggest-7.3.1.jar minimal-json-0.9.5.jar +moshi-1.5.0.jar +msgpack-core-0.8.16.jar netty-3.10.5.Final.jar netty-buffer-4.1.42.Final.jar netty-codec-4.1.42.Final.jar @@ -127,9 +134,7 @@ netty-resolver-4.1.42.Final.jar netty-resolver-dns-4.1.42.Final.jar netty-tcnative-boringssl-static-2.0.26.Final.jar netty-transport-4.1.42.Final.jar -okhttp-2.7.5.jar okhttp-3.9.0.jar -okhttp-ws-2.7.5.jar okio-1.13.0.jar opencensus-api-0.24.0.jar opencensus-contrib-grpc-metrics-0.24.0.jar @@ -141,7 +146,7 @@ protobuf-java-util-3.11.4.jar rank-eval-client-6.3.2.jar reactive-streams-1.0.2.jar reflectasm-1.11.7.jar -resourcecify-annotations-0.9.2.jar +resourcecify-annotations-0.21.0.jar retrofit-2.3.0.jar simpleclient-0.6.0.jar simpleclient_common-0.6.0.jar @@ -149,16 +154,11 @@ simpleclient_hotspot-0.6.0.jar simpleclient_httpserver-0.6.0.jar slf4j-api-1.7.25.jar snakeyaml-1.18.jar -sundr-codegen-0.9.2.jar -sundr-core-0.9.2.jar -swagger-annotations-1.5.12.jar +sundr-codegen-0.21.0.jar +sundr-core-0.21.0.jar +swagger-annotations-1.5.22.jar t-digest-3.2.jar -zipkin-2.9.1.jar -zookeeper-3.4.10.jar -converter-moshi-2.5.0.jar -influxdb-java-2.15.jar -logging-interceptor-3.13.1.jar -moshi-1.5.0.jar -msgpack-core-0.8.16.jar vavr-0.10.3.jar vavr-match-0.10.3.jar +zipkin-2.9.1.jar +zookeeper-3.4.10.jar