From 7ec21321625e0d22b4a391a1305fd7ca06af9fc9 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Wed, 11 Mar 2015 15:25:35 +0200 Subject: [PATCH 01/14] Send heartbeats at 2/3 of ttl + List service nodes filtered by alive, when getting the server list --- pom.xml | 21 +-- .../cloud/consul/model/SerfStatusEnum.java | 22 +++ .../ConsulDiscoveryClientConfiguration.java | 5 + .../consul/discovery/ConsulLifecycle.java | 12 +- .../cloud/consul/discovery/ConsulServer.java | 12 ++ .../cloud/consul/discovery/TtlScheduler.java | 126 ++++++++++++++++++ spring-cloud-consul-utils/pom.xml | 39 ++++++ .../consul/alive/AliveFilteringContext.java | 58 ++++++++ .../consul/alive/AliveServerListFilter.java | 37 +++++ .../consul/alive/FilteringAgentClient.java | 21 +++ .../alive/FilteringAgentClientImpl.java | 44 ++++++ .../alive/ServiceCheckServerListFilter.java | 62 +++++++++ 12 files changed, 446 insertions(+), 13 deletions(-) create mode 100644 spring-cloud-consul-core/src/main/java/org/springframework/cloud/consul/model/SerfStatusEnum.java create mode 100644 spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java create mode 100644 spring-cloud-consul-utils/pom.xml create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java diff --git a/pom.xml b/pom.xml index a4b147e7..62f620e2 100644 --- a/pom.xml +++ b/pom.xml @@ -18,12 +18,12 @@ - - https://github.com/spring-cloud/spring-cloud-consul - scm:git:git://github.com/spring-cloud/spring-cloud-consul.git - scm:git:ssh://git@github.com/spring-cloud/spring-cloud-consul.git - HEAD - + + https://github.com/spring-cloud/spring-cloud-consul + scm:git:git://github.com/spring-cloud/spring-cloud-consul.git + scm:git:ssh://git@github.com/spring-cloud/spring-cloud-consul.git + HEAD + spring-cloud-consul-core @@ -32,6 +32,7 @@ spring-cloud-consul-bus spring-cloud-consul-sample spring-cloud-consul-tests + spring-cloud-consul-utils @@ -56,9 +57,9 @@ 1.0.0.BUILD-SNAPSHOT - org.springframework.cloud - spring-cloud-consul-core - 1.0.0.BUILD-SNAPSHOT + org.springframework.cloud + spring-cloud-consul-core + 1.0.0.BUILD-SNAPSHOT org.springframework.cloud @@ -110,7 +111,7 @@ com.ecwid.consul consul-api - 0.1 + 1.0.8 javax.servlet diff --git a/spring-cloud-consul-core/src/main/java/org/springframework/cloud/consul/model/SerfStatusEnum.java b/spring-cloud-consul-core/src/main/java/org/springframework/cloud/consul/model/SerfStatusEnum.java new file mode 100644 index 00000000..53e741b4 --- /dev/null +++ b/spring-cloud-consul-core/src/main/java/org/springframework/cloud/consul/model/SerfStatusEnum.java @@ -0,0 +1,22 @@ +package org.springframework.cloud.consul.model; + +/** + * Gossip pool (serf) statuses. + * Created by nicu on 10.03.2015. + */ +public enum SerfStatusEnum { + StatusAlive(1), + StatusLeaving(2), + StatusLeft(3), + StatusFailed(4); + private final int code; + + SerfStatusEnum(int code) { + this.code=code; + } + + public int getCode() { + return code; + } + +} diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java index 561d93dd..27cb85ae 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java @@ -13,6 +13,11 @@ public class ConsulDiscoveryClientConfiguration { return new ConsulLifecycle(); } + @Bean + public TtlScheduler ttlScheduler() { + return new TtlScheduler(); + } + /*@Bean public ConsulLoadBalancerClient consulLoadBalancerClient() { return new ConsulLoadBalancerClient(); diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java index b1909ca0..7d42a594 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java @@ -19,6 +19,9 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { @Autowired private ConsulProperties consulProperties; + @Autowired + private TtlScheduler ttlScheduler; + @Override protected void register() { NewService service = new NewService(); @@ -30,8 +33,9 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { Integer port = new Integer(getEnvironment().getProperty("server.port", "8080")); service.setPort(port); service.setTags(consulProperties.getTags()); - - //TODO: add support for Check + NewService.Check check = new NewService.Check(); + check.setTtl(ttlScheduler.getTTL() + "s"); + service.setCheck(check); register(service); } @@ -49,6 +53,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { protected void register(NewService service) { log.info("Registering service with consul: {}", service.toString()); client.agentServiceRegister(service); + ttlScheduler.add(service); } @Override @@ -57,7 +62,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { } @Override - protected void deregister(){ + protected void deregister() { deregister(getContext().getId()); } @@ -67,6 +72,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { } private void deregister(String serviceId) { + ttlScheduler.remove(serviceId); client.agentServiceDeregister(serviceId); } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulServer.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulServer.java index 9d1c2b82..a3b7eeff 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulServer.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulServer.java @@ -9,9 +9,13 @@ import com.netflix.loadbalancer.Server; public class ConsulServer extends Server { private final MetaInfo metaInfo; + private final String address; + private final String node; public ConsulServer(final CatalogService service) { super(service.getNode(), service.getServicePort()); + address = service.getAddress(); + node = service.getNode(); metaInfo = new MetaInfo() { @Override public String getAppName() { @@ -39,4 +43,12 @@ public class ConsulServer extends Server { public MetaInfo getMetaInfo() { return metaInfo; } + + public String getAddress() { + return address; + } + + public String getNode() { + return node; + } } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java new file mode 100644 index 00000000..38e799e3 --- /dev/null +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -0,0 +1,126 @@ +package org.springframework.cloud.consul.discovery; + +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.agent.model.NewService; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; + +import javax.annotation.PreDestroy; +import java.util.Collections; +import java.util.Set; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Created by nicu on 11.03.2015. + */ +@Slf4j +public class TtlScheduler implements ApplicationContextAware { + private static final AtomicInteger _instances = new AtomicInteger(); + private static final int DEFAULT_TTL = 3; // must be > 1 + public static final int HEARTBEAT_INTERVAL_RATIO = 2 / 3; + + private final ScheduledExecutorService ttlExecutor; + private volatile ScheduledFuture scheduledFuture; + private final Thread shutdownThread; + private final Set serviceIds; + private final AtomicBoolean isShuttingDown; + private final Runnable ttlPingThread; + private volatile int ttl; + private volatile int heartbeatInterval; + + @Autowired + private ConsulClient client; + + public TtlScheduler() { + if (_instances.addAndGet(1) > 1) { + throw new IllegalStateException("Expecting this to be used in singleton mode!"); + } + //one thread to refresh ttl for each registered app to local agent is ok + ttlExecutor = Executors.newSingleThreadScheduledExecutor(); + shutdownThread = new Thread(new Runnable() { + public void run() { + log.info("Shutting down the Executor Pool for TtlScheduler"); + shutdownExecutorPool(); + } + }); + Runtime.getRuntime().addShutdownHook(shutdownThread); + serviceIds = Collections.newSetFromMap(new ConcurrentHashMap()); + isShuttingDown = new AtomicBoolean(); + ttlPingThread = new TtlPingThread(); + } + + @Override + public void setApplicationContext(ApplicationContext context) throws BeansException { + ttl = context.getEnvironment().getProperty("consul.ttl", Integer.class, DEFAULT_TTL); + ttl = Math.min(2, ttl); + // heartbeat at 2/3 ttl, but no later than ttl -1s and, (under lesser priority), no sooner than 1s from now + heartbeatInterval = Math.max(ttl - 1, Math.min(ttl * HEARTBEAT_INTERVAL_RATIO, 1)); + scheduleTtlHeartbeat(); + } + + /** + * Add a service to the checks loop. + */ + public void add(final NewService service) { + serviceIds.add(service.getId()); + } + + public void remove(String serviceId) { + serviceIds.remove(serviceId); + } + + public int getTTL() { + return ttl; + } + + + private void schedule(int millis, Runnable command) { + ttlExecutor.schedule(command, millis, TimeUnit.MILLISECONDS); + } + + private void scheduleTtlHeartbeat() { + scheduledFuture = ttlExecutor.scheduleAtFixedRate( + ttlPingThread, + 0, heartbeatInterval, + TimeUnit.SECONDS); + } + + private void shutdownExecutorPool() { + isShuttingDown.set(true); + ttlExecutor.shutdown(); + try { + Runtime.getRuntime().removeShutdownHook(shutdownThread); + } catch (IllegalStateException ignored) { + } + } + + class TtlPingThread implements Runnable { + public void run() { + try { + heartbeatServices(); + } catch (Throwable e) { + log.error("Exception while trying to send heartbeat from application to consul local agent in due TTL", e); + } + } + } + + void heartbeatServices() { + for (String serviceId : serviceIds) { + if (!isShuttingDown.get()) { + client.agentCheckPass(serviceId); + } + } + } + + @PreDestroy + public void shutdown() { + if (scheduledFuture != null) { + scheduledFuture.cancel(true); + } + } +} \ No newline at end of file diff --git a/spring-cloud-consul-utils/pom.xml b/spring-cloud-consul-utils/pom.xml new file mode 100644 index 00000000..e883d293 --- /dev/null +++ b/spring-cloud-consul-utils/pom.xml @@ -0,0 +1,39 @@ + + + + spring-cloud-consul + org.springframework.cloud + 1.0.0.BUILD-SNAPSHOT + ../spring-cloud-consul-sample/pom.xml + + 4.0.0 + Spring Cloud Consul Utils + + spring-cloud-consul-utils + + + + com.ecwid.consul + consul-api + + + org.springframework.cloud + spring-cloud-consul-config + + + org.springframework.cloud + spring-cloud-consul-discovery + + + org.springframework.cloud + spring-cloud-consul-bus + + + org.projectlombok + lombok + + + + \ No newline at end of file diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java new file mode 100644 index 00000000..16c1708c --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java @@ -0,0 +1,58 @@ +package org.springframework.cloud.consul.alive; + +import com.netflix.loadbalancer.Server; +import com.netflix.loadbalancer.ServerListFilter; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.discovery.ConsulServer; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.Configuration; + +import java.util.ArrayList; +import java.util.List; + +/** + * Injects a server list filter for giving servers hosting a service only for the live servers per serf status. + * @author nicu marasoiu on 10.03.2015. + */ +@Configuration +@ComponentScan +public class AliveFilteringContext { + @Bean + @Autowired + public ServerListFilter aliveServerListFilter(FilteringAgentClient filteringAgentClient) { + return adapter(new AliveServerListFilter(filteringAgentClient)); + } + + @Bean + @Autowired + public ServerListFilter ttlServerListFilter() { + return adapter(new ServiceCheckServerListFilter()); + } + + private ServerListFilter adapter(final ServerListFilter consulServerList) { + + return new ServerListFilter() { + @Override + public List getFilteredListOfServers(List servers) { + return adapt2(consulServerList.getFilteredListOfServers(adapt1(servers))); + } + + private List adapt1(List servers) { + List consulServers = new ArrayList(servers.size()); + for (Server consulServer : servers) { + consulServers.add((ConsulServer) consulServer); + } + return consulServers; + } + + private List adapt2(List consulServers) { + List servers = new ArrayList(consulServers.size()); + for (ConsulServer consulServer : consulServers) { + servers.add(consulServer); + } + return servers; + } + }; + } +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java new file mode 100644 index 00000000..ba481874 --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java @@ -0,0 +1,37 @@ +package org.springframework.cloud.consul.alive; + +import com.netflix.loadbalancer.ServerListFilter; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.discovery.ConsulServer; + +import java.util.ArrayList; +import java.util.List; +import java.util.Set; + +/** + * Server filter: returns only alive servers. + * Each consul agent runs a serf agent which is a member of the serf gossip pool. + * The serf status (alive/failed/etc) is reflected in 2 consul APIs: in the agent API and in the catalog API. + * We prefer the agent API because it is most up to date (or perhaps we should intersect them and pick members that are liv in both). + * @author nicu marasoiu on 10.03.2015. + */ +public class AliveServerListFilter implements ServerListFilter{ + private FilteringAgentClient filteringAgentClient; + + @Autowired + public AliveServerListFilter(FilteringAgentClient filteringAgentClient) { + this.filteringAgentClient = filteringAgentClient; + } + + @Override + public List getFilteredListOfServers(List servers) { + Set liveNodes = filteringAgentClient.getAliveAgentsAddresses(); + List filteredServers = new ArrayList<>(); + for (ConsulServer server : servers) { + if (liveNodes.contains(server.getAddress())) { + filteredServers.add(server); + } + } + return filteredServers; + } +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java new file mode 100644 index 00000000..f60e6d6d --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java @@ -0,0 +1,21 @@ +package org.springframework.cloud.consul.alive; + +import com.ecwid.consul.v1.agent.model.Member; + +import java.util.List; +import java.util.Set; + +/** + * A CatalogClient which decorates some methods of CatalogClient with filtering, retaining info pertaining to live nodes. + * @author nicu on 10.03.2015. + */ +public interface FilteringAgentClient { + /** + * @return the set of alive gossip pool members (client or server consul agents). + */ + List getAliveAgents(); + /** + * @return the set of alive gossip pool members addresses. + */ + Set getAliveAgentsAddresses(); +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java new file mode 100644 index 00000000..dc5bd6c0 --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java @@ -0,0 +1,44 @@ +package org.springframework.cloud.consul.alive; + +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.agent.model.Member; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.model.SerfStatusEnum; +import org.springframework.stereotype.Service; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +@Service +public class FilteringAgentClientImpl implements FilteringAgentClient { + public static final int ALIVE_STATUS = SerfStatusEnum.StatusAlive.getCode(); + private final ConsulClient client; + + @Autowired + public FilteringAgentClientImpl(ConsulClient client) { + this.client = client; + } + + @Override + public List getAliveAgents() { + List members = client.getAgentMembers().getValue(); + List liveMembers = new ArrayList<>(members.size()); + for (Member peer : members) { + if (peer.getStatus() == ALIVE_STATUS) { + liveMembers.add(peer); + } + } + return liveMembers; + } + + @Override + public Set getAliveAgentsAddresses() { + Set addresses = new HashSet(); + for (Member server : getAliveAgents()) { + addresses.add(server.getAddress()); + } + return addresses; + } +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java new file mode 100644 index 00000000..0d4147a5 --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java @@ -0,0 +1,62 @@ +package org.springframework.cloud.consul.alive; + +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.QueryParams; +import com.ecwid.consul.v1.health.model.Check; +import com.netflix.loadbalancer.ServerListFilter; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.discovery.ConsulServer; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +/** + * Created by nicu on 12.03.2015. + */ +public class ServiceCheckServerListFilter implements ServerListFilter { + + @Autowired + private ConsulClient client; + + @Override + public List getFilteredListOfServers(List servers) { + Set passingServiceIds = getPassingServiceIds(servers); + List okServers = new ArrayList<>(servers.size()); + for(ConsulServer consulServer: servers){ + String serviceId = consulServer.getMetaInfo().getInstanceId(); + if(passingServiceIds.contains(serviceId)){ + List nodeChecks = client.getHealthChecksForNode(consulServer.getNode(), QueryParams.DEFAULT).getValue(); + boolean passingNodeChecks = true; + for(Check check: nodeChecks) { + if(check.getStatus() != Check.CheckStatus.PASSING) { + passingNodeChecks = false; + break; + } + } + if(passingNodeChecks){ + okServers.add(consulServer); + } + } + } + return okServers; + } + + private Set getPassingServiceIds(List servers) { + Set serviceIds = new HashSet<>(1); + for(ConsulServer server:servers){ + serviceIds.add(server.getMetaInfo().getInstanceId()); + } + for(String serviceId: serviceIds) { + List serviceChecks = client.getHealthChecksForService(serviceId, QueryParams.DEFAULT).getValue(); + for(Check check: serviceChecks){ + if(check.getStatus() != Check.CheckStatus.PASSING) { + serviceIds.remove(check.getServiceId()); + } + } + } + return serviceIds; + } + +} From 01013fbfce916bbeb4ecae6d73374b9801be2774 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Thu, 12 Mar 2015 18:09:50 +0200 Subject: [PATCH 02/14] fix parent pom --- spring-cloud-consul-utils/pom.xml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/spring-cloud-consul-utils/pom.xml b/spring-cloud-consul-utils/pom.xml index e883d293..62704da2 100644 --- a/spring-cloud-consul-utils/pom.xml +++ b/spring-cloud-consul-utils/pom.xml @@ -3,10 +3,10 @@ xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> - spring-cloud-consul org.springframework.cloud + spring-cloud-consul 1.0.0.BUILD-SNAPSHOT - ../spring-cloud-consul-sample/pom.xml + .. 4.0.0 Spring Cloud Consul Utils From c3f172d34413c0d4b677012a1c28a39021e7202e Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 10:28:43 +0200 Subject: [PATCH 03/14] add ecwid repo; @Scheduled; other fixes --- pom.xml | 12 +++ spring-cloud-consul-core/pom.xml | 4 + .../cloud/consul/discovery/TtlScheduler.java | 89 ++++--------------- spring-cloud-consul-utils/pom.xml | 4 + 4 files changed, 35 insertions(+), 74 deletions(-) diff --git a/pom.xml b/pom.xml index 62f620e2..33c75656 100644 --- a/pom.xml +++ b/pom.xml @@ -25,6 +25,13 @@ HEAD + + + ecwid + http://nexus.ecwid.com/content/groups/public + + + spring-cloud-consul-core spring-cloud-consul-config @@ -119,6 +126,11 @@ + + com.google.code.gson + gson + 2.3.1 + org.apache.httpcomponents diff --git a/spring-cloud-consul-core/pom.xml b/spring-cloud-consul-core/pom.xml index ed025c00..fb4369c1 100644 --- a/spring-cloud-consul-core/pom.xml +++ b/spring-cloud-consul-core/pom.xml @@ -33,6 +33,10 @@ com.ecwid.consul consul-api + + com.google.code.gson + gson + org.projectlombok lombok diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index 38e799e3..201d13a5 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -7,120 +7,61 @@ import org.springframework.beans.BeansException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; +import org.springframework.scheduling.annotation.Scheduled; -import javax.annotation.PreDestroy; -import java.util.Collections; -import java.util.Set; -import java.util.concurrent.*; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicInteger; /** * Created by nicu on 11.03.2015. */ @Slf4j public class TtlScheduler implements ApplicationContextAware { - private static final AtomicInteger _instances = new AtomicInteger(); private static final int DEFAULT_TTL = 3; // must be > 1 public static final int HEARTBEAT_INTERVAL_RATIO = 2 / 3; - private final ScheduledExecutorService ttlExecutor; - private volatile ScheduledFuture scheduledFuture; - private final Thread shutdownThread; - private final Set serviceIds; - private final AtomicBoolean isShuttingDown; - private final Runnable ttlPingThread; + private final Map serviceHeartbeats = new ConcurrentHashMap<>(); + private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); private volatile int ttl; private volatile int heartbeatInterval; @Autowired private ConsulClient client; - public TtlScheduler() { - if (_instances.addAndGet(1) > 1) { - throw new IllegalStateException("Expecting this to be used in singleton mode!"); - } - //one thread to refresh ttl for each registered app to local agent is ok - ttlExecutor = Executors.newSingleThreadScheduledExecutor(); - shutdownThread = new Thread(new Runnable() { - public void run() { - log.info("Shutting down the Executor Pool for TtlScheduler"); - shutdownExecutorPool(); - } - }); - Runtime.getRuntime().addShutdownHook(shutdownThread); - serviceIds = Collections.newSetFromMap(new ConcurrentHashMap()); - isShuttingDown = new AtomicBoolean(); - ttlPingThread = new TtlPingThread(); - } - @Override public void setApplicationContext(ApplicationContext context) throws BeansException { ttl = context.getEnvironment().getProperty("consul.ttl", Integer.class, DEFAULT_TTL); ttl = Math.min(2, ttl); // heartbeat at 2/3 ttl, but no later than ttl -1s and, (under lesser priority), no sooner than 1s from now heartbeatInterval = Math.max(ttl - 1, Math.min(ttl * HEARTBEAT_INTERVAL_RATIO, 1)); - scheduleTtlHeartbeat(); } /** * Add a service to the checks loop. */ public void add(final NewService service) { - serviceIds.add(service.getId()); + serviceHeartbeats.put(service.getId(), 0L); } public void remove(String serviceId) { - serviceIds.remove(serviceId); + serviceHeartbeats.remove(serviceId); } public int getTTL() { return ttl; } - - private void schedule(int millis, Runnable command) { - ttlExecutor.schedule(command, millis, TimeUnit.MILLISECONDS); - } - - private void scheduleTtlHeartbeat() { - scheduledFuture = ttlExecutor.scheduleAtFixedRate( - ttlPingThread, - 0, heartbeatInterval, - TimeUnit.SECONDS); - } - - private void shutdownExecutorPool() { - isShuttingDown.set(true); - ttlExecutor.shutdown(); - try { - Runtime.getRuntime().removeShutdownHook(shutdownThread); - } catch (IllegalStateException ignored) { - } - } - - class TtlPingThread implements Runnable { - public void run() { - try { - heartbeatServices(); - } catch (Throwable e) { - log.error("Exception while trying to send heartbeat from application to consul local agent in due TTL", e); - } - } - } - + @Scheduled(initialDelay = 0, fixedRate = 100) void heartbeatServices() { - for (String serviceId : serviceIds) { - if (!isShuttingDown.get()) { - client.agentCheckPass(serviceId); + if (heartbeatingNow.compareAndSet(false, true)) { + for (String serviceId : serviceHeartbeats.keySet()) { + long latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); + if(latestHeartbeatDoneForService + heartbeatInterval <= System.currentTimeMillis()) { + client.agentCheckPass(serviceId); + serviceHeartbeats.put(serviceId, System.currentTimeMillis()); + } } } } - - @PreDestroy - public void shutdown() { - if (scheduledFuture != null) { - scheduledFuture.cancel(true); - } - } } \ No newline at end of file diff --git a/spring-cloud-consul-utils/pom.xml b/spring-cloud-consul-utils/pom.xml index 62704da2..dfd085b1 100644 --- a/spring-cloud-consul-utils/pom.xml +++ b/spring-cloud-consul-utils/pom.xml @@ -18,6 +18,10 @@ com.ecwid.consul consul-api + + com.google.code.gson + gson + org.springframework.cloud spring-cloud-consul-config From 4534de04287f081f224fcfbf0c95318c68f5eced Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 10:31:10 +0200 Subject: [PATCH 04/14] revert reformatting --- pom.xml | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pom.xml b/pom.xml index 33c75656..5de0c65a 100644 --- a/pom.xml +++ b/pom.xml @@ -18,12 +18,12 @@ - - https://github.com/spring-cloud/spring-cloud-consul - scm:git:git://github.com/spring-cloud/spring-cloud-consul.git - scm:git:ssh://git@github.com/spring-cloud/spring-cloud-consul.git - HEAD - + + https://github.com/spring-cloud/spring-cloud-consul + scm:git:git://github.com/spring-cloud/spring-cloud-consul.git + scm:git:ssh://git@github.com/spring-cloud/spring-cloud-consul.git + HEAD + From d443a198a2954fbe0d943a6977337581f866b62d Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 10:51:12 +0200 Subject: [PATCH 05/14] @Value for properties --- .../cloud/consul/discovery/TtlScheduler.java | 85 ++++++++++--------- 1 file changed, 46 insertions(+), 39 deletions(-) diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index 201d13a5..db64ef51 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -5,10 +5,13 @@ import com.ecwid.consul.v1.agent.model.NewService; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.BeansException; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.scheduling.annotation.Scheduled; +import javax.annotation.PostConstruct; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; @@ -17,51 +20,55 @@ import java.util.concurrent.atomic.AtomicBoolean; * Created by nicu on 11.03.2015. */ @Slf4j -public class TtlScheduler implements ApplicationContextAware { - private static final int DEFAULT_TTL = 3; // must be > 1 - public static final int HEARTBEAT_INTERVAL_RATIO = 2 / 3; +@ConfigurationProperties +public class TtlScheduler { - private final Map serviceHeartbeats = new ConcurrentHashMap<>(); - private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); - private volatile int ttl; - private volatile int heartbeatInterval; + private final Map serviceHeartbeats = new ConcurrentHashMap<>(); + private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); - @Autowired - private ConsulClient client; + @Value("${consul.ttl:3}") + private volatile int ttl; - @Override - public void setApplicationContext(ApplicationContext context) throws BeansException { - ttl = context.getEnvironment().getProperty("consul.ttl", Integer.class, DEFAULT_TTL); - ttl = Math.min(2, ttl); - // heartbeat at 2/3 ttl, but no later than ttl -1s and, (under lesser priority), no sooner than 1s from now - heartbeatInterval = Math.max(ttl - 1, Math.min(ttl * HEARTBEAT_INTERVAL_RATIO, 1)); + @Value("${consul.heartbeatIntervalRatio:0.6}") + private volatile int heartbeatIntervalRatio; + + private volatile int heartbeatInterval; + + @Autowired + private ConsulClient client; + + @PostConstruct + public void computeHeartbeatInterval() { + // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser priority), no sooner than 1s from now + heartbeatInterval = Math.max(ttl - 1, Math.min(ttl * heartbeatIntervalRatio, 1)); } - /** - * Add a service to the checks loop. - */ - public void add(final NewService service) { - serviceHeartbeats.put(service.getId(), 0L); - } + /** + * Add a service to the checks loop. + */ + public void add(final NewService service) { + serviceHeartbeats.put(service.getId(), 0L); + } - public void remove(String serviceId) { - serviceHeartbeats.remove(serviceId); - } + public void remove(String serviceId) { + serviceHeartbeats.remove(serviceId); + } - public int getTTL() { - return ttl; - } + public int getTTL() { + return ttl; + } - @Scheduled(initialDelay = 0, fixedRate = 100) - void heartbeatServices() { - if (heartbeatingNow.compareAndSet(false, true)) { - for (String serviceId : serviceHeartbeats.keySet()) { - long latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); - if(latestHeartbeatDoneForService + heartbeatInterval <= System.currentTimeMillis()) { - client.agentCheckPass(serviceId); - serviceHeartbeats.put(serviceId, System.currentTimeMillis()); - } - } - } - } + @Scheduled(initialDelay = 0, fixedRate = 100) + private void heartbeatServices() { + if (heartbeatingNow.compareAndSet(false, true)) { + for (String serviceId : serviceHeartbeats.keySet()) { + long latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); + if (latestHeartbeatDoneForService + heartbeatInterval <= System + .currentTimeMillis()) { + client.agentCheckPass(serviceId); + serviceHeartbeats.put(serviceId, System.currentTimeMillis()); + } + } + } + } } \ No newline at end of file From 014c6bb9b5c9bf9991b46e5e3e294244771c109b Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 10:52:50 +0200 Subject: [PATCH 06/14] remove ComponentScan --- .../cloud/consul/alive/AliveFilteringContext.java | 1 - 1 file changed, 1 deletion(-) diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java index 16c1708c..59a7721f 100644 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java @@ -16,7 +16,6 @@ import java.util.List; * @author nicu marasoiu on 10.03.2015. */ @Configuration -@ComponentScan public class AliveFilteringContext { @Bean @Autowired From 91133fd2c185d73792978a8c6bcc97b8fef2b966 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 12:04:44 +0200 Subject: [PATCH 07/14] fix + rename package, format eclipse --- .../cloud/consul/discovery/TtlScheduler.java | 6 +- spring-cloud-consul-utils/pom.xml | 4 +- .../consul/alive/AliveFilteringContext.java | 57 ---------------- .../consul/alive/AliveServerListFilter.java | 37 ----------- .../consul/alive/FilteringAgentClient.java | 21 ------ .../alive/FilteringAgentClientImpl.java | 44 ------------- .../alive/ServiceCheckServerListFilter.java | 62 ------------------ .../AliveServerListFilter.java | 39 +++++++++++ .../FilteringAgentClient.java | 23 +++++++ .../FilteringAgentClientImpl.java | 45 +++++++++++++ .../ServiceCheckServerListFilter.java | 65 +++++++++++++++++++ ...nsulServerListFiltersFilteringContext.java | 61 +++++++++++++++++ 12 files changed, 238 insertions(+), 226 deletions(-) delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java create mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index db64ef51..aaf6c43f 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -29,8 +29,8 @@ public class TtlScheduler { @Value("${consul.ttl:3}") private volatile int ttl; - @Value("${consul.heartbeatIntervalRatio:0.6}") - private volatile int heartbeatIntervalRatio; + @Value("${consul.heartbeatIntervalRatio:0.66}") + private volatile float heartbeatIntervalRatio; private volatile int heartbeatInterval; @@ -40,7 +40,7 @@ public class TtlScheduler { @PostConstruct public void computeHeartbeatInterval() { // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser priority), no sooner than 1s from now - heartbeatInterval = Math.max(ttl - 1, Math.min(ttl * heartbeatIntervalRatio, 1)); + heartbeatInterval = Math.round(Math.max(ttl - 1, Math.min(ttl * heartbeatIntervalRatio, 1))); } /** diff --git a/spring-cloud-consul-utils/pom.xml b/spring-cloud-consul-utils/pom.xml index dfd085b1..256aac90 100644 --- a/spring-cloud-consul-utils/pom.xml +++ b/spring-cloud-consul-utils/pom.xml @@ -1,6 +1,6 @@ - org.springframework.cloud diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java deleted file mode 100644 index 59a7721f..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveFilteringContext.java +++ /dev/null @@ -1,57 +0,0 @@ -package org.springframework.cloud.consul.alive; - -import com.netflix.loadbalancer.Server; -import com.netflix.loadbalancer.ServerListFilter; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.cloud.consul.discovery.ConsulServer; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.ComponentScan; -import org.springframework.context.annotation.Configuration; - -import java.util.ArrayList; -import java.util.List; - -/** - * Injects a server list filter for giving servers hosting a service only for the live servers per serf status. - * @author nicu marasoiu on 10.03.2015. - */ -@Configuration -public class AliveFilteringContext { - @Bean - @Autowired - public ServerListFilter aliveServerListFilter(FilteringAgentClient filteringAgentClient) { - return adapter(new AliveServerListFilter(filteringAgentClient)); - } - - @Bean - @Autowired - public ServerListFilter ttlServerListFilter() { - return adapter(new ServiceCheckServerListFilter()); - } - - private ServerListFilter adapter(final ServerListFilter consulServerList) { - - return new ServerListFilter() { - @Override - public List getFilteredListOfServers(List servers) { - return adapt2(consulServerList.getFilteredListOfServers(adapt1(servers))); - } - - private List adapt1(List servers) { - List consulServers = new ArrayList(servers.size()); - for (Server consulServer : servers) { - consulServers.add((ConsulServer) consulServer); - } - return consulServers; - } - - private List adapt2(List consulServers) { - List servers = new ArrayList(consulServers.size()); - for (ConsulServer consulServer : consulServers) { - servers.add(consulServer); - } - return servers; - } - }; - } -} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java deleted file mode 100644 index ba481874..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/AliveServerListFilter.java +++ /dev/null @@ -1,37 +0,0 @@ -package org.springframework.cloud.consul.alive; - -import com.netflix.loadbalancer.ServerListFilter; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.cloud.consul.discovery.ConsulServer; - -import java.util.ArrayList; -import java.util.List; -import java.util.Set; - -/** - * Server filter: returns only alive servers. - * Each consul agent runs a serf agent which is a member of the serf gossip pool. - * The serf status (alive/failed/etc) is reflected in 2 consul APIs: in the agent API and in the catalog API. - * We prefer the agent API because it is most up to date (or perhaps we should intersect them and pick members that are liv in both). - * @author nicu marasoiu on 10.03.2015. - */ -public class AliveServerListFilter implements ServerListFilter{ - private FilteringAgentClient filteringAgentClient; - - @Autowired - public AliveServerListFilter(FilteringAgentClient filteringAgentClient) { - this.filteringAgentClient = filteringAgentClient; - } - - @Override - public List getFilteredListOfServers(List servers) { - Set liveNodes = filteringAgentClient.getAliveAgentsAddresses(); - List filteredServers = new ArrayList<>(); - for (ConsulServer server : servers) { - if (liveNodes.contains(server.getAddress())) { - filteredServers.add(server); - } - } - return filteredServers; - } -} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java deleted file mode 100644 index f60e6d6d..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClient.java +++ /dev/null @@ -1,21 +0,0 @@ -package org.springframework.cloud.consul.alive; - -import com.ecwid.consul.v1.agent.model.Member; - -import java.util.List; -import java.util.Set; - -/** - * A CatalogClient which decorates some methods of CatalogClient with filtering, retaining info pertaining to live nodes. - * @author nicu on 10.03.2015. - */ -public interface FilteringAgentClient { - /** - * @return the set of alive gossip pool members (client or server consul agents). - */ - List getAliveAgents(); - /** - * @return the set of alive gossip pool members addresses. - */ - Set getAliveAgentsAddresses(); -} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java deleted file mode 100644 index dc5bd6c0..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/FilteringAgentClientImpl.java +++ /dev/null @@ -1,44 +0,0 @@ -package org.springframework.cloud.consul.alive; - -import com.ecwid.consul.v1.ConsulClient; -import com.ecwid.consul.v1.agent.model.Member; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.cloud.consul.model.SerfStatusEnum; -import org.springframework.stereotype.Service; - -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; - -@Service -public class FilteringAgentClientImpl implements FilteringAgentClient { - public static final int ALIVE_STATUS = SerfStatusEnum.StatusAlive.getCode(); - private final ConsulClient client; - - @Autowired - public FilteringAgentClientImpl(ConsulClient client) { - this.client = client; - } - - @Override - public List getAliveAgents() { - List members = client.getAgentMembers().getValue(); - List liveMembers = new ArrayList<>(members.size()); - for (Member peer : members) { - if (peer.getStatus() == ALIVE_STATUS) { - liveMembers.add(peer); - } - } - return liveMembers; - } - - @Override - public Set getAliveAgentsAddresses() { - Set addresses = new HashSet(); - for (Member server : getAliveAgents()) { - addresses.add(server.getAddress()); - } - return addresses; - } -} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java deleted file mode 100644 index 0d4147a5..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/alive/ServiceCheckServerListFilter.java +++ /dev/null @@ -1,62 +0,0 @@ -package org.springframework.cloud.consul.alive; - -import com.ecwid.consul.v1.ConsulClient; -import com.ecwid.consul.v1.QueryParams; -import com.ecwid.consul.v1.health.model.Check; -import com.netflix.loadbalancer.ServerListFilter; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.cloud.consul.discovery.ConsulServer; - -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; - -/** - * Created by nicu on 12.03.2015. - */ -public class ServiceCheckServerListFilter implements ServerListFilter { - - @Autowired - private ConsulClient client; - - @Override - public List getFilteredListOfServers(List servers) { - Set passingServiceIds = getPassingServiceIds(servers); - List okServers = new ArrayList<>(servers.size()); - for(ConsulServer consulServer: servers){ - String serviceId = consulServer.getMetaInfo().getInstanceId(); - if(passingServiceIds.contains(serviceId)){ - List nodeChecks = client.getHealthChecksForNode(consulServer.getNode(), QueryParams.DEFAULT).getValue(); - boolean passingNodeChecks = true; - for(Check check: nodeChecks) { - if(check.getStatus() != Check.CheckStatus.PASSING) { - passingNodeChecks = false; - break; - } - } - if(passingNodeChecks){ - okServers.add(consulServer); - } - } - } - return okServers; - } - - private Set getPassingServiceIds(List servers) { - Set serviceIds = new HashSet<>(1); - for(ConsulServer server:servers){ - serviceIds.add(server.getMetaInfo().getInstanceId()); - } - for(String serviceId: serviceIds) { - List serviceChecks = client.getHealthChecksForService(serviceId, QueryParams.DEFAULT).getValue(); - for(Check check: serviceChecks){ - if(check.getStatus() != Check.CheckStatus.PASSING) { - serviceIds.remove(check.getServiceId()); - } - } - } - return serviceIds; - } - -} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java new file mode 100644 index 00000000..185d4778 --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java @@ -0,0 +1,39 @@ +package org.springframework.cloud.consul.serverlistfilters; + +import java.util.ArrayList; +import java.util.List; +import java.util.Set; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.discovery.ConsulServer; + +import com.netflix.loadbalancer.ServerListFilter; + +/** + * Server filter: returns only alive servers. Each consul agent runs a serf agent which is + * a member of the serf gossip pool. The serf status (alive/failed/etc) is reflected in 2 + * consul APIs: in the agent API and in the catalog API. We prefer the agent API because + * it is most up to date (or perhaps we should intersect them and pick members that are + * liv in both). + * @author nicu marasoiu on 10.03.2015. + */ +public class AliveServerListFilter implements ServerListFilter { + private FilteringAgentClient filteringAgentClient; + + @Autowired + public AliveServerListFilter(FilteringAgentClient filteringAgentClient) { + this.filteringAgentClient = filteringAgentClient; + } + + @Override + public List getFilteredListOfServers(List servers) { + Set liveNodes = filteringAgentClient.getAliveAgentsAddresses(); + List filteredServers = new ArrayList<>(); + for (ConsulServer server : servers) { + if (liveNodes.contains(server.getAddress())) { + filteredServers.add(server); + } + } + return filteredServers; + } +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java new file mode 100644 index 00000000..81e69efe --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java @@ -0,0 +1,23 @@ +package org.springframework.cloud.consul.serverlistfilters; + +import java.util.List; +import java.util.Set; + +import com.ecwid.consul.v1.agent.model.Member; + +/** + * A CatalogClient which decorates some methods of CatalogClient with filtering, retaining + * info pertaining to live nodes. + * @author nicu on 10.03.2015. + */ +public interface FilteringAgentClient { + /** + * @return the set of alive gossip pool members (client or server consul agents). + */ + List getAliveAgents(); + + /** + * @return the set of alive gossip pool members addresses. + */ + Set getAliveAgentsAddresses(); +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java new file mode 100644 index 00000000..e939e2dc --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java @@ -0,0 +1,45 @@ +package org.springframework.cloud.consul.serverlistfilters; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.model.SerfStatusEnum; +import org.springframework.stereotype.Service; + +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.agent.model.Member; + +@Service +public class FilteringAgentClientImpl implements FilteringAgentClient { + public static final int ALIVE_STATUS = SerfStatusEnum.StatusAlive.getCode(); + private final ConsulClient client; + + @Autowired + public FilteringAgentClientImpl(ConsulClient client) { + this.client = client; + } + + @Override + public List getAliveAgents() { + List members = client.getAgentMembers().getValue(); + List liveMembers = new ArrayList<>(members.size()); + for (Member peer : members) { + if (peer.getStatus() == ALIVE_STATUS) { + liveMembers.add(peer); + } + } + return liveMembers; + } + + @Override + public Set getAliveAgentsAddresses() { + Set addresses = new HashSet(); + for (Member server : getAliveAgents()) { + addresses.add(server.getAddress()); + } + return addresses; + } +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java new file mode 100644 index 00000000..b9170c67 --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java @@ -0,0 +1,65 @@ +package org.springframework.cloud.consul.serverlistfilters; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.discovery.ConsulServer; + +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.QueryParams; +import com.ecwid.consul.v1.health.model.Check; +import com.netflix.loadbalancer.ServerListFilter; + +/** + * Created by nicu on 12.03.2015. + */ +public class ServiceCheckServerListFilter implements ServerListFilter { + + @Autowired + private ConsulClient client; + + @Override + public List getFilteredListOfServers(List servers) { + Set passingServiceIds = getPassingServiceIds(servers); + List okServers = new ArrayList<>(servers.size()); + for (ConsulServer consulServer : servers) { + String serviceId = consulServer.getMetaInfo().getInstanceId(); + if (passingServiceIds.contains(serviceId)) { + List nodeChecks = client.getHealthChecksForNode( + consulServer.getNode(), QueryParams.DEFAULT).getValue(); + boolean passingNodeChecks = true; + for (Check check : nodeChecks) { + if (check.getStatus() != Check.CheckStatus.PASSING) { + passingNodeChecks = false; + break; + } + } + if (passingNodeChecks) { + okServers.add(consulServer); + } + } + } + return okServers; + } + + private Set getPassingServiceIds(List servers) { + Set serviceIds = new HashSet<>(1); + for (ConsulServer server : servers) { + serviceIds.add(server.getMetaInfo().getInstanceId()); + } + for (String serviceId : serviceIds) { + List serviceChecks = client.getHealthChecksForService(serviceId, + QueryParams.DEFAULT).getValue(); + for (Check check : serviceChecks) { + if (check.getStatus() != Check.CheckStatus.PASSING) { + serviceIds.remove(check.getServiceId()); + } + } + } + return serviceIds; + } + +} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java new file mode 100644 index 00000000..413ff95b --- /dev/null +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java @@ -0,0 +1,61 @@ +package org.springframework.cloud.consul.serverlistfilters; + +import java.util.ArrayList; +import java.util.List; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.consul.discovery.ConsulServer; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +import com.netflix.loadbalancer.Server; +import com.netflix.loadbalancer.ServerListFilter; + +/** + * Injects a server list filter for giving servers hosting a service only for the live + * servers per serf status. + * @author nicu marasoiu on 10.03.2015. + */ +@Configuration +public class UsefulConsulServerListFiltersFilteringContext { + @Bean + @Autowired + public ServerListFilter aliveServerListFilter( + FilteringAgentClient filteringAgentClient) { + return adapter(new AliveServerListFilter(filteringAgentClient)); + } + + @Bean + @Autowired + public ServerListFilter ttlServerListFilter() { + return adapter(new ServiceCheckServerListFilter()); + } + + private ServerListFilter adapter( + final ServerListFilter consulServerList) { + + return new ServerListFilter() { + @Override + public List getFilteredListOfServers(List servers) { + return adapt2(consulServerList.getFilteredListOfServers(adapt1(servers))); + } + + private List adapt1(List servers) { + List consulServers = new ArrayList( + servers.size()); + for (Server consulServer : servers) { + consulServers.add((ConsulServer) consulServer); + } + return consulServers; + } + + private List adapt2(List consulServers) { + List servers = new ArrayList(consulServers.size()); + for (ConsulServer consulServer : consulServers) { + servers.add(consulServer); + } + return servers; + } + }; + } +} From a2879eae4c2e6df02960116def504cf6aec09703 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 12:09:21 +0200 Subject: [PATCH 08/14] cleanup --- pom.xml | 8 +++---- .../cloud/consul/discovery/TtlScheduler.java | 24 +++++++++---------- 2 files changed, 16 insertions(+), 16 deletions(-) diff --git a/pom.xml b/pom.xml index 5de0c65a..2b7b2ef1 100644 --- a/pom.xml +++ b/pom.xml @@ -64,9 +64,9 @@ 1.0.0.BUILD-SNAPSHOT - org.springframework.cloud - spring-cloud-consul-core - 1.0.0.BUILD-SNAPSHOT + org.springframework.cloud + spring-cloud-consul-core + 1.0.0.BUILD-SNAPSHOT org.springframework.cloud @@ -118,7 +118,7 @@ com.ecwid.consul consul-api - 1.0.8 + 1.0 javax.servlet diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index aaf6c43f..e743c0ca 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -1,21 +1,21 @@ package org.springframework.cloud.consul.discovery; -import com.ecwid.consul.v1.ConsulClient; -import com.ecwid.consul.v1.agent.model.NewService; -import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.BeansException; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.boot.context.properties.ConfigurationProperties; -import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; -import org.springframework.scheduling.annotation.Scheduled; - -import javax.annotation.PostConstruct; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; +import javax.annotation.PostConstruct; + +import lombok.extern.slf4j.Slf4j; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.scheduling.annotation.Scheduled; + +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.agent.model.NewService; + /** * Created by nicu on 11.03.2015. */ From 88d98de7ae8f34d117ed23f2268c688cb0722264 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 12:37:38 +0200 Subject: [PATCH 09/14] refactor Ttl scheduler with joda; use ConfProperties fully --- pom.xml | 7 ++- spring-cloud-consul-discovery/pom.xml | 6 ++- .../cloud/consul/discovery/TtlScheduler.java | 45 ++++++++++++------- 3 files changed, 39 insertions(+), 19 deletions(-) diff --git a/pom.xml b/pom.xml index 2b7b2ef1..3132beac 100644 --- a/pom.xml +++ b/pom.xml @@ -118,7 +118,7 @@ com.ecwid.consul consul-api - 1.0 + 1.0.8 javax.servlet @@ -177,6 +177,11 @@ ribbon-httpclient ${ribbon.version} + + joda-time + joda-time + 2.7 + org.projectlombok lombok diff --git a/spring-cloud-consul-discovery/pom.xml b/spring-cloud-consul-discovery/pom.xml index df57be82..a6695ec4 100644 --- a/spring-cloud-consul-discovery/pom.xml +++ b/spring-cloud-consul-discovery/pom.xml @@ -47,6 +47,10 @@ spring-boot-starter-test test - + + joda-time + joda-time + + diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index e743c0ca..1fdc96a8 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -5,9 +5,15 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; import javax.annotation.PostConstruct; +import javax.validation.constraints.DecimalMax; +import javax.validation.constraints.DecimalMin; +import javax.validation.constraints.Max; +import javax.validation.constraints.Min; import lombok.extern.slf4j.Slf4j; +import org.joda.time.DateTime; +import org.joda.time.Period; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.context.properties.ConfigurationProperties; @@ -20,34 +26,39 @@ import com.ecwid.consul.v1.agent.model.NewService; * Created by nicu on 11.03.2015. */ @Slf4j -@ConfigurationProperties +@ConfigurationProperties(prefix = "consul.heartbeat") public class TtlScheduler { - private final Map serviceHeartbeats = new ConcurrentHashMap<>(); + public static final DateTime EXPIRED_DATE = new DateTime(0); + private final Map serviceHeartbeats = new ConcurrentHashMap<>(); private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); - @Value("${consul.ttl:3}") - private volatile int ttl; + @Min(1) + @Max(10) + private volatile int ttl = 3; - @Value("${consul.heartbeatIntervalRatio:0.66}") - private volatile float heartbeatIntervalRatio; + @DecimalMin("0.1") + @DecimalMax("0.9") + private volatile double intervalRatio = 2.0/3.0; - private volatile int heartbeatInterval; + private volatile Period heartbeatInterval; @Autowired private ConsulClient client; @PostConstruct - public void computeHeartbeatInterval() { - // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser priority), no sooner than 1s from now - heartbeatInterval = Math.round(Math.max(ttl - 1, Math.min(ttl * heartbeatIntervalRatio, 1))); - } + public void computeHeartbeatInterval() { + // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser + // priority), no sooner than 1s from now + heartbeatInterval = new Period(Math.round(1000 * Math.max(ttl - 1, + Math.min(ttl * intervalRatio, 1)))); + } /** * Add a service to the checks loop. */ public void add(final NewService service) { - serviceHeartbeats.put(service.getId(), 0L); + serviceHeartbeats.put(service.getId(), EXPIRED_DATE); } public void remove(String serviceId) { @@ -58,15 +69,15 @@ public class TtlScheduler { return ttl; } - @Scheduled(initialDelay = 0, fixedRate = 100) + @Scheduled(initialDelay = 0, fixedRateString = "${consul.heartbeat.fixedRate:201000}") private void heartbeatServices() { if (heartbeatingNow.compareAndSet(false, true)) { for (String serviceId : serviceHeartbeats.keySet()) { - long latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); - if (latestHeartbeatDoneForService + heartbeatInterval <= System - .currentTimeMillis()) { + DateTime latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); + if (latestHeartbeatDoneForService.plus(heartbeatInterval).isBefore( + new DateTime())) { client.agentCheckPass(serviceId); - serviceHeartbeats.put(serviceId, System.currentTimeMillis()); + serviceHeartbeats.put(serviceId, new DateTime()); } } } From 893a6b0fd2b9495c01d926c776ae7a6fc3cdf50c Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 14:09:16 +0200 Subject: [PATCH 10/14] injecting proper instances in example configuration to import --- ...efulConsulServerListFiltersFilteringContext.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java index 413ff95b..ea7eb553 100644 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java +++ b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java @@ -3,6 +3,7 @@ package org.springframework.cloud.consul.serverlistfilters; import java.util.ArrayList; import java.util.List; +import com.ecwid.consul.v1.ConsulClient; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.consul.discovery.ConsulServer; import org.springframework.context.annotation.Bean; @@ -31,7 +32,17 @@ public class UsefulConsulServerListFiltersFilteringContext { return adapter(new ServiceCheckServerListFilter()); } - private ServerListFilter adapter( + @Bean + public FilteringAgentClient filteringAgentClient(){ + return new FilteringAgentClientImpl(consulClient()); + } + + @Bean + public ConsulClient consulClient() { + return new ConsulClient(); + } + + private ServerListFilter adapter( final ServerListFilter consulServerList) { return new ServerListFilter() { From 12d3695d7aae26418f2302be839662da24a79c76 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 16:12:09 +0200 Subject: [PATCH 11/14] cleanup: use more appropriately ConfigurationProperties making a DTO for TTL heartbeating --- .../ConsulDiscoveryClientConfiguration.java | 5 + .../consul/discovery/ConsulLifecycle.java | 5 +- .../discovery/TtlHeartbeatConfiguration.java | 37 +++++++ .../cloud/consul/discovery/TtlScheduler.java | 103 ++++++------------ 4 files changed, 82 insertions(+), 68 deletions(-) create mode 100644 spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java index 27cb85ae..54467306 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java @@ -18,6 +18,11 @@ public class ConsulDiscoveryClientConfiguration { return new TtlScheduler(); } + @Bean + public TtlHeartbeatConfiguration ttlHeartbeatConfiguration() { + return new TtlHeartbeatConfiguration(); + } + /*@Bean public ConsulLoadBalancerClient consulLoadBalancerClient() { return new ConsulLoadBalancerClient(); diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java index 7d42a594..3b30f52a 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java @@ -22,6 +22,9 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { @Autowired private TtlScheduler ttlScheduler; + @Autowired + private TtlHeartbeatConfiguration ttlConfig; + @Override protected void register() { NewService service = new NewService(); @@ -34,7 +37,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { service.setPort(port); service.setTags(consulProperties.getTags()); NewService.Check check = new NewService.Check(); - check.setTtl(ttlScheduler.getTTL() + "s"); + check.setTtl(ttlConfig.getTtlAsString()); service.setCheck(check); register(service); } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java new file mode 100644 index 00000000..a1120d9a --- /dev/null +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java @@ -0,0 +1,37 @@ +package org.springframework.cloud.consul.discovery; + +import lombok.Data; +import org.joda.time.Period; +import org.springframework.boot.context.properties.ConfigurationProperties; + +import javax.annotation.PostConstruct; +import javax.validation.constraints.DecimalMax; +import javax.validation.constraints.DecimalMin; +import javax.validation.constraints.Max; +import javax.validation.constraints.Min; + +@ConfigurationProperties(prefix = "consul.heartbeat") +@Data +public class TtlHeartbeatConfiguration { + @Min(1) + @Max(10) + private volatile int ttl = 3; + + @DecimalMin("0.1") + @DecimalMax("0.9") + private volatile double intervalRatio = 2.0 / 3.0; + + private volatile Period heartbeatInterval; + + @PostConstruct + public void computeHeartbeatInterval() { + // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser + // priority), no sooner than 1s from now + heartbeatInterval = new Period(Math.round(1000 * Math.max(ttl - 1, + Math.min(ttl * intervalRatio, 1)))); + } + + public String getTtlAsString() { + return ttl + "s"; + } +} diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index 1fdc96a8..7dd2be1a 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -1,85 +1,54 @@ package org.springframework.cloud.consul.discovery; +import com.ecwid.consul.v1.ConsulClient; +import com.ecwid.consul.v1.agent.model.NewService; +import lombok.extern.slf4j.Slf4j; +import org.joda.time.DateTime; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; + import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; -import javax.annotation.PostConstruct; -import javax.validation.constraints.DecimalMax; -import javax.validation.constraints.DecimalMin; -import javax.validation.constraints.Max; -import javax.validation.constraints.Min; - -import lombok.extern.slf4j.Slf4j; - -import org.joda.time.DateTime; -import org.joda.time.Period; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.boot.context.properties.ConfigurationProperties; -import org.springframework.scheduling.annotation.Scheduled; - -import com.ecwid.consul.v1.ConsulClient; -import com.ecwid.consul.v1.agent.model.NewService; - /** * Created by nicu on 11.03.2015. */ @Slf4j -@ConfigurationProperties(prefix = "consul.heartbeat") public class TtlScheduler { - public static final DateTime EXPIRED_DATE = new DateTime(0); - private final Map serviceHeartbeats = new ConcurrentHashMap<>(); - private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); + public static final DateTime EXPIRED_DATE = new DateTime(0); + private final Map serviceHeartbeats = new ConcurrentHashMap<>(); + private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); - @Min(1) - @Max(10) - private volatile int ttl = 3; + @Autowired + private TtlHeartbeatConfiguration configuration; - @DecimalMin("0.1") - @DecimalMax("0.9") - private volatile double intervalRatio = 2.0/3.0; + @Autowired + private ConsulClient client; - private volatile Period heartbeatInterval; + /** + * Add a service to the checks loop. + */ + public void add(final NewService service) { + serviceHeartbeats.put(service.getId(), EXPIRED_DATE); + } - @Autowired - private ConsulClient client; + public void remove(String serviceId) { + serviceHeartbeats.remove(serviceId); + } - @PostConstruct - public void computeHeartbeatInterval() { - // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser - // priority), no sooner than 1s from now - heartbeatInterval = new Period(Math.round(1000 * Math.max(ttl - 1, - Math.min(ttl * intervalRatio, 1)))); - } - - /** - * Add a service to the checks loop. - */ - public void add(final NewService service) { - serviceHeartbeats.put(service.getId(), EXPIRED_DATE); - } - - public void remove(String serviceId) { - serviceHeartbeats.remove(serviceId); - } - - public int getTTL() { - return ttl; - } - - @Scheduled(initialDelay = 0, fixedRateString = "${consul.heartbeat.fixedRate:201000}") - private void heartbeatServices() { - if (heartbeatingNow.compareAndSet(false, true)) { - for (String serviceId : serviceHeartbeats.keySet()) { - DateTime latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); - if (latestHeartbeatDoneForService.plus(heartbeatInterval).isBefore( - new DateTime())) { - client.agentCheckPass(serviceId); - serviceHeartbeats.put(serviceId, new DateTime()); - } - } - } - } + @Scheduled(initialDelay = 0, fixedRateString = "${consul.heartbeat.fixedRate:201000}") + private void heartbeatServices() { + if (heartbeatingNow.compareAndSet(false, true)) { + for (String serviceId : serviceHeartbeats.keySet()) { + DateTime latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); + if (latestHeartbeatDoneForService.plus(configuration.getHeartbeatInterval()) + .isBefore(DateTime.now())) { + client.agentCheckPass(serviceId); + serviceHeartbeats.put(serviceId, DateTime.now()); + } + } + } + } } \ No newline at end of file From c7301f55e4f338503d6594842860998e075e7da6 Mon Sep 17 00:00:00 2001 From: nicu marasoiu Date: Fri, 13 Mar 2015 16:15:28 +0200 Subject: [PATCH 12/14] rename to TtlProperties --- .../consul/discovery/ConsulDiscoveryClientConfiguration.java | 4 ++-- .../cloud/consul/discovery/ConsulLifecycle.java | 2 +- ...eartbeatConfiguration.java => TtlHeartbeatProperties.java} | 2 +- .../springframework/cloud/consul/discovery/TtlScheduler.java | 2 +- 4 files changed, 5 insertions(+), 5 deletions(-) rename spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/{TtlHeartbeatConfiguration.java => TtlHeartbeatProperties.java} (96%) diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java index 54467306..9301a11a 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java @@ -19,8 +19,8 @@ public class ConsulDiscoveryClientConfiguration { } @Bean - public TtlHeartbeatConfiguration ttlHeartbeatConfiguration() { - return new TtlHeartbeatConfiguration(); + public TtlHeartbeatProperties ttlHeartbeatConfiguration() { + return new TtlHeartbeatProperties(); } /*@Bean diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java index 3b30f52a..8a4929e6 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java @@ -23,7 +23,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { private TtlScheduler ttlScheduler; @Autowired - private TtlHeartbeatConfiguration ttlConfig; + private TtlHeartbeatProperties ttlConfig; @Override protected void register() { diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatProperties.java similarity index 96% rename from spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java rename to spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatProperties.java index a1120d9a..e3b78d48 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatProperties.java @@ -12,7 +12,7 @@ import javax.validation.constraints.Min; @ConfigurationProperties(prefix = "consul.heartbeat") @Data -public class TtlHeartbeatConfiguration { +public class TtlHeartbeatProperties { @Min(1) @Max(10) private volatile int ttl = 3; diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index 7dd2be1a..e5acb87c 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -22,7 +22,7 @@ public class TtlScheduler { private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); @Autowired - private TtlHeartbeatConfiguration configuration; + private TtlHeartbeatProperties configuration; @Autowired private ConsulClient client; From e35ac90108257790c5f653840436c63c25c3105d Mon Sep 17 00:00:00 2001 From: Spencer Gibb Date: Fri, 13 Mar 2015 13:33:11 -0600 Subject: [PATCH 13/14] move new filters to discovery, use ServiceCheckServerListFilter as the default filter --- pom.xml | 2 +- spring-cloud-consul-core/pom.xml | 1 + spring-cloud-consul-discovery/pom.xml | 1 - .../ConsulRibbonClientConfiguration.java | 6 ++ .../filters}/AliveServerListFilter.java | 4 +- .../filters/FilteringAgentClient.java | 18 ++--- .../ServiceCheckServerListFilter.java | 8 ++- spring-cloud-consul-utils/pom.xml | 43 ----------- .../FilteringAgentClient.java | 23 ------ ...nsulServerListFiltersFilteringContext.java | 72 ------------------- 10 files changed, 21 insertions(+), 157 deletions(-) rename {spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters => spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters}/AliveServerListFilter.java (90%) rename spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java => spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/FilteringAgentClient.java (62%) rename {spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters => spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters}/ServiceCheckServerListFilter.java (92%) delete mode 100644 spring-cloud-consul-utils/pom.xml delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java delete mode 100644 spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java diff --git a/pom.xml b/pom.xml index bfe7d6c4..5b5717dd 100644 --- a/pom.xml +++ b/pom.xml @@ -39,7 +39,6 @@ spring-cloud-consul-bus spring-cloud-consul-sample spring-cloud-consul-tests - spring-cloud-consul-utils docs @@ -127,6 +126,7 @@ + com.google.code.gson gson diff --git a/spring-cloud-consul-core/pom.xml b/spring-cloud-consul-core/pom.xml index fb4369c1..05052729 100644 --- a/spring-cloud-consul-core/pom.xml +++ b/spring-cloud-consul-core/pom.xml @@ -33,6 +33,7 @@ com.ecwid.consul consul-api + com.google.code.gson gson diff --git a/spring-cloud-consul-discovery/pom.xml b/spring-cloud-consul-discovery/pom.xml index 7cf5ed45..e92ddc87 100644 --- a/spring-cloud-consul-discovery/pom.xml +++ b/spring-cloud-consul-discovery/pom.xml @@ -51,7 +51,6 @@ joda-time joda-time - org.springframework.boot spring-boot-starter-web diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java index 96913c10..b173a84e 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java @@ -24,6 +24,7 @@ import javax.annotation.PostConstruct; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.cloud.consul.discovery.filters.ServiceCheckServerListFilter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -69,6 +70,11 @@ public class ConsulRibbonClientConfiguration { return serverList; } + @Bean + public ServiceCheckServerListFilter ribbonServerListFilter() { + return new ServiceCheckServerListFilter(client); + } + @PostConstruct public void preprocess() { setProp(this.serviceId, DeploymentContextBasedVipAddresses.key(), this.serviceId); diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java similarity index 90% rename from spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java rename to spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java index 185d4778..423c0d27 100644 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/AliveServerListFilter.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java @@ -1,10 +1,9 @@ -package org.springframework.cloud.consul.serverlistfilters; +package org.springframework.cloud.consul.discovery.filters; import java.util.ArrayList; import java.util.List; import java.util.Set; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.consul.discovery.ConsulServer; import com.netflix.loadbalancer.ServerListFilter; @@ -20,7 +19,6 @@ import com.netflix.loadbalancer.ServerListFilter; public class AliveServerListFilter implements ServerListFilter { private FilteringAgentClient filteringAgentClient; - @Autowired public AliveServerListFilter(FilteringAgentClient filteringAgentClient) { this.filteringAgentClient = filteringAgentClient; } diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/FilteringAgentClient.java similarity index 62% rename from spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java rename to spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/FilteringAgentClient.java index e939e2dc..cb76c8d1 100644 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClientImpl.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/FilteringAgentClient.java @@ -1,28 +1,25 @@ -package org.springframework.cloud.consul.serverlistfilters; +package org.springframework.cloud.consul.discovery.filters; import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.consul.model.SerfStatusEnum; -import org.springframework.stereotype.Service; import com.ecwid.consul.v1.ConsulClient; import com.ecwid.consul.v1.agent.model.Member; -@Service -public class FilteringAgentClientImpl implements FilteringAgentClient { - public static final int ALIVE_STATUS = SerfStatusEnum.StatusAlive.getCode(); +public class FilteringAgentClient { + + private static final int ALIVE_STATUS = SerfStatusEnum.StatusAlive.getCode(); + private final ConsulClient client; - @Autowired - public FilteringAgentClientImpl(ConsulClient client) { + public FilteringAgentClient(ConsulClient client) { this.client = client; } - @Override public List getAliveAgents() { List members = client.getAgentMembers().getValue(); List liveMembers = new ArrayList<>(members.size()); @@ -34,9 +31,8 @@ public class FilteringAgentClientImpl implements FilteringAgentClient { return liveMembers; } - @Override public Set getAliveAgentsAddresses() { - Set addresses = new HashSet(); + Set addresses = new HashSet<>(); for (Member server : getAliveAgents()) { addresses.add(server.getAddress()); } diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java similarity index 92% rename from spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java rename to spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java index b9170c67..6d2c5e5f 100644 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/ServiceCheckServerListFilter.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java @@ -1,11 +1,10 @@ -package org.springframework.cloud.consul.serverlistfilters; +package org.springframework.cloud.consul.discovery.filters; import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cloud.consul.discovery.ConsulServer; import com.ecwid.consul.v1.ConsulClient; @@ -18,9 +17,12 @@ import com.netflix.loadbalancer.ServerListFilter; */ public class ServiceCheckServerListFilter implements ServerListFilter { - @Autowired private ConsulClient client; + public ServiceCheckServerListFilter(ConsulClient client) { + this.client = client; + } + @Override public List getFilteredListOfServers(List servers) { Set passingServiceIds = getPassingServiceIds(servers); diff --git a/spring-cloud-consul-utils/pom.xml b/spring-cloud-consul-utils/pom.xml deleted file mode 100644 index 256aac90..00000000 --- a/spring-cloud-consul-utils/pom.xml +++ /dev/null @@ -1,43 +0,0 @@ - - - - org.springframework.cloud - spring-cloud-consul - 1.0.0.BUILD-SNAPSHOT - .. - - 4.0.0 - Spring Cloud Consul Utils - - spring-cloud-consul-utils - - - - com.ecwid.consul - consul-api - - - com.google.code.gson - gson - - - org.springframework.cloud - spring-cloud-consul-config - - - org.springframework.cloud - spring-cloud-consul-discovery - - - org.springframework.cloud - spring-cloud-consul-bus - - - org.projectlombok - lombok - - - - \ No newline at end of file diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java deleted file mode 100644 index 81e69efe..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/FilteringAgentClient.java +++ /dev/null @@ -1,23 +0,0 @@ -package org.springframework.cloud.consul.serverlistfilters; - -import java.util.List; -import java.util.Set; - -import com.ecwid.consul.v1.agent.model.Member; - -/** - * A CatalogClient which decorates some methods of CatalogClient with filtering, retaining - * info pertaining to live nodes. - * @author nicu on 10.03.2015. - */ -public interface FilteringAgentClient { - /** - * @return the set of alive gossip pool members (client or server consul agents). - */ - List getAliveAgents(); - - /** - * @return the set of alive gossip pool members addresses. - */ - Set getAliveAgentsAddresses(); -} diff --git a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java b/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java deleted file mode 100644 index ea7eb553..00000000 --- a/spring-cloud-consul-utils/src/main/java/org/springframework/cloud/consul/serverlistfilters/UsefulConsulServerListFiltersFilteringContext.java +++ /dev/null @@ -1,72 +0,0 @@ -package org.springframework.cloud.consul.serverlistfilters; - -import java.util.ArrayList; -import java.util.List; - -import com.ecwid.consul.v1.ConsulClient; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.cloud.consul.discovery.ConsulServer; -import org.springframework.context.annotation.Bean; -import org.springframework.context.annotation.Configuration; - -import com.netflix.loadbalancer.Server; -import com.netflix.loadbalancer.ServerListFilter; - -/** - * Injects a server list filter for giving servers hosting a service only for the live - * servers per serf status. - * @author nicu marasoiu on 10.03.2015. - */ -@Configuration -public class UsefulConsulServerListFiltersFilteringContext { - @Bean - @Autowired - public ServerListFilter aliveServerListFilter( - FilteringAgentClient filteringAgentClient) { - return adapter(new AliveServerListFilter(filteringAgentClient)); - } - - @Bean - @Autowired - public ServerListFilter ttlServerListFilter() { - return adapter(new ServiceCheckServerListFilter()); - } - - @Bean - public FilteringAgentClient filteringAgentClient(){ - return new FilteringAgentClientImpl(consulClient()); - } - - @Bean - public ConsulClient consulClient() { - return new ConsulClient(); - } - - private ServerListFilter adapter( - final ServerListFilter consulServerList) { - - return new ServerListFilter() { - @Override - public List getFilteredListOfServers(List servers) { - return adapt2(consulServerList.getFilteredListOfServers(adapt1(servers))); - } - - private List adapt1(List servers) { - List consulServers = new ArrayList( - servers.size()); - for (Server consulServer : servers) { - consulServers.add((ConsulServer) consulServer); - } - return consulServers; - } - - private List adapt2(List consulServers) { - List servers = new ArrayList(consulServers.size()); - for (ConsulServer consulServer : consulServers) { - servers.add(consulServer); - } - return servers; - } - }; - } -} From d8a6e7e9adb352f6fe0ed7f815dbbb9a4eab4f27 Mon Sep 17 00:00:00 2001 From: Spencer Gibb Date: Fri, 13 Mar 2015 15:19:38 -0600 Subject: [PATCH 14/14] cleanup ttl implementation --- .../ConsulDiscoveryClientConfiguration.java | 17 +++++----- .../consul/discovery/ConsulLifecycle.java | 4 +-- .../ConsulRibbonClientConfiguration.java | 4 ++- ...operties.java => HeartbeatProperties.java} | 31 ++++++++++--------- .../cloud/consul/discovery/TtlScheduler.java | 30 ++++++++++-------- .../filters/AliveServerListFilter.java | 12 ++++--- .../filters/ServiceCheckServerListFilter.java | 18 ++++++----- 7 files changed, 65 insertions(+), 51 deletions(-) rename spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/{TtlHeartbeatProperties.java => HeartbeatProperties.java} (66%) diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java index 9301a11a..70e2e4e4 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulDiscoveryClientConfiguration.java @@ -1,5 +1,7 @@ package org.springframework.cloud.consul.discovery; +import com.ecwid.consul.v1.ConsulClient; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -8,6 +10,10 @@ import org.springframework.context.annotation.Configuration; */ @Configuration public class ConsulDiscoveryClientConfiguration { + + @Autowired + private ConsulClient consulClient; + @Bean public ConsulLifecycle consulLifecycle() { return new ConsulLifecycle(); @@ -15,19 +21,14 @@ public class ConsulDiscoveryClientConfiguration { @Bean public TtlScheduler ttlScheduler() { - return new TtlScheduler(); + return new TtlScheduler(heartbeatProperties(), consulClient); } @Bean - public TtlHeartbeatProperties ttlHeartbeatConfiguration() { - return new TtlHeartbeatProperties(); + public HeartbeatProperties heartbeatProperties() { + return new HeartbeatProperties(); } - /*@Bean - public ConsulLoadBalancerClient consulLoadBalancerClient() { - return new ConsulLoadBalancerClient(); - }*/ - @Bean public ConsulDiscoveryClient consulDiscoveryClient() { return new ConsulDiscoveryClient(); diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java index 8a4929e6..606e754d 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulLifecycle.java @@ -23,7 +23,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { private TtlScheduler ttlScheduler; @Autowired - private TtlHeartbeatProperties ttlConfig; + private HeartbeatProperties ttlConfig; @Override protected void register() { @@ -37,7 +37,7 @@ public class ConsulLifecycle extends AbstractDiscoveryLifecycle { service.setPort(port); service.setTags(consulProperties.getTags()); NewService.Check check = new NewService.Check(); - check.setTtl(ttlConfig.getTtlAsString()); + check.setTtl(ttlConfig.getTtl()); service.setCheck(check); register(service); } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java index b173a84e..00a264ba 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/ConsulRibbonClientConfiguration.java @@ -21,6 +21,8 @@ import static com.netflix.client.config.CommonClientConfigKey.EnableZoneAffinity import javax.annotation.PostConstruct; +import com.netflix.loadbalancer.Server; +import com.netflix.loadbalancer.ServerListFilter; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; @@ -71,7 +73,7 @@ public class ConsulRibbonClientConfiguration { } @Bean - public ServiceCheckServerListFilter ribbonServerListFilter() { + public ServerListFilter ribbonServerListFilter() { return new ServiceCheckServerListFilter(client); } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatProperties.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/HeartbeatProperties.java similarity index 66% rename from spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatProperties.java rename to spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/HeartbeatProperties.java index e3b78d48..ab0ba03e 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlHeartbeatProperties.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/HeartbeatProperties.java @@ -1,37 +1,40 @@ package org.springframework.cloud.consul.discovery; -import lombok.Data; -import org.joda.time.Period; -import org.springframework.boot.context.properties.ConfigurationProperties; - import javax.annotation.PostConstruct; import javax.validation.constraints.DecimalMax; import javax.validation.constraints.DecimalMin; -import javax.validation.constraints.Max; import javax.validation.constraints.Min; +import javax.validation.constraints.NotNull; + +import lombok.Data; + +import org.joda.time.Period; +import org.springframework.boot.context.properties.ConfigurationProperties; @ConfigurationProperties(prefix = "consul.heartbeat") @Data -public class TtlHeartbeatProperties { +public class HeartbeatProperties { @Min(1) - @Max(10) - private volatile int ttl = 3; + private int ttlValue = 30; + + @NotNull + private String ttlUnit = "s"; @DecimalMin("0.1") @DecimalMax("0.9") - private volatile double intervalRatio = 2.0 / 3.0; + private double intervalRatio = 2.0 / 3.0; - private volatile Period heartbeatInterval; + private Period heartbeatInterval; @PostConstruct public void computeHeartbeatInterval() { // heartbeat rate at ratio * ttl, but no later than ttl -1s and, (under lesser // priority), no sooner than 1s from now - heartbeatInterval = new Period(Math.round(1000 * Math.max(ttl - 1, - Math.min(ttl * intervalRatio, 1)))); + heartbeatInterval = new Period(Math.round(1000 * Math.max(ttlValue - 1, + Math.min(ttlValue * intervalRatio, 1)))); } - public String getTtlAsString() { - return ttl + "s"; + public String getTtl() { + return ttlValue + ttlUnit; } } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java index e5acb87c..5563b9de 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/TtlScheduler.java @@ -4,7 +4,6 @@ import com.ecwid.consul.v1.ConsulClient; import com.ecwid.consul.v1.agent.model.NewService; import lombok.extern.slf4j.Slf4j; import org.joda.time.DateTime; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.scheduling.annotation.Scheduled; import java.util.Map; @@ -19,14 +18,16 @@ public class TtlScheduler { public static final DateTime EXPIRED_DATE = new DateTime(0); private final Map serviceHeartbeats = new ConcurrentHashMap<>(); - private final AtomicBoolean heartbeatingNow = new AtomicBoolean(); - @Autowired - private TtlHeartbeatProperties configuration; + private HeartbeatProperties configuration; - @Autowired private ConsulClient client; + public TtlScheduler(HeartbeatProperties configuration, ConsulClient client) { + this.configuration = configuration; + this.client = client; + } + /** * Add a service to the checks loop. */ @@ -38,16 +39,19 @@ public class TtlScheduler { serviceHeartbeats.remove(serviceId); } - @Scheduled(initialDelay = 0, fixedRateString = "${consul.heartbeat.fixedRate:201000}") + @Scheduled(initialDelay = 0, fixedRateString = "${consul.heartbeat.fixedRate:15000}") private void heartbeatServices() { - if (heartbeatingNow.compareAndSet(false, true)) { - for (String serviceId : serviceHeartbeats.keySet()) { - DateTime latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); - if (latestHeartbeatDoneForService.plus(configuration.getHeartbeatInterval()) - .isBefore(DateTime.now())) { - client.agentCheckPass(serviceId); - serviceHeartbeats.put(serviceId, DateTime.now()); + for (String serviceId : serviceHeartbeats.keySet()) { + DateTime latestHeartbeatDoneForService = serviceHeartbeats.get(serviceId); + if (latestHeartbeatDoneForService.plus(configuration.getHeartbeatInterval()) + .isBefore(DateTime.now())) { + String checkId = serviceId; + if (!checkId.startsWith("service:")) { + checkId = "service:"+checkId; } + client.agentCheckPass(checkId); + log.info("Sending consul heartbeat for: "+serviceId); + serviceHeartbeats.put(serviceId, DateTime.now()); } } } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java index 423c0d27..47e6c91a 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/AliveServerListFilter.java @@ -4,6 +4,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Set; +import com.netflix.loadbalancer.Server; import org.springframework.cloud.consul.discovery.ConsulServer; import com.netflix.loadbalancer.ServerListFilter; @@ -16,7 +17,7 @@ import com.netflix.loadbalancer.ServerListFilter; * liv in both). * @author nicu marasoiu on 10.03.2015. */ -public class AliveServerListFilter implements ServerListFilter { +public class AliveServerListFilter implements ServerListFilter { private FilteringAgentClient filteringAgentClient; public AliveServerListFilter(FilteringAgentClient filteringAgentClient) { @@ -24,11 +25,12 @@ public class AliveServerListFilter implements ServerListFilter { } @Override - public List getFilteredListOfServers(List servers) { + public List getFilteredListOfServers(List servers) { Set liveNodes = filteringAgentClient.getAliveAgentsAddresses(); - List filteredServers = new ArrayList<>(); - for (ConsulServer server : servers) { - if (liveNodes.contains(server.getAddress())) { + List filteredServers = new ArrayList<>(); + for (Server server : servers) { + ConsulServer consulServer = ConsulServer.class.cast(server); + if (liveNodes.contains(consulServer.getAddress())) { filteredServers.add(server); } } diff --git a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java index 6d2c5e5f..211e63b7 100644 --- a/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java +++ b/spring-cloud-consul-discovery/src/main/java/org/springframework/cloud/consul/discovery/filters/ServiceCheckServerListFilter.java @@ -5,6 +5,7 @@ import java.util.HashSet; import java.util.List; import java.util.Set; +import com.netflix.loadbalancer.Server; import org.springframework.cloud.consul.discovery.ConsulServer; import com.ecwid.consul.v1.ConsulClient; @@ -15,7 +16,7 @@ import com.netflix.loadbalancer.ServerListFilter; /** * Created by nicu on 12.03.2015. */ -public class ServiceCheckServerListFilter implements ServerListFilter { +public class ServiceCheckServerListFilter implements ServerListFilter { private ConsulClient client; @@ -24,12 +25,13 @@ public class ServiceCheckServerListFilter implements ServerListFilter getFilteredListOfServers(List servers) { + public List getFilteredListOfServers(List servers) { Set passingServiceIds = getPassingServiceIds(servers); - List okServers = new ArrayList<>(servers.size()); - for (ConsulServer consulServer : servers) { - String serviceId = consulServer.getMetaInfo().getInstanceId(); + List okServers = new ArrayList<>(servers.size()); + for (Server server : servers) { + String serviceId = server.getMetaInfo().getInstanceId(); if (passingServiceIds.contains(serviceId)) { + ConsulServer consulServer = ConsulServer.class.cast(server); List nodeChecks = client.getHealthChecksForNode( consulServer.getNode(), QueryParams.DEFAULT).getValue(); boolean passingNodeChecks = true; @@ -40,16 +42,16 @@ public class ServiceCheckServerListFilter implements ServerListFilter getPassingServiceIds(List servers) { + private Set getPassingServiceIds(List servers) { Set serviceIds = new HashSet<>(1); - for (ConsulServer server : servers) { + for (Server server : servers) { serviceIds.add(server.getMetaInfo().getInstanceId()); } for (String serviceId : serviceIds) {