From 34e8e7d0f09b94c1990f62cda8bc43056e989825 Mon Sep 17 00:00:00 2001 From: Oleg Vyukov Date: Sat, 24 Feb 2018 16:02:24 +0300 Subject: [PATCH] KubernetesCatalogWatch --- spring-cloud-kubernetes-discovery/pom.xml | 7 ++ .../discovery/KubernetesCatalogWatch.java | 66 ++++++++++++++++++ ...ubernetesDiscoveryClientConfiguration.java | 8 ++- .../KubernetesDiscoveryProperties.java | 10 +++ ...atalogServicesWatchtConfigurationTest.java | 69 +++++++++++++++++++ .../discovery/KubernetesCatalogWatchTest.java | 51 ++++++++++++++ 6 files changed, 210 insertions(+), 1 deletion(-) create mode 100644 spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java create mode 100644 spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogServicesWatchtConfigurationTest.java create mode 100644 spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java 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)); + } +}