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..43a1870f --- /dev/null +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java @@ -0,0 +1,68 @@ +/* + * 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 org.slf4j.Logger; +import org.slf4j.LoggerFactory; +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 + */ +public class KubernetesCatalogWatch implements ApplicationEventPublisherAware { + + private static final Logger logger = LoggerFactory.getLogger(KubernetesCatalogWatch.class); + + 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)) { + logger.trace("Received services update from kubernetesDiscoveryClient: {}", services); + publisher.publishEvent(new HeartbeatEvent(this, services)); + } + } catch (Exception e) { + logger.error("Error watching Kubernetes Services", e); + } + } +} diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientAutoConfiguration.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientAutoConfiguration.java index 036583af..7295daeb 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientAutoConfiguration.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientAutoConfiguration.java @@ -1,6 +1,8 @@ package org.springframework.cloud.kubernetes.discovery; import io.fabric8.kubernetes.client.KubernetesClient; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.cloud.client.discovery.DiscoveryClient; import org.springframework.cloud.kubernetes.registry.KubernetesRegistration; import org.springframework.cloud.kubernetes.registry.KubernetesServiceRegistry; @@ -33,4 +35,11 @@ public class KubernetesDiscoveryClientAutoConfiguration { public KubernetesDiscoveryProperties getKubernetesDiscoveryProperties() { return new KubernetesDiscoveryProperties(); } + + @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/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchConfigurationTest.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchConfigurationTest.java new file mode 100644 index 00000000..29d3f211 --- /dev/null +++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchConfigurationTest.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 KubernetesCatalogServicesWatchConfigurationTest { + + 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, + KubernetesDiscoveryClientAutoConfiguration.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)); + } +}