diff --git a/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscovery.java b/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscovery.java
index 223f167b..6f4e0dd1 100644
--- a/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscovery.java
+++ b/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscovery.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2015 the original author or authors.
+ * Copyright 2013-2016 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -19,6 +19,7 @@ package org.springframework.cloud.netflix.turbine;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
+import java.util.Map;
import org.springframework.cloud.client.ServiceInstance;
import org.springframework.cloud.client.discovery.DiscoveryClient;
@@ -33,10 +34,10 @@ import lombok.extern.apachecommons.CommonsLog;
/**
* Class that encapsulates an {@link InstanceDiscovery}
- * implementation that uses Eureka (see https://github.com/Netflix/eureka) The plugin
- * requires a list of applications configured. It then queries the set of instances for
- * each application. Instance information retrieved from Eureka must be translated to
- * something that Turbine can understand i.e the
+ * implementation that uses Spring Cloud Commons (see https://github.com/spring-cloud/spring-cloud-commons)
+ * The plugin requires a list of applications configured. It then queries the set of
+ * instances for * each application. Instance information retrieved from the {@link DiscoveryClient}
+ * must be translated to * something that Turbine can understand i.e the
* {@link Instance} class.
*
* All the logic to perform this translation can be overriden here, so that you can
@@ -48,10 +49,14 @@ import lombok.extern.apachecommons.CommonsLog;
public class CommonsInstanceDiscovery implements InstanceDiscovery {
private static final String DEFAULT_CLUSTER_NAME_EXPRESSION = "serviceId";
+ protected static final String PORT_KEY = "port";
+ protected static final String SECURE_PORT_KEY = "securePort";
+ protected static final String FUSED_HOST_PORT_KEY = "fusedHostPort";
private final Expression clusterNameExpression;
private DiscoveryClient discoveryClient;
private TurbineProperties turbineProperties;
+ private final boolean combineHostPort;
public CommonsInstanceDiscovery(TurbineProperties turbineProperties, DiscoveryClient discoveryClient) {
this(turbineProperties, DEFAULT_CLUSTER_NAME_EXPRESSION);
@@ -67,6 +72,7 @@ public class CommonsInstanceDiscovery implements InstanceDiscovery {
clusterNameExpression = defaultExpression;
}
this.clusterNameExpression = parser.parseExpression(clusterNameExpression);
+ this.combineHostPort = turbineProperties.isCombineHostPort();
}
protected Expression getClusterNameExpression() {
@@ -77,8 +83,12 @@ public class CommonsInstanceDiscovery implements InstanceDiscovery {
return turbineProperties;
}
+ protected boolean isCombineHostPort() {
+ return combineHostPort;
+ }
+
/**
- * Method that queries Eureka service for a list of configured application names
+ * Method that queries DiscoveryClient for a list of configured application names
* @return Collection
*/
@Override
@@ -141,31 +151,23 @@ public class CommonsInstanceDiscovery implements InstanceDiscovery {
/**
* Private helper that marshals the information from each instance into something that
- * Turbine can understand. Override this method for your own implementation for
- * parsing Eureka info.
+ * Turbine can understand. Override this method for your own implementation.
* @param serviceInstance
* @return Instance
*/
- private Instance marshall(ServiceInstance serviceInstance) {
+ Instance marshall(ServiceInstance serviceInstance) {
String hostname = serviceInstance.getHost();
+ String port = String.valueOf(serviceInstance.getPort());
String cluster = getClusterName(serviceInstance);
Boolean status = Boolean.TRUE; //TODO: where to get?
if (hostname != null && cluster != null && status != null) {
- Instance instance = new Instance(hostname, cluster, status);
+ Instance instance = getInstance(hostname, port, cluster, status);
- // TODO: reimplement when metadata is in commons
- // add metadata
- /*Map metadata = instanceInfo.getMetadata();
- if (metadata != null) {
- instance.getAttributes().putAll(metadata);
- }*/
-
- // add ports
- instance.getAttributes().put("port", String.valueOf(serviceInstance.getPort()));
+ Map metadata = serviceInstance.getMetadata();
boolean securePortEnabled = serviceInstance.isSecure();
- if (securePortEnabled) {
- instance.getAttributes().put("securePort", String.valueOf(serviceInstance.getPort()));
- }
+
+ addMetadata(instance, hostname, port, securePortEnabled, port, metadata);
+
return instance;
}
else {
@@ -173,9 +175,31 @@ public class CommonsInstanceDiscovery implements InstanceDiscovery {
}
}
+ protected void addMetadata(Instance instance, String hostname, String port, boolean securePortEnabled, String securePort, Map metadata) {
+ // add metadata
+ if (metadata != null) {
+ instance.getAttributes().putAll(metadata);
+ }
+
+ // add ports
+ instance.getAttributes().put(PORT_KEY, port);
+ if (securePortEnabled) {
+ instance.getAttributes().put(SECURE_PORT_KEY, securePort);
+ }
+ if (this.isCombineHostPort()) {
+ String fusedHostPort = securePortEnabled ? hostname+":"+securePort : instance.getHostname() ;
+ instance.getAttributes().put(FUSED_HOST_PORT_KEY, fusedHostPort);
+ }
+ }
+
+ protected Instance getInstance(String hostname, String port, String cluster, Boolean status) {
+ String hostPart = this.isCombineHostPort() ? hostname+":"+port : hostname;
+ return new Instance(hostPart, cluster, status);
+ }
+
/**
- * Helper that fetches the cluster name. Cluster is a Turbine concept and not a Eureka
- * concept. By default we choose the amazon asg name as the cluster. A custom
+ * Helper that fetches the cluster name. Cluster is a Turbine concept and not a commons
+ * concept. By default we choose the amazon serviceId as the cluster. A custom
* implementation can be plugged in by overriding this method.
*/
protected String getClusterName(Object object) {
diff --git a/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/EurekaInstanceDiscovery.java b/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/EurekaInstanceDiscovery.java
index 32d7b53b..1f2ff8ae 100644
--- a/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/EurekaInstanceDiscovery.java
+++ b/spring-cloud-netflix-turbine/src/main/java/org/springframework/cloud/netflix/turbine/EurekaInstanceDiscovery.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013-2015 the original author or authors.
+ * Copyright 2013-2016 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -47,14 +47,14 @@ import lombok.extern.apachecommons.CommonsLog;
public class EurekaInstanceDiscovery extends CommonsInstanceDiscovery {
private static final String EUREKA_DEFAULT_CLUSTER_NAME_EXPRESSION = "appName";
+ private static final String ASG_KEY = "asg";
private final EurekaClient eurekaClient;
- private final boolean combineHostPort;
+
public EurekaInstanceDiscovery(TurbineProperties turbineProperties, EurekaClient eurekaClient) {
super(turbineProperties, EUREKA_DEFAULT_CLUSTER_NAME_EXPRESSION);
this.eurekaClient = eurekaClient;
- this.combineHostPort = turbineProperties.isCombineHostPort();
}
/**
@@ -104,37 +104,26 @@ public class EurekaInstanceDiscovery extends CommonsInstanceDiscovery {
String cluster = getClusterName(instanceInfo);
Boolean status = parseInstanceStatus(instanceInfo.getStatus());
if (hostname != null && cluster != null && status != null) {
- String hostPart = combineHostPort ? hostname+":"+port : hostname;
- Instance instance = new Instance(hostPart, cluster, status);
+ Instance instance = getInstance(hostname, port, cluster, status);
- // add metadata
Map metadata = instanceInfo.getMetadata();
- if (metadata != null) {
- instance.getAttributes().putAll(metadata);
- }
+ boolean securePortEnabled = instanceInfo.isPortEnabled(InstanceInfo.PortType.SECURE);
+ String securePort = String.valueOf(instanceInfo.getSecurePort());
+
+ addMetadata(instance, hostname, port, securePortEnabled, securePort, metadata);
// add amazon metadata
String asgName = instanceInfo.getASGName();
if (asgName != null) {
- instance.getAttributes().put("asg", asgName);
+ instance.getAttributes().put(ASG_KEY, asgName);
}
+
DataCenterInfo dcInfo = instanceInfo.getDataCenterInfo();
if (dcInfo != null && dcInfo.getName().equals(DataCenterInfo.Name.Amazon)) {
AmazonInfo amznInfo = (AmazonInfo) dcInfo;
instance.getAttributes().putAll(amznInfo.getMetadata());
}
- // add ports
- instance.getAttributes().put("port", String.valueOf(instanceInfo.getPort()));
- boolean securePortEnabled = instanceInfo.isPortEnabled(InstanceInfo.PortType.SECURE);
- if (securePortEnabled) {
- instance.getAttributes().put("securePort", String.valueOf(instanceInfo.getSecurePort()));
- }
-
- if (combineHostPort) {
- String fusedHostPort = securePortEnabled ? hostname+":"+String.valueOf(instanceInfo.getSecurePort()) : hostPart ;
- instance.getAttributes().put("fusedHostPort", fusedHostPort);
- }
return instance;
}
else {
diff --git a/spring-cloud-netflix-turbine/src/test/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscoveryTest.java b/spring-cloud-netflix-turbine/src/test/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscoveryTest.java
new file mode 100644
index 00000000..be44dd30
--- /dev/null
+++ b/spring-cloud-netflix-turbine/src/test/java/org/springframework/cloud/netflix/turbine/CommonsInstanceDiscoveryTest.java
@@ -0,0 +1,146 @@
+/*
+ * Copyright 2013-2016 the original author or authors.
+ *
+ * Licensed 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.springframework.cloud.netflix.turbine;
+
+import java.util.Collections;
+
+import org.junit.Before;
+import org.junit.Test;
+import org.springframework.cloud.client.DefaultServiceInstance;
+import org.springframework.cloud.client.discovery.DiscoveryClient;
+
+import com.netflix.turbine.discovery.Instance;
+
+import static org.junit.Assert.assertEquals;
+import static org.mockito.Mockito.mock;
+
+/**
+ * @author Spencer Gibb
+ */
+public class CommonsInstanceDiscoveryTest {
+
+ private DiscoveryClient discoveryClient;
+ private TurbineProperties turbineProperties;
+
+ @Before
+ public void setUp() throws Exception {
+ this.discoveryClient = mock(DiscoveryClient.class);
+ this.turbineProperties = new TurbineProperties();
+ }
+
+ @Test
+ public void testSecureCombineHostPort() {
+ turbineProperties.setCombineHostPort(true);
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String appName = "testAppName";
+ int port = 8443;
+ String hostName = "myhost";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, hostName, port, true);
+ Instance instance = discovery.marshall(serviceInstance);
+ assertEquals("port is wrong", String.valueOf(port), instance.getAttributes().get("port"));
+ assertEquals("securePort is wrong", String.valueOf(port), instance.getAttributes().get("securePort"));
+
+ String urlPath = SpringClusterMonitor.ClusterConfigBasedUrlClosure.getUrlPath(instance);
+ assertEquals("url is wrong", "https://"+hostName+":"+port+"/hystrix.stream", urlPath);
+ }
+
+ @Test
+ public void testCombineHostPort() {
+ turbineProperties.setCombineHostPort(true);
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String appName = "testAppName";
+ int port = 8080;
+ String hostName = "myhost";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, hostName, port, false);
+ Instance instance = discovery.marshall(serviceInstance);
+ assertEquals("hostname is wrong", hostName+":"+port, instance.getHostname());
+ assertEquals("port is wrong", String.valueOf(port), instance.getAttributes().get("port"));
+
+ String urlPath = SpringClusterMonitor.ClusterConfigBasedUrlClosure.getUrlPath(instance);
+ assertEquals("url is wrong", "http://"+hostName+":"+port+"/hystrix.stream", urlPath);
+
+ String clusterName = discovery.getClusterName(serviceInstance);
+ assertEquals("clusterName is wrong", appName, clusterName);
+ }
+
+ @Test
+ public void testGetClusterName() {
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String appName = "testAppName";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, "myhost", 8080, false);
+ String clusterName = discovery.getClusterName(serviceInstance);
+ assertEquals("clusterName is wrong", appName, clusterName);
+ }
+
+ @Test
+ public void testGetPort() {
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String appName = "testAppName";
+ int port = 8080;
+ String hostName = "myhost";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, hostName, port, false);
+ Instance instance = discovery.marshall(serviceInstance);
+ assertEquals("port is wrong", String.valueOf(port), instance.getAttributes().get("port"));
+
+ String urlPath = SpringClusterMonitor.ClusterConfigBasedUrlClosure.getUrlPath(instance);
+ assertEquals("url is wrong", "http://"+hostName+":"+port+"/hystrix.stream", urlPath);
+ }
+
+ @Test
+ public void testGetSecurePort() {
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String appName = "testAppName";
+ //int port = 8080;
+ int port = 8443;
+ String hostName = "myhost";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, hostName, port, true);
+ Instance instance = discovery.marshall(serviceInstance);
+ assertEquals("port is wrong", String.valueOf(port), instance.getAttributes().get("port"));
+ assertEquals("securePort is wrong", String.valueOf(port), instance.getAttributes().get("securePort"));
+
+ String urlPath = SpringClusterMonitor.ClusterConfigBasedUrlClosure.getUrlPath(instance);
+ assertEquals("url is wrong", "https://"+hostName+":"+port+"/hystrix.stream", urlPath);
+ }
+
+ @Test
+ public void testGetClusterNameCustomExpression() {
+ turbineProperties.setClusterNameExpression("host");
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String appName = "testAppName";
+ String hostName = "myhost";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, hostName, 8080, true);
+ String clusterName = discovery.getClusterName(serviceInstance);
+ assertEquals("clusterName is wrong", hostName, clusterName);
+ }
+
+ @Test
+ public void testGetClusterNameInstanceMetadataMapExpression() {
+ turbineProperties.setClusterNameExpression("metadata['cluster']");
+ CommonsInstanceDiscovery discovery = createDiscovery();
+ String metadataProperty = "myCluster";
+ String appName = "testAppName";
+ String hostName = "myhost";
+ DefaultServiceInstance serviceInstance = new DefaultServiceInstance(appName, hostName, 8080, true, Collections.singletonMap("cluster", metadataProperty));
+ String clusterName = discovery.getClusterName(serviceInstance);
+ assertEquals("clusterName is wrong", metadataProperty, clusterName);
+ }
+
+ private CommonsInstanceDiscovery createDiscovery() {
+ return new CommonsInstanceDiscovery(turbineProperties, discoveryClient);
+ }
+
+}