From ed99079d20fa321f0dedd764ffa7538416aa6b0e Mon Sep 17 00:00:00 2001
From: zhangwei <26432832+arugal@users.noreply.github.com>
Date: Wed, 17 Jul 2019 14:40:45 +0800
Subject: [PATCH] Fix service cluster plugin bug (#3074)
* fix 3069
---
oap-server/pom.xml | 1 +
.../cluster-consul-plugin/pom.xml | 100 +++++++
.../plugin/consul/ConsulCoordinator.java | 10 +-
.../plugin/consul/ConsulCoordinatorTest.java | 3 +-
...terModuleConsulProviderFunctionalTest.java | 263 ++++++++++++++++++
...usterModuleEtcdProviderFunctionalTest.java | 232 +++++++++++++++
.../cluster-nacos-plugin/pom.xml | 96 +++++++
.../plugin/nacos/NacosCoordinator.java | 11 +-
...sterModuleNacosProviderFunctionalTest.java | 198 +++++++++++++
.../plugin/nacos/NacosCoordinatorTest.java | 3 +-
.../cluster-zookeeper-plugin/pom.xml | 95 +++++++
.../ClusterModuleZookeeperProvider.java | 3 +-
.../zookeeper/ZookeeperCoordinator.java | 40 ++-
.../ClusterModuleZookeeperProviderTest.java | 56 ++++
...lusterModuleZookeeperProviderTestCase.java | 77 -----
...ModuleZookeeperProviderFunctionalTest.java | 235 ++++++++++++++++
.../zookeeper/ZookeeperCoordinatorTest.java | 7 +-
.../configuration-zookeeper/pom.xml | 4 -
18 files changed, 1309 insertions(+), 125 deletions(-)
create mode 100644 oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ITClusterModuleConsulProviderFunctionalTest.java
create mode 100644 oap-server/server-cluster-plugin/cluster-etcd-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/etcd/ITClusterModuleEtcdProviderFunctionalTest.java
create mode 100644 oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/ITClusterModuleNacosProviderFunctionalTest.java
create mode 100644 oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTest.java
delete mode 100644 oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTestCase.java
create mode 100644 oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ITClusterModuleZookeeperProviderFunctionalTest.java
diff --git a/oap-server/pom.xml b/oap-server/pom.xml
index 7c58df042..8816919a9 100644
--- a/oap-server/pom.xml
+++ b/oap-server/pom.xml
@@ -83,6 +83,7 @@
2.12.0
2.17.0
v3.2.3
+ 3.5
diff --git a/oap-server/server-cluster-plugin/cluster-consul-plugin/pom.xml b/oap-server/server-cluster-plugin/cluster-consul-plugin/pom.xml
index 7549e4494..d4ef6b256 100644
--- a/oap-server/server-cluster-plugin/cluster-consul-plugin/pom.xml
+++ b/oap-server/server-cluster-plugin/cluster-consul-plugin/pom.xml
@@ -28,6 +28,10 @@
cluster-consul-plugin
jar
+
+ 0.9
+
+
org.apache.skywalking
@@ -50,4 +54,100 @@
+
+
+
+ CI-with-IT
+
+
+
+ io.fabric8
+ docker-maven-plugin
+
+ all
+ default
+ true
+ true
+ IfNotPresent
+
+
+
+ prepare-consul
+ pre-integration-test
+
+ start
+
+
+
+
+ consul:${consul.image.version}
+ cluster-consul-plugin-integration-test-cluster
+
+ agent -server -bootstrap-expect=1 -client=0.0.0.0
+
+ consul.port:8500
+
+
+ Synced node info
+
+
+
+
+
+
+
+
+ prepare-consul-stop
+ post-integration-test
+
+ stop
+
+
+
+
+
+ org.codehaus.gmaven
+ gmaven-plugin
+ ${gmaven-plugin.version}
+
+
+ add-default-properties
+ initialize
+
+ execute
+
+
+ 2.0
+
+ project.properties.setProperty('docker.hostname', 'localhost')
+
+ log.info("Docker hostname is " + project.properties['docker.hostname'])
+
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-failsafe-plugin
+
+
+
+ ${docker.hostname}:${consul.port}
+
+
+
+
+
+
+ integration-test
+ verify
+
+
+
+
+
+
+
+
diff --git a/oap-server/server-cluster-plugin/cluster-consul-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinator.java b/oap-server/server-cluster-plugin/cluster-consul-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinator.java
index 1f6c23212..f1ea9b912 100644
--- a/oap-server/server-cluster-plugin/cluster-consul-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinator.java
+++ b/oap-server/server-cluster-plugin/cluster-consul-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinator.java
@@ -54,13 +54,11 @@ public class ConsulCoordinator implements ClusterRegister, ClusterNodesQuery {
if (CollectionUtils.isNotEmpty(nodes)) {
nodes.forEach(node -> {
if (!Strings.isNullOrEmpty(node.getService().getAddress())) {
- if (Objects.nonNull(selfAddress)) {
- if (selfAddress.getHost().equals(node.getService().getAddress()) && selfAddress.getPort() == node.getService().getPort()) {
- remoteInstances.add(new RemoteInstance(new Address(node.getService().getAddress(), node.getService().getPort(), true)));
- } else {
- remoteInstances.add(new RemoteInstance(new Address(node.getService().getAddress(), node.getService().getPort(), false)));
- }
+ Address address = new Address(node.getService().getAddress(), node.getService().getPort(), false);
+ if (address.equals(selfAddress)) {
+ address.setSelf(true);
}
+ remoteInstances.add(new RemoteInstance(address));
}
});
}
diff --git a/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinatorTest.java b/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinatorTest.java
index 181e00393..d4954199c 100644
--- a/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinatorTest.java
+++ b/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ConsulCoordinatorTest.java
@@ -104,7 +104,8 @@ public class ConsulCoordinatorTest {
List serviceHealths = mockHealth();
when(consulResponse.getResponse()).thenReturn(serviceHealths);
List remoteInstances = coordinator.queryRemoteNodes();
- assertTrue(remoteInstances.isEmpty());
+ // filter empty address
+ assertEquals(2, remoteInstances.size());
}
@Test
diff --git a/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ITClusterModuleConsulProviderFunctionalTest.java b/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ITClusterModuleConsulProviderFunctionalTest.java
new file mode 100644
index 000000000..48a170ffb
--- /dev/null
+++ b/oap-server/server-cluster-plugin/cluster-consul-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/consul/ITClusterModuleConsulProviderFunctionalTest.java
@@ -0,0 +1,263 @@
+/*
+ * 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.consul;
+
+import com.google.common.base.Strings;
+import com.orbitz.consul.AgentClient;
+import com.orbitz.consul.Consul;
+import com.orbitz.consul.model.agent.ImmutableRegistration;
+import com.orbitz.consul.model.agent.Registration;
+import org.apache.skywalking.apm.util.StringUtil;
+import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery;
+import org.apache.skywalking.oap.server.core.cluster.ClusterRegister;
+import org.apache.skywalking.oap.server.core.cluster.RemoteInstance;
+import org.apache.skywalking.oap.server.core.remote.client.Address;
+import org.apache.skywalking.oap.server.library.module.ModuleProvider;
+import org.apache.skywalking.oap.server.telemetry.api.TelemetryRelatedContext;
+import org.junit.Before;
+import org.junit.Test;
+import org.powermock.reflect.Whitebox;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * @author zhangwei
+ */
+public class ITClusterModuleConsulProviderFunctionalTest {
+
+ private String consulAddress;
+
+ @Before
+ public void before() {
+ consulAddress = System.getProperty("consul.address");
+ assertFalse(StringUtil.isEmpty(consulAddress));
+ }
+
+ @Test
+ public void registerRemote() throws Exception {
+ final String serviceName = "register_remote";
+ ModuleProvider provider = createProvider(serviceName);
+
+ Address selfAddress = new Address("127.0.0.1", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfInternal() throws Exception {
+ final String serviceName = "register_remote_internal";
+ ModuleProvider provider =
+ createProvider(serviceName, "127.0.1.2", 1001);
+
+ Address selfAddress = new Address("127.0.0.2", 1002, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+
+ ClusterModuleConsulConfig config = (ClusterModuleConsulConfig) provider.createConfigBeanIfAbsent();
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(config.getInternalComHost(), queryAddress.getHost());
+ assertEquals(config.getInternalComPort(), queryAddress.getPort());
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfReceiver() throws Exception {
+ final String serviceName = "register_remote_receiver";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ // Mixed or Aggregator
+ Address selfAddress = new Address("127.0.0.3", 1003, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(providerA).registerRemote(instance);
+
+ // Receiver
+ List remoteInstances = queryRemoteNodes(providerB, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertFalse(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfCluster() throws Exception {
+ final String serviceName = "register_remote_cluster";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ Address addressA = new Address("127.0.0.4", 1004, true);
+ Address addressB = new Address("127.0.0.5", 1005, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+ }
+
+ @Test
+ public void unregisterRemoteOfCluster() throws Exception {
+ final String serviceName = "unregister_remote_cluster";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ Address addressA = new Address("127.0.0.6", 1006, true);
+ Address addressB = new Address("127.0.0.7", 1007, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+
+ // unregister A
+ Consul client = Whitebox.getInternalState(providerA, "client");
+ AgentClient agentClient = client.agentClient();
+ agentClient.deregister(instanceA.getAddress().toString());
+
+ // only B
+ remoteInstancesOfB = queryRemoteNodes(providerB, 1, 120);
+ assertEquals(1, remoteInstancesOfB.size());
+ Address address = remoteInstancesOfB.get(0).getAddress();
+ assertEquals(address, addressB);
+ assertTrue(addressB.isSelf());
+ }
+
+ private ClusterModuleConsulProvider createProvider(String serviceName) throws Exception {
+ return createProvider(serviceName, null, 0);
+ }
+
+ private ClusterModuleConsulProvider createProvider(String serviceName, String internalComHost, int internalComPort) throws Exception {
+ ClusterModuleConsulProvider provider = new ClusterModuleConsulProvider();
+
+ ClusterModuleConsulConfig config = (ClusterModuleConsulConfig) provider.createConfigBeanIfAbsent();
+
+ config.setHostPort(consulAddress);
+ config.setServiceName(serviceName);
+
+ if (!StringUtil.isEmpty(internalComHost)) {
+ config.setInternalComHost(internalComHost);
+ }
+
+ if (internalComPort > 0) {
+ config.setInternalComPort(internalComPort);
+ }
+
+ provider.prepare();
+ provider.start();
+ provider.notifyAfterCompleted();
+
+ ConsulCoordinator consulCoordinator = (ConsulCoordinator) provider.getService(ClusterRegister.class);
+
+ // ignore health check
+ ClusterRegister register = remoteInstance -> {
+ if (needUsingInternalAddr(config)) {
+ remoteInstance = new RemoteInstance(new Address(config.getInternalComHost(), config.getInternalComPort(), true));
+ }
+
+ Consul client = Whitebox.getInternalState(consulCoordinator, "client");
+ AgentClient agentClient = client.agentClient();
+ Whitebox.setInternalState(consulCoordinator, "selfAddress", remoteInstance.getAddress());
+ TelemetryRelatedContext.INSTANCE.setId(remoteInstance.getAddress().toString());
+ Registration registration = ImmutableRegistration.builder()
+ .id(remoteInstance.getAddress().toString())
+ .name(serviceName)
+ .address(remoteInstance.getAddress().getHost())
+ .port(remoteInstance.getAddress().getPort())
+ .build();
+
+ agentClient.register(registration);
+ };
+
+ provider.registerServiceImplementation(ClusterRegister.class, register);
+ return provider;
+ }
+
+ private boolean needUsingInternalAddr(ClusterModuleConsulConfig config) {
+ return !Strings.isNullOrEmpty(config.getInternalComHost()) && config.getInternalComPort() > 0;
+ }
+
+ private ClusterRegister getClusterRegister(ModuleProvider provider) {
+ return provider.getService(ClusterRegister.class);
+ }
+
+ private ClusterNodesQuery getClusterNodesQuery(ModuleProvider provider) {
+ return provider.getService(ClusterNodesQuery.class);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals) throws InterruptedException {
+ return queryRemoteNodes(provider, goals, 20);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals, int cyclic) throws InterruptedException {
+ do {
+ List instances = getClusterNodesQuery(provider).queryRemoteNodes();
+ if (instances.size() == goals) {
+ return instances;
+ } else {
+ Thread.sleep(1000);
+ }
+ } while (--cyclic > 0);
+ return Collections.EMPTY_LIST;
+ }
+
+ private void validateServiceInstance(Address selfAddress, Address otherAddress, List queryResult) {
+ assertEquals(2, queryResult.size());
+
+ boolean selfExist = false, otherExist = false;
+
+ for (RemoteInstance instance : queryResult) {
+ Address queryAddress = instance.getAddress();
+ if (queryAddress.equals(selfAddress) && queryAddress.isSelf()) {
+ selfExist = true;
+ } else if (queryAddress.equals(otherAddress) && !queryAddress.isSelf()) {
+ otherExist = true;
+ }
+ }
+
+ assertTrue(selfExist);
+ assertTrue(otherExist);
+ }
+}
diff --git a/oap-server/server-cluster-plugin/cluster-etcd-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/etcd/ITClusterModuleEtcdProviderFunctionalTest.java b/oap-server/server-cluster-plugin/cluster-etcd-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/etcd/ITClusterModuleEtcdProviderFunctionalTest.java
new file mode 100644
index 000000000..ceb442e7d
--- /dev/null
+++ b/oap-server/server-cluster-plugin/cluster-etcd-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/etcd/ITClusterModuleEtcdProviderFunctionalTest.java
@@ -0,0 +1,232 @@
+/*
+ * 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.etcd;
+
+import mousio.etcd4j.EtcdClient;
+import org.apache.skywalking.apm.util.StringUtil;
+import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery;
+import org.apache.skywalking.oap.server.core.cluster.ClusterRegister;
+import org.apache.skywalking.oap.server.core.cluster.RemoteInstance;
+import org.apache.skywalking.oap.server.core.remote.client.Address;
+import org.apache.skywalking.oap.server.library.module.ModuleProvider;
+import org.apache.skywalking.oap.server.library.module.ModuleStartException;
+import org.junit.Before;
+import org.junit.Test;
+import org.powermock.reflect.Whitebox;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * @author zhangwei
+ */
+public class ITClusterModuleEtcdProviderFunctionalTest {
+
+ private String etcdAddress;
+
+ @Before
+ public void before() {
+ String etcdHost = System.getProperty("etcd.host");
+ String port = System.getProperty("etcd.port");
+ assertTrue(!StringUtil.isEmpty(etcdHost) && !StringUtil.isEmpty(port));
+ etcdAddress = etcdHost + ":" + port;
+ }
+
+ @Test
+ public void registerRemote() throws Exception {
+ final String serviceName = "register_remote";
+ ModuleProvider provider = createProvider(serviceName);
+
+ Address selfAddress = new Address("127.0.0.1", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfInternal() throws Exception {
+ final String serviceName = "register_remote_internal";
+ ModuleProvider provider =
+ createProvider(serviceName, "127.0.1.2", 1000);
+
+ Address selfAddress = new Address("127.0.0.2", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+
+ ClusterModuleEtcdConfig config = (ClusterModuleEtcdConfig) provider.createConfigBeanIfAbsent();
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(config.getInternalComHost(), queryAddress.getHost());
+ assertEquals(config.getInternalComPort(), queryAddress.getPort());
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfReceiver() throws Exception {
+ final String serviceName = "register_remote_receiver";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ // Mixed or Aggregator
+ Address selfAddress = new Address("127.0.0.3", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(providerA).registerRemote(instance);
+
+ // Receiver
+ List remoteInstances = queryRemoteNodes(providerB, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertFalse(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfCluster() throws Exception {
+ final String serviceName = "register_remote_cluster";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ Address addressA = new Address("127.0.0.4", 1000, true);
+ Address addressB = new Address("127.0.0.5", 1000, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+ }
+
+ @Test
+ public void unregisterRemoteOfCluster() throws Exception {
+ final String serviceName = "unregister_remote_cluster";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ Address addressA = new Address("127.0.0.4", 1000, true);
+ Address addressB = new Address("127.0.0.5", 1000, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+
+ // unregister A
+ EtcdClient client = Whitebox.getInternalState(providerA, "client");
+ client.close();
+
+ // only B
+ remoteInstancesOfB = queryRemoteNodes(providerB, 1, 120);
+ assertEquals(1, remoteInstancesOfB.size());
+ Address address = remoteInstancesOfB.get(0).getAddress();
+ assertEquals(address, addressB);
+ assertTrue(addressB.isSelf());
+ }
+
+ private ClusterModuleEtcdProvider createProvider(String serviceName) throws ModuleStartException {
+ return createProvider(serviceName, null, 0);
+ }
+
+ private ClusterModuleEtcdProvider createProvider(String serviceName, String internalComHost, int internalComPort) throws ModuleStartException {
+ ClusterModuleEtcdProvider provider = new ClusterModuleEtcdProvider();
+
+ ClusterModuleEtcdConfig config = (ClusterModuleEtcdConfig) provider.createConfigBeanIfAbsent();
+
+ config.setHostPort(etcdAddress);
+ config.setServiceName(serviceName);
+
+ if (!StringUtil.isEmpty(internalComHost)) {
+ config.setInternalComHost(internalComHost);
+ }
+
+ if (internalComPort > 0) {
+ config.setInternalComPort(internalComPort);
+ }
+
+ provider.prepare();
+ provider.start();
+ provider.notifyAfterCompleted();
+ return provider;
+ }
+
+ private ClusterRegister getClusterRegister(ModuleProvider provider) {
+ return provider.getService(ClusterRegister.class);
+ }
+
+ private ClusterNodesQuery getClusterNodesQuery(ModuleProvider provider) {
+ return provider.getService(ClusterNodesQuery.class);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals) throws InterruptedException {
+ return queryRemoteNodes(provider, goals, 20);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals, int cyclic) throws InterruptedException {
+ do {
+ List instances = getClusterNodesQuery(provider).queryRemoteNodes();
+ if (instances.size() == goals) {
+ return instances;
+ } else {
+ Thread.sleep(1000);
+ }
+ } while (--cyclic > 0);
+ return Collections.EMPTY_LIST;
+ }
+
+ private void validateServiceInstance(Address selfAddress, Address otherAddress, List queryResult) {
+ assertEquals(2, queryResult.size());
+
+ boolean selfExist = false, otherExist = false;
+
+ for (RemoteInstance instance : queryResult) {
+ Address queryAddress = instance.getAddress();
+ if (queryAddress.equals(selfAddress) && queryAddress.isSelf()) {
+ selfExist = true;
+ } else if (queryAddress.equals(otherAddress) && !queryAddress.isSelf()) {
+ otherExist = true;
+ }
+ }
+
+ assertTrue(selfExist);
+ assertTrue(otherExist);
+ }
+}
diff --git a/oap-server/server-cluster-plugin/cluster-nacos-plugin/pom.xml b/oap-server/server-cluster-plugin/cluster-nacos-plugin/pom.xml
index 7af37c45f..ab4af2123 100644
--- a/oap-server/server-cluster-plugin/cluster-nacos-plugin/pom.xml
+++ b/oap-server/server-cluster-plugin/cluster-nacos-plugin/pom.xml
@@ -63,4 +63,100 @@
+
+
+
+ CI-with-IT
+
+
+
+ io.fabric8
+ docker-maven-plugin
+
+ all
+ true
+ default
+ true
+ IfNotPresent
+
+
+ nacos/nacos-server:${nacos.version}
+ cluster-nacos-plugin-integration-test-nacos
+
+
+ standalone
+
+
+ nacos.port:8848
+
+
+ Nacos started successfully
+
+
+
+
+
+
+
+
+ start
+ pre-integration-test
+
+ start
+
+
+
+ stop
+ post-integration-test
+
+ stop
+
+
+
+
+
+ org.codehaus.gmaven
+ gmaven-plugin
+ ${gmaven-plugin.version}
+
+
+ add-default-properties
+ initialize
+
+ execute
+
+
+ 2.0
+
+ project.properties.setProperty('docker.hostname', 'localhost')
+
+ log.info("Docker hostname is " + project.properties['docker.hostname'])
+
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-failsafe-plugin
+
+
+
+ ${docker.hostname}:${nacos.port}
+
+
+
+
+
+
+ integration-test
+ verify
+
+
+
+
+
+
+
+
diff --git a/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinator.java b/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinator.java
index 10d56f4ab..4225a753a 100644
--- a/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinator.java
+++ b/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinator.java
@@ -27,7 +27,6 @@ import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import java.util.ArrayList;
import java.util.List;
-import java.util.Objects;
/**
* @author caoyixiong
@@ -50,13 +49,11 @@ public class NacosCoordinator implements ClusterRegister, ClusterNodesQuery {
List instances = namingService.selectInstances(config.getServiceName(), true);
if (CollectionUtils.isNotEmpty(instances)) {
instances.forEach(instance -> {
- if (Objects.nonNull(selfAddress)) {
- if (selfAddress.getHost().equals(instance.getIp()) && selfAddress.getPort() == instance.getPort()) {
- result.add(new RemoteInstance(new Address(instance.getIp(), instance.getPort(), true)));
- } else {
- result.add(new RemoteInstance(new Address(instance.getIp(), instance.getPort(), false)));
- }
+ Address address = new Address(instance.getIp(), instance.getPort(), false);
+ if (address.equals(selfAddress)) {
+ address.setSelf(true);
}
+ result.add(new RemoteInstance(address));
});
}
} catch (NacosException e) {
diff --git a/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/ITClusterModuleNacosProviderFunctionalTest.java b/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/ITClusterModuleNacosProviderFunctionalTest.java
new file mode 100644
index 000000000..12200ffc7
--- /dev/null
+++ b/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/ITClusterModuleNacosProviderFunctionalTest.java
@@ -0,0 +1,198 @@
+/*
+ * 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.nacos;
+
+import com.alibaba.nacos.api.naming.NamingService;
+import org.apache.skywalking.apm.util.StringUtil;
+import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery;
+import org.apache.skywalking.oap.server.core.cluster.ClusterRegister;
+import org.apache.skywalking.oap.server.core.cluster.RemoteInstance;
+import org.apache.skywalking.oap.server.core.remote.client.Address;
+import org.apache.skywalking.oap.server.library.module.ModuleProvider;
+import org.apache.skywalking.oap.server.library.module.ModuleStartException;
+import org.junit.Before;
+import org.junit.Test;
+import org.powermock.reflect.Whitebox;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * @author zhangwei
+ */
+public class ITClusterModuleNacosProviderFunctionalTest {
+
+
+ private String nacosAddress;
+
+ @Before
+ public void before() {
+ nacosAddress = System.getProperty("nacos.address");
+ assertFalse(StringUtil.isEmpty(nacosAddress));
+ }
+
+ @Test
+ public void registerRemote() throws Exception {
+ final String serviceName = "register_remote";
+ ModuleProvider provider = createProvider(serviceName);
+
+ Address selfAddress = new Address("127.0.0.1", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfReceiver() throws Exception {
+ final String serviceName = "register_remote_receiver";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ // Mixed or Aggregator
+ Address selfAddress = new Address("127.0.0.3", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(providerA).registerRemote(instance);
+
+ // Receiver
+ List remoteInstances = queryRemoteNodes(providerB, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertFalse(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfCluster() throws Exception {
+ final String serviceName = "register_remote_cluster";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ Address addressA = new Address("127.0.0.4", 1000, true);
+ Address addressB = new Address("127.0.0.5", 1000, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+ }
+
+ @Test
+ public void deregisterRemoteOfCluster() throws Exception {
+ final String serviceName = "deregister_remote_cluster";
+ ModuleProvider providerA = createProvider(serviceName);
+ ModuleProvider providerB = createProvider(serviceName);
+
+ Address addressA = new Address("127.0.0.6", 1000, true);
+ Address addressB = new Address("127.0.0.7", 1000, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+
+ // deregister A
+ ClusterRegister register = getClusterRegister(providerA);
+ NamingService namingServiceA = Whitebox.getInternalState(register, "namingService");
+ namingServiceA.deregisterInstance(serviceName, addressA.getHost(), addressA.getPort());
+
+ // only B
+ remoteInstancesOfB = queryRemoteNodes(providerB, 1);
+ assertEquals(1, remoteInstancesOfB.size());
+ Address address = remoteInstancesOfB.get(0).getAddress();
+ assertEquals(addressB, address);
+ assertTrue(address.isSelf());
+ }
+
+ private ClusterModuleNacosProvider createProvider(String servicName) throws ModuleStartException {
+ ClusterModuleNacosProvider provider = new ClusterModuleNacosProvider();
+
+ ClusterModuleNacosConfig config = (ClusterModuleNacosConfig) provider.createConfigBeanIfAbsent();
+
+ config.setHostPort(nacosAddress);
+ config.setServiceName(servicName);
+
+ provider.prepare();
+ provider.start();
+ provider.notifyAfterCompleted();
+ return provider;
+ }
+
+ private ClusterRegister getClusterRegister(ModuleProvider provider) {
+ return provider.getService(ClusterRegister.class);
+ }
+
+ private ClusterNodesQuery getClusterNodesQuery(ModuleProvider provider) {
+ return provider.getService(ClusterNodesQuery.class);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals) throws InterruptedException {
+ int i = 20;
+ do {
+ List instances = getClusterNodesQuery(provider).queryRemoteNodes();
+ if (instances.size() == goals) {
+ return instances;
+ } else {
+ Thread.sleep(1000);
+ }
+ } while (--i > 0);
+ return Collections.EMPTY_LIST;
+ }
+
+ private void validateServiceInstance(Address selfAddress, Address otherAddress, List queryResult) {
+ assertEquals(2, queryResult.size());
+
+ boolean selfExist = false, otherExist = false;
+
+ for (RemoteInstance instance : queryResult) {
+ Address queryAddress = instance.getAddress();
+ if (queryAddress.equals(selfAddress) && queryAddress.isSelf()) {
+ selfExist = true;
+ } else if (queryAddress.equals(otherAddress) && !queryAddress.isSelf()) {
+ otherExist = true;
+ }
+ }
+
+ assertTrue(selfExist);
+ assertTrue(otherExist);
+ }
+
+}
diff --git a/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinatorTest.java b/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinatorTest.java
index b362d5a54..f73eed57b 100644
--- a/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinatorTest.java
+++ b/oap-server/server-cluster-plugin/cluster-nacos-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/nacos/NacosCoordinatorTest.java
@@ -32,7 +32,6 @@ import java.util.Collections;
import java.util.List;
import static org.junit.Assert.assertEquals;
-import static org.junit.Assert.assertTrue;
import static org.mockito.Matchers.anyBoolean;
import static org.mockito.Matchers.anyString;
import static org.mockito.Mockito.mock;
@@ -85,7 +84,7 @@ public class NacosCoordinatorTest {
List instances = mockInstance();
when(namingService.selectInstances(anyString(), anyBoolean())).thenReturn(instances);
List remoteInstances = coordinator.queryRemoteNodes();
- assertTrue(remoteInstances.isEmpty());
+ assertEquals(remoteInstances.size(), instances.size());
}
@Test
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/pom.xml b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/pom.xml
index 67c5e04cf..6defdd646 100644
--- a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/pom.xml
+++ b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/pom.xml
@@ -51,4 +51,99 @@
curator-test
+
+
+
+ CI-with-IT
+
+
+
+ io.fabric8
+ docker-maven-plugin
+
+ all
+ default
+ true
+ true
+ IfNotPresent
+
+
+
+ prepare-zookeeper
+ pre-integration-test
+
+ start
+
+
+
+
+ zookeeper:${zookeeper.image.version}
+ cluster-zookeeper-plugin-integration-test-zookeeper
+
+
+ zk-port:2181
+
+
+ binding to port
+
+
+
+
+
+
+
+
+ prepare-zookeeper-start
+ post-integration-test
+
+ stop
+
+
+
+
+
+ org.codehaus.gmaven
+ gmaven-plugin
+ ${gmaven-plugin.version}
+
+
+ add-default-properties
+ initialize
+
+ execute
+
+
+ 2.0
+
+ project.properties.setProperty('docker.hostname', 'localhost')
+
+ log.info("Docker hostname is " + project.properties['docker.hostname'])
+
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-failsafe-plugin
+
+
+
+ ${docker.hostname}:${zk-port}
+
+
+
+
+
+
+ integration-test
+ verify
+
+
+
+
+
+
+
+
\ No newline at end of file
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProvider.java b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProvider.java
index b7be978d5..f1f274057 100644
--- a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProvider.java
+++ b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProvider.java
@@ -70,16 +70,17 @@ public class ClusterModuleZookeeperProvider extends ModuleProvider {
.watchInstances(true)
.serializer(new SWInstanceSerializer()).build();
+ ZookeeperCoordinator coordinator;
try {
client.start();
client.blockUntilConnected();
serviceDiscovery.start();
+ coordinator = new ZookeeperCoordinator(config, serviceDiscovery);
} catch (Exception e) {
logger.error(e.getMessage(), e);
throw new ModuleStartException(e.getMessage(), e);
}
- ZookeeperCoordinator coordinator = new ZookeeperCoordinator(config, serviceDiscovery);
this.registerServiceImplementation(ClusterRegister.class, coordinator);
this.registerServiceImplementation(ClusterNodesQuery.class, coordinator);
}
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinator.java b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinator.java
index 06c95bdf8..5ab36cb36 100644
--- a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinator.java
+++ b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/main/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinator.java
@@ -32,25 +32,28 @@ import org.slf4j.*;
public class ZookeeperCoordinator implements ClusterRegister, ClusterNodesQuery {
private static final Logger logger = LoggerFactory.getLogger(ZookeeperCoordinator.class);
+ private static final String REMOTE_NAME_PATH = "remote";
+
private final ClusterModuleZookeeperConfig config;
private final ServiceDiscovery serviceDiscovery;
- private volatile ServiceCache serviceCache;
+ private final ServiceCache serviceCache;
private volatile Address selfAddress;
- ZookeeperCoordinator(ClusterModuleZookeeperConfig config, ServiceDiscovery serviceDiscovery) {
+ ZookeeperCoordinator(ClusterModuleZookeeperConfig config, ServiceDiscovery serviceDiscovery) throws Exception {
this.config = config;
this.serviceDiscovery = serviceDiscovery;
+ this.serviceCache = serviceDiscovery.serviceCacheBuilder().name(REMOTE_NAME_PATH).build();
+ this.serviceCache.start();
}
@Override public synchronized void registerRemote(RemoteInstance remoteInstance) throws ServiceRegisterException {
try {
- String remoteNamePath = "remote";
if (needUsingInternalAddr()) {
remoteInstance = new RemoteInstance(new Address(config.getInternalComHost(), config.getInternalComPort(), true));
}
ServiceInstance thisInstance = ServiceInstance.builder()
- .name(remoteNamePath)
+ .name(REMOTE_NAME_PATH)
.id(UUID.randomUUID().toString())
.address(remoteInstance.getAddress().getHost())
.port(remoteInstance.getAddress().getPort())
@@ -59,12 +62,6 @@ public class ZookeeperCoordinator implements ClusterRegister, ClusterNodesQuery
serviceDiscovery.registerService(thisInstance);
- serviceCache = serviceDiscovery.serviceCacheBuilder()
- .name(remoteNamePath)
- .build();
-
- serviceCache.start();
-
this.selfAddress = remoteInstance.getAddress();
TelemetryRelatedContext.INSTANCE.setId(selfAddress.toString());
} catch (Exception e) {
@@ -74,19 +71,16 @@ public class ZookeeperCoordinator implements ClusterRegister, ClusterNodesQuery
@Override public List queryRemoteNodes() {
List remoteInstanceDetails = new ArrayList<>(20);
- if (Objects.nonNull(serviceCache)) {
- List> serviceInstances = serviceCache.getInstances();
-
- serviceInstances.forEach(serviceInstance -> {
- RemoteInstance instance = serviceInstance.getPayload();
- if (instance.getAddress().equals(selfAddress)) {
- instance.getAddress().setSelf(true);
- } else {
- instance.getAddress().setSelf(false);
- }
- remoteInstanceDetails.add(instance);
- });
- }
+ List> serviceInstances = serviceCache.getInstances();
+ serviceInstances.forEach(serviceInstance -> {
+ RemoteInstance instance = serviceInstance.getPayload();
+ if (instance.getAddress().equals(selfAddress)) {
+ instance.getAddress().setSelf(true);
+ } else {
+ instance.getAddress().setSelf(false);
+ }
+ remoteInstanceDetails.add(instance);
+ });
return remoteInstanceDetails;
}
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTest.java b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTest.java
new file mode 100644
index 000000000..8a6ac730d
--- /dev/null
+++ b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTest.java
@@ -0,0 +1,56 @@
+/*
+ * 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.zookeeper;
+
+import org.apache.skywalking.oap.server.core.cluster.ClusterModule;
+import org.apache.skywalking.oap.server.library.module.ModuleConfig;
+import org.junit.Test;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * @author peng-yongsheng zhangwei
+ */
+public class ClusterModuleZookeeperProviderTest {
+
+ private ClusterModuleZookeeperProvider provider = new ClusterModuleZookeeperProvider();
+
+ @Test
+ public void name() {
+ assertEquals("zookeeper", provider.name());
+ }
+
+ @Test
+ public void module() {
+ assertEquals(ClusterModule.class, provider.module());
+ }
+
+ @Test
+ public void createConfigBeanIfAbsent() {
+ ModuleConfig moduleConfig = provider.createConfigBeanIfAbsent();
+ assertTrue(moduleConfig instanceof ClusterModuleZookeeperConfig);
+ }
+
+ @Test
+ public void requiredModules() {
+ String[] modules = provider.requiredModules();
+ assertEquals(0, modules.length);
+ }
+}
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTestCase.java b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTestCase.java
deleted file mode 100644
index c8f64461e..000000000
--- a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ClusterModuleZookeeperProviderTestCase.java
+++ /dev/null
@@ -1,77 +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.zookeeper;
-
-import java.io.IOException;
-import java.util.List;
-import org.apache.curator.test.TestingServer;
-import org.apache.skywalking.oap.server.core.cluster.*;
-import org.apache.skywalking.oap.server.core.remote.client.Address;
-import org.apache.skywalking.oap.server.library.module.*;
-import org.junit.*;
-
-/**
- * @author peng-yongsheng
- */
-public class ClusterModuleZookeeperProviderTestCase {
-
- private TestingServer server;
-
- @Before
- public void before() throws Exception {
- server = new TestingServer(12181, true);
- server.start();
- }
-
- @Test
- public void testStart() throws ServiceNotProvidedException, ModuleStartException, ServiceRegisterException, InterruptedException {
- ClusterModuleZookeeperProvider provider = new ClusterModuleZookeeperProvider();
- ClusterModuleZookeeperConfig moduleConfig = (ClusterModuleZookeeperConfig)provider.createConfigBeanIfAbsent();
- moduleConfig.setHostPort(server.getConnectString());
- moduleConfig.setBaseSleepTimeMs(3000);
- moduleConfig.setMaxRetries(4);
-
- provider.prepare();
- provider.start();
-
- ClusterRegister moduleRegister = provider.getService(ClusterRegister.class);
- ClusterNodesQuery clusterNodesQuery = provider.getService(ClusterNodesQuery.class);
-
- RemoteInstance remoteInstance = new RemoteInstance(new Address("ProviderAHost", 1000, true));
-
- moduleRegister.registerRemote(remoteInstance);
-
- for (int i = 0; i < 20; i++) {
- List detailsList = clusterNodesQuery.queryRemoteNodes();
- if (detailsList.size() == 0) {
- Thread.sleep(500);
- continue;
- }
- Assert.assertEquals(1, detailsList.size());
- Assert.assertEquals("ProviderAHost", detailsList.get(0).getAddress().getHost());
- Assert.assertEquals(1000, detailsList.get(0).getAddress().getPort());
- }
-
- }
-
- @After
- public void after() throws IOException {
- server.stop();
- }
-}
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ITClusterModuleZookeeperProviderFunctionalTest.java b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ITClusterModuleZookeeperProviderFunctionalTest.java
new file mode 100644
index 000000000..a3df25ea5
--- /dev/null
+++ b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ITClusterModuleZookeeperProviderFunctionalTest.java
@@ -0,0 +1,235 @@
+/*
+ * 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.zookeeper;
+
+import org.apache.curator.x.discovery.ServiceDiscovery;
+import org.apache.skywalking.apm.util.StringUtil;
+import org.apache.skywalking.oap.server.core.cluster.ClusterNodesQuery;
+import org.apache.skywalking.oap.server.core.cluster.ClusterRegister;
+import org.apache.skywalking.oap.server.core.cluster.RemoteInstance;
+import org.apache.skywalking.oap.server.core.remote.client.Address;
+import org.apache.skywalking.oap.server.library.module.ModuleProvider;
+import org.junit.Before;
+import org.junit.Test;
+import org.powermock.reflect.Whitebox;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+/**
+ * @author zhangwei
+ */
+public class ITClusterModuleZookeeperProviderFunctionalTest {
+
+ private String zkAddress;
+
+ @Before
+ public void before() {
+ zkAddress = System.getProperty("zk.address");
+ assertFalse(StringUtil.isEmpty(zkAddress));
+ }
+
+ @Test
+ public void registerRemote() throws Exception {
+ final String namespace = "register_remote";
+ ModuleProvider provider = createProvider(namespace);
+
+ Address selfAddress = new Address("127.0.0.1", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfInternal() throws Exception {
+ final String namespace = "register_remote_internal";
+ ModuleProvider provider =
+ createProvider(namespace, "127.0.1.2", 1000);
+
+ Address selfAddress = new Address("127.0.0.2", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(provider).registerRemote(instance);
+
+ List remoteInstances = queryRemoteNodes(provider, 1);
+
+ ClusterModuleZookeeperConfig config = (ClusterModuleZookeeperConfig) provider.createConfigBeanIfAbsent();
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(config.getInternalComHost(), queryAddress.getHost());
+ assertEquals(config.getInternalComPort(), queryAddress.getPort());
+ assertTrue(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfReceiver() throws Exception {
+ final String namespace = "register_remote_receiver";
+ ModuleProvider providerA = createProvider(namespace);
+ ModuleProvider providerB = createProvider(namespace);
+
+ // Mixed or Aggregator
+ Address selfAddress = new Address("127.0.0.3", 1000, true);
+ RemoteInstance instance = new RemoteInstance(selfAddress);
+ getClusterRegister(providerA).registerRemote(instance);
+
+ // Receiver
+ List remoteInstances = queryRemoteNodes(providerB, 1);
+ assertEquals(1, remoteInstances.size());
+ Address queryAddress = remoteInstances.get(0).getAddress();
+ assertEquals(selfAddress, queryAddress);
+ assertFalse(queryAddress.isSelf());
+ }
+
+ @Test
+ public void registerRemoteOfCluster() throws Exception {
+ final String namespace = "register_remote_cluster";
+ ModuleProvider providerA = createProvider(namespace);
+ ModuleProvider providerB = createProvider(namespace);
+
+ Address addressA = new Address("127.0.0.4", 1000, true);
+ Address addressB = new Address("127.0.0.5", 1000, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+ }
+
+ @Test
+ public void unregisterRemoteOfCluster() throws Exception {
+ final String namespace = "unregister_remote_cluster";
+ ModuleProvider providerA = createProvider(namespace);
+ ModuleProvider providerB = createProvider(namespace);
+
+ Address addressA = new Address("127.0.0.4", 1000, true);
+ Address addressB = new Address("127.0.0.5", 1000, true);
+
+ RemoteInstance instanceA = new RemoteInstance(addressA);
+ RemoteInstance instanceB = new RemoteInstance(addressB);
+
+ getClusterRegister(providerA).registerRemote(instanceA);
+ getClusterRegister(providerB).registerRemote(instanceB);
+
+ List remoteInstancesOfA = queryRemoteNodes(providerA, 2);
+ validateServiceInstance(addressA, addressB, remoteInstancesOfA);
+
+ List remoteInstancesOfB = queryRemoteNodes(providerB, 2);
+ validateServiceInstance(addressB, addressA, remoteInstancesOfB);
+
+ // unregister A
+ ClusterRegister register = getClusterRegister(providerA);
+ ServiceDiscovery discoveryA = Whitebox.getInternalState(register, "serviceDiscovery");
+ discoveryA.close();
+
+ // only B
+ remoteInstancesOfB = queryRemoteNodes(providerB, 1, 120);
+ assertEquals(1, remoteInstancesOfB.size());
+ Address address = remoteInstancesOfB.get(0).getAddress();
+ assertEquals(address, addressB);
+ assertTrue(addressB.isSelf());
+ }
+
+ private ClusterModuleZookeeperProvider createProvider(String namespace) throws Exception {
+ return createProvider(namespace, null, 0);
+ }
+
+ private ClusterModuleZookeeperProvider createProvider(String namespace, String internalComHost, int internalComPort) throws Exception {
+ ClusterModuleZookeeperProvider provider = new ClusterModuleZookeeperProvider();
+
+ ClusterModuleZookeeperConfig moduleConfig = (ClusterModuleZookeeperConfig) provider.createConfigBeanIfAbsent();
+ moduleConfig.setHostPort(zkAddress);
+ moduleConfig.setBaseSleepTimeMs(3000);
+ moduleConfig.setMaxRetries(3);
+
+ if (!StringUtil.isEmpty(namespace)) {
+ moduleConfig.setNameSpace(namespace);
+ }
+
+ if (!StringUtil.isEmpty(internalComHost)) {
+ moduleConfig.setInternalComHost(internalComHost);
+ }
+
+ if (internalComPort > 0) {
+ moduleConfig.setInternalComPort(internalComPort);
+ }
+
+ provider.prepare();
+ provider.start();
+ provider.notifyAfterCompleted();
+
+ return provider;
+ }
+
+ private ClusterRegister getClusterRegister(ModuleProvider provider) {
+ return provider.getService(ClusterRegister.class);
+ }
+
+ private ClusterNodesQuery getClusterNodesQuery(ModuleProvider provider) {
+ return provider.getService(ClusterNodesQuery.class);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals) throws InterruptedException {
+ return queryRemoteNodes(provider, goals, 20);
+ }
+
+ private List queryRemoteNodes(ModuleProvider provider, int goals, int cyclic) throws InterruptedException {
+ do {
+ List instances = getClusterNodesQuery(provider).queryRemoteNodes();
+ if (instances.size() == goals) {
+ return instances;
+ } else {
+ Thread.sleep(1000);
+ }
+ } while (--cyclic > 0);
+ return Collections.EMPTY_LIST;
+ }
+
+ private void validateServiceInstance(Address selfAddress, Address otherAddress, List queryResult) {
+ assertEquals(2, queryResult.size());
+
+ boolean selfExist = false, otherExist = false;
+
+ for (RemoteInstance instance : queryResult) {
+ Address queryAddress = instance.getAddress();
+ if (queryAddress.equals(selfAddress) && queryAddress.isSelf()) {
+ selfExist = true;
+ } else if (queryAddress.equals(otherAddress) && !queryAddress.isSelf()) {
+ otherExist = true;
+ }
+ }
+
+ assertTrue(selfExist);
+ assertTrue(otherExist);
+ }
+}
diff --git a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinatorTest.java b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinatorTest.java
index cc7e58aef..d10330e35 100644
--- a/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinatorTest.java
+++ b/oap-server/server-cluster-plugin/cluster-zookeeper-plugin/src/test/java/org/apache/skywalking/oap/server/cluster/plugin/zookeeper/ZookeeperCoordinatorTest.java
@@ -52,14 +52,13 @@ public class ZookeeperCoordinatorTest {
@Before
public void setUp() throws Exception {
- config.setHostPort(address.getHost() + ":" + address.getPort());
- coordinator = new ZookeeperCoordinator(config, serviceDiscovery);
- when(serviceDiscovery.serviceCacheBuilder()).thenReturn(cacheBuilder);
when(cacheBuilder.name("remote")).thenReturn(cacheBuilder);
when(cacheBuilder.build()).thenReturn(serviceCache);
doNothing().when(serviceCache).start();
-
doNothing().when(serviceDiscovery).registerService(any());
+ when(serviceDiscovery.serviceCacheBuilder()).thenReturn(cacheBuilder);
+ config.setHostPort(address.getHost() + ":" + address.getPort());
+ coordinator = new ZookeeperCoordinator(config, serviceDiscovery);
}
@Test
diff --git a/oap-server/server-configuration/configuration-zookeeper/pom.xml b/oap-server/server-configuration/configuration-zookeeper/pom.xml
index 9cacd4d3d..70bb80ef8 100644
--- a/oap-server/server-configuration/configuration-zookeeper/pom.xml
+++ b/oap-server/server-configuration/configuration-zookeeper/pom.xml
@@ -26,10 +26,6 @@
4.0.0
configuration-zookeeper
-
-
- 3.5
-