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; + } + +}