diff --git a/spring-cloud-kubernetes-discovery/pom.xml b/spring-cloud-kubernetes-discovery/pom.xml
index c2f2dd12..aaea904a 100644
--- a/spring-cloud-kubernetes-discovery/pom.xml
+++ b/spring-cloud-kubernetes-discovery/pom.xml
@@ -52,6 +52,13 @@
true
+
+ org.projectlombok
+ lombok
+
+ provided
+
+
org.springframework.boot
diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java
new file mode 100644
index 00000000..0e2e6f40
--- /dev/null
+++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java
@@ -0,0 +1,66 @@
+/*
+ * Copyright (C) 2016 to the original authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.cloud.kubernetes.discovery;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.cloud.client.discovery.event.HeartbeatEvent;
+import org.springframework.context.ApplicationEventPublisher;
+import org.springframework.context.ApplicationEventPublisherAware;
+import org.springframework.scheduling.annotation.Scheduled;
+
+import java.util.List;
+import java.util.concurrent.atomic.AtomicReference;
+
+
+/**
+ * @author Oleg Vyukov
+ */
+@Slf4j
+public class KubernetesCatalogWatch implements ApplicationEventPublisherAware {
+
+ private final KubernetesDiscoveryClient kubernetesDiscoveryClient;
+ private final AtomicReference> catalogServicesState = new AtomicReference<>();
+ private ApplicationEventPublisher publisher;
+
+ public KubernetesCatalogWatch(KubernetesDiscoveryClient kubernetesDiscoveryClient) {
+ this.kubernetesDiscoveryClient = kubernetesDiscoveryClient;
+ }
+
+ @Override
+ public void setApplicationEventPublisher(ApplicationEventPublisher publisher) {
+ this.publisher = publisher;
+ }
+
+ @Scheduled(fixedDelayString = "${spring.cloud.kubernetes.discovery.catalogServicesWatchDelay:30000}")
+ public void catalogServicesWatch() {
+ try {
+ List previousState = catalogServicesState.get();
+
+ List services = kubernetesDiscoveryClient.getServices();
+
+ services.sort(String::compareTo);
+ catalogServicesState.set(services);
+
+ if (!services.equals(previousState)) {
+ log.trace("Received services update from kubernetesDiscoveryClient: {}", services);
+ publisher.publishEvent(new HeartbeatEvent(this, services));
+ }
+ } catch (Exception e) {
+ log.error("Error watching Kubernetes Services", e);
+ }
+ }
+}
diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientConfiguration.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientConfiguration.java
index 824edf61..925532a5 100644
--- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientConfiguration.java
+++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientConfiguration.java
@@ -37,6 +37,12 @@ public class KubernetesDiscoveryClientConfiguration {
@Bean
public KubernetesDiscoveryLifecycle kubernetesDiscoveryLifecycle(KubernetesClient client, KubernetesDiscoveryProperties properties) {
return new KubernetesDiscoveryLifecycle(client, properties);
+ }
- }
+ @Bean
+ @ConditionalOnMissingBean
+ @ConditionalOnProperty(name = "spring.cloud.kubernetes.discovery.catalog-services-watch.enabled", matchIfMissing = true)
+ public KubernetesCatalogWatch kubernetesCatalogWatch(KubernetesDiscoveryClient client) {
+ return new KubernetesCatalogWatch(client);
+ }
}
diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java
index 26b5a21f..35ef9378 100644
--- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java
+++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java
@@ -27,6 +27,8 @@ public class KubernetesDiscoveryProperties {
@Value("${spring.application.name:unknown}")
private String serviceName = "unknown";
+ private int catalogServicesWatchDelay = 30000;
+
public boolean isEnabled() {
return enabled;
}
@@ -38,4 +40,12 @@ public class KubernetesDiscoveryProperties {
public String getServiceName() {
return serviceName;
}
+
+ public int getCatalogServicesWatchDelay() {
+ return catalogServicesWatchDelay;
+ }
+
+ public void setCatalogServicesWatchDelay(int catalogServicesWatchDelay) {
+ this.catalogServicesWatchDelay = catalogServicesWatchDelay;
+ }
}
diff --git a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchtConfigurationTest.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchtConfigurationTest.java
new file mode 100644
index 00000000..c5a49942
--- /dev/null
+++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchtConfigurationTest.java
@@ -0,0 +1,69 @@
+package org.springframework.cloud.kubernetes.discovery;
+
+import io.fabric8.kubernetes.client.KubernetesClient;
+import org.junit.After;
+import org.junit.Ignore;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.autoconfigure.context.PropertyPlaceholderAutoConfiguration;
+import org.springframework.boot.builder.SpringApplicationBuilder;
+import org.springframework.boot.test.context.ConfigFileApplicationContextInitializer;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.boot.test.context.TestConfiguration;
+import org.springframework.boot.test.mock.mockito.MockBean;
+import org.springframework.context.ApplicationContext;
+import org.springframework.context.ConfigurableApplicationContext;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.test.context.junit4.SpringRunner;
+
+import static org.junit.Assert.*;
+import static org.mockito.Mockito.mock;
+
+/**
+ * @author Oleg Vyukov
+ */
+public class KubernetesCatalogServicesWatchtConfigurationTest {
+
+ private ConfigurableApplicationContext context;
+
+ @After
+ public void close() {
+ if (this.context != null) {
+ this.context.close();
+ }
+ }
+
+ @Test
+ public void kubernetesCatalogWatchDisabled() throws Exception {
+ setup("spring.cloud.kubernetes.discovery.catalog-services-watch.enabled=false");
+ assertFalse(context.containsBean("kubernetesCatalogWatch"));
+ }
+
+ @Test
+ public void kubernetesCatalogWatchDefaultEnabled() throws Exception {
+ setup();
+ assertTrue(context.containsBean("kubernetesCatalogWatch"));
+ }
+
+ private void setup(String... env) {
+ this.context = new SpringApplicationBuilder(
+ PropertyPlaceholderAutoConfiguration.class,
+ KubernetesClientTestConfiguration.class,
+ KubernetesDiscoveryClientConfiguration.class).web(false)
+ .properties(env).run();
+ }
+
+
+ @Configuration
+ static class KubernetesClientTestConfiguration {
+
+ @Bean
+ KubernetesClient kubernetesClient() {
+ return mock(KubernetesClient.class);
+ }
+
+
+ }
+}
diff --git a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java
new file mode 100644
index 00000000..71580e2e
--- /dev/null
+++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java
@@ -0,0 +1,51 @@
+package org.springframework.cloud.kubernetes.discovery;
+
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.runners.MockitoJUnitRunner;
+import org.springframework.cloud.client.discovery.event.HeartbeatEvent;
+import org.springframework.context.ApplicationEventPublisher;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.Assert.*;
+import static org.mockito.Mockito.*;
+
+/**
+ * @author Oleg Vyukov
+ */
+@RunWith(MockitoJUnitRunner.class)
+public class KubernetesCatalogWatchTest {
+
+ @Mock
+ private KubernetesDiscoveryClient kubernetesDiscoveryClient;
+
+ @Mock
+ private ApplicationEventPublisher applicationEventPublisher;
+
+ @InjectMocks
+ private KubernetesCatalogWatch underTest;
+
+ @Before
+ public void setUp() throws Exception {
+ underTest.setApplicationEventPublisher(applicationEventPublisher);
+ }
+
+ @Test
+ public void testRandomOrder() throws Exception {
+ final List services = Arrays.asList("api", "api", "other");
+ final List shuffleServices = Arrays.asList("api", "other", "api");
+ when(kubernetesDiscoveryClient.getServices()).thenReturn(services);
+
+ underTest.catalogServicesWatch();
+ // second execution on shuffleServices
+ underTest.catalogServicesWatch();
+
+ verify(applicationEventPublisher).publishEvent(any(HeartbeatEvent.class));
+ }
+}