From 085d2cd5c336a1d89fae6e00d91d6c4ed342159a Mon Sep 17 00:00:00 2001 From: Gytis Trikleris Date: Fri, 25 May 2018 13:22:18 +0200 Subject: [PATCH] KubernetesLockRegistry --- .../lock/ConfigMapLockRepository.java | 18 +++---- .../cloud/kubernetes/lock/KubernetesLock.java | 12 +++-- .../lock/KubernetesLockRegistry.java | 53 +++++++++++++++++++ .../lock/ConfigMapLockRepositoryIT.java | 15 +++--- .../lock/ConfigMapLockRepositoryTest.java | 14 ++--- .../lock/KubernetesLockRegistryTest.java | 50 +++++++++++++++++ .../kubernetes/lock/KubernetesLockTest.java | 21 ++++---- 7 files changed, 145 insertions(+), 38 deletions(-) create mode 100644 spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistry.java create mode 100644 spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistryTest.java diff --git a/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepository.java b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepository.java index 37f17377..731bda59 100644 --- a/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepository.java +++ b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepository.java @@ -31,7 +31,7 @@ public class ConfigMapLockRepository { static final String HOLDER_KEY = "holder"; - static final String EXPIRATION_KEY = "expiration"; + static final String CREATED_AT_KEY = "created_at"; static final String PROVIDER_LABEL = "provider"; @@ -62,16 +62,16 @@ public class ConfigMapLockRepository { return Optional.ofNullable(configMap); } - public boolean create(String name, String holder, long expiration) { + public boolean create(String name, String holder, long createdAt) { String configMapName = getConfigMapName(name); - String expirationString = String.valueOf(expiration); + String createdAtString = String.valueOf(createdAt); ConfigMap configMap = new ConfigMapBuilder().withNewMetadata() .withName(configMapName) .addToLabels(PROVIDER_LABEL, PROVIDER_LABEL_VALUE) .addToLabels(KIND_LABEL, KIND_LABEL_VALUE) .endMetadata() .addToData(HOLDER_KEY, holder) - .addToData(EXPIRATION_KEY, expirationString) + .addToData(CREATED_AT_KEY, createdAtString) .build(); try { @@ -94,16 +94,16 @@ public class ConfigMapLockRepository { .delete(); } - public void deleteIfExpired(String name) { + public void deleteIfOlderThan(String name, long age) { get(name) - .filter(this::isExpired) + .filter(m -> this.isOlder(m, age)) // TODO what if someone else deletes and creates a lock in this gap? .ifPresent(c -> delete(name)); } - private boolean isExpired(ConfigMap configMap) { - String expirationString = configMap.getData().get(EXPIRATION_KEY); - return Long.valueOf(expirationString) < System.currentTimeMillis(); + private boolean isOlder(ConfigMap configMap, long age) { + String createdAtString = configMap.getData().get(CREATED_AT_KEY); + return System.currentTimeMillis() - Long.valueOf(createdAtString) > age; } private String getConfigMapName(String name) { diff --git a/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLock.java b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLock.java index 0119f247..f13386c0 100644 --- a/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLock.java +++ b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLock.java @@ -22,6 +22,8 @@ import java.util.concurrent.locks.Lock; public class KubernetesLock implements Lock { + static final long DEFAULT_TTL = 10000; + private static final int LOCK_RETRY_INTERVAL = 100; private final ConfigMapLockRepository repository; @@ -30,13 +32,13 @@ public class KubernetesLock implements Lock { private final String holder; - private final long expiration; + private final long createdAt; - public KubernetesLock(ConfigMapLockRepository repository, String name, String holder, long expiration) { + public KubernetesLock(ConfigMapLockRepository repository, String name, String holder, long createdAt) { this.repository = repository; this.name = name; this.holder = holder; - this.expiration = expiration; + this.createdAt = createdAt; } @Override @@ -65,8 +67,8 @@ public class KubernetesLock implements Lock { @Override public boolean tryLock() { - repository.deleteIfExpired(name); - return repository.create(name, holder, expiration); + repository.deleteIfOlderThan(name, DEFAULT_TTL); + return repository.create(name, holder, createdAt); } @Override diff --git a/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistry.java b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistry.java new file mode 100644 index 00000000..01eafa13 --- /dev/null +++ b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistry.java @@ -0,0 +1,53 @@ +/* + * Copyright (C) 2018 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.lock; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.locks.Lock; + +import org.springframework.integration.support.locks.ExpirableLockRegistry; +import org.springframework.util.Assert; + +public class KubernetesLockRegistry implements ExpirableLockRegistry { + + private final ConfigMapLockRepository repository; + + private final String id; + + private final Map locks; + + public KubernetesLockRegistry(ConfigMapLockRepository repository, String id) { + this.repository = repository; + this.id = id; + this.locks = new HashMap<>(); + } + + @Override + public Lock obtain(Object key) { + Assert.isInstanceOf(String.class, key); + String name = (String) key; + + return locks.computeIfAbsent(name, n -> new KubernetesLock(repository, n, id, System.currentTimeMillis())); + } + + @Override + public void expireUnusedOlderThan(long age) { + locks.forEach((n, l) -> repository.deleteIfOlderThan(n, age)); + } + +} diff --git a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryIT.java b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryIT.java index 27104203..e973e865 100644 --- a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryIT.java +++ b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryIT.java @@ -15,8 +15,9 @@ import org.junit.Test; import org.junit.runner.RunWith; import static org.assertj.core.api.Assertions.assertThat; -import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.EXPIRATION_KEY; +import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.CREATED_AT_KEY; import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.HOLDER_KEY; +import static org.springframework.cloud.kubernetes.lock.KubernetesLock.DEFAULT_TTL; @RunWith(ArquillianConditionalRunner.class) @RequiresKubernetes @@ -53,7 +54,7 @@ public class ConfigMapLockRepositoryIT { Map data = optionalConfigMap.get().getData(); assertThat(data).containsEntry(HOLDER_KEY, HOLDER); - assertThat(data).containsEntry(EXPIRATION_KEY, String.valueOf(1000)); + assertThat(data).containsEntry(CREATED_AT_KEY, String.valueOf(1000)); } @Test @@ -64,22 +65,22 @@ public class ConfigMapLockRepositoryIT { @Test public void shouldDelete() { - repository.create(NAME, HOLDER, System.currentTimeMillis() + 10000); + repository.create(NAME, HOLDER, 0); repository.delete(NAME); assertThat(repository.get(NAME).isPresent()).isFalse(); } @Test public void shouldDeleteExpired() { - repository.create(NAME, HOLDER, System.currentTimeMillis() - 1); - repository.deleteIfExpired(NAME); + repository.create(NAME, HOLDER, 0); + repository.deleteIfOlderThan(NAME, DEFAULT_TTL); assertThat(repository.get(NAME).isPresent()).isFalse(); } @Test public void shouldKeepNotExpired() { - repository.create(NAME, HOLDER, System.currentTimeMillis() + 10000); - repository.deleteIfExpired(NAME); + repository.create(NAME, HOLDER, System.currentTimeMillis()); + repository.deleteIfOlderThan(NAME, DEFAULT_TTL); assertThat(repository.get(NAME).isPresent()).isTrue(); } diff --git a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryTest.java b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryTest.java index 46e4db08..39a8e7d8 100644 --- a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryTest.java +++ b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/ConfigMapLockRepositoryTest.java @@ -24,7 +24,7 @@ import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.CONFIG_MAP_PREFIX; -import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.EXPIRATION_KEY; +import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.CREATED_AT_KEY; import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.HOLDER_KEY; import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.KIND_LABEL; import static org.springframework.cloud.kubernetes.lock.ConfigMapLockRepository.KIND_LABEL_VALUE; @@ -106,7 +106,7 @@ public class ConfigMapLockRepositoryTest { .addToLabels(KIND_LABEL, KIND_LABEL_VALUE) .endMetadata() .addToData(HOLDER_KEY, HOLDER) - .addToData(EXPIRATION_KEY, "1000") + .addToData(CREATED_AT_KEY, "1000") .build(); verify(mockInNamespaceOperation).create(eq(expectedConfigMap)); } @@ -127,19 +127,19 @@ public class ConfigMapLockRepositoryTest { } @Test - public void shouldDeleteExpired() { - given(mockData.get(EXPIRATION_KEY)).willReturn(String.valueOf(System.currentTimeMillis() - 1)); + public void shouldDeleteOld() { + given(mockData.get(CREATED_AT_KEY)).willReturn(String.valueOf(System.currentTimeMillis() - 1000)); - repository.deleteIfExpired(NAME); + repository.deleteIfOlderThan(NAME, 100); verify(mockWithNameResource).delete(); } @Test public void shouldNotDeleteNonExpired() { - given(mockData.get(EXPIRATION_KEY)).willReturn(String.valueOf(System.currentTimeMillis() + 10000)); + given(mockData.get(CREATED_AT_KEY)).willReturn(String.valueOf(System.currentTimeMillis()) + 1000); - repository.deleteIfExpired(NAME); + repository.deleteIfOlderThan(NAME, 100); verify(mockWithNameResource, times(0)).delete(); } diff --git a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistryTest.java b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistryTest.java new file mode 100644 index 00000000..153bbdf0 --- /dev/null +++ b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockRegistryTest.java @@ -0,0 +1,50 @@ +package org.springframework.cloud.kubernetes.lock; + +import java.util.concurrent.locks.Lock; + +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.junit.MockitoJUnitRunner; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.verify; + +@RunWith(MockitoJUnitRunner.class) +public class KubernetesLockRegistryTest { + + @Mock + private ConfigMapLockRepository mockRepository; + + private KubernetesLockRegistry registry; + + @Before + public void before() { + registry = new KubernetesLockRegistry(mockRepository, "test-id"); + } + + @Test + public void shouldObtain() { + Lock actualLock = registry.obtain("test-name"); + KubernetesLock expectedLock = new KubernetesLock(mockRepository, "test-name", "test-id", System.currentTimeMillis()); + + assertThat(actualLock).isEqualToComparingOnlyGivenFields(expectedLock, "name", "holder"); + } + + @Test(expected = IllegalArgumentException.class) + public void shouldFailToObtainWithNonStringKey() { + registry.obtain(new Object()); + } + + @Test + public void expireUnusedOlderThanShouldDelegate() { + registry.obtain("test-name-1"); + registry.obtain("test-name-2"); + registry.expireUnusedOlderThan(100); + + verify(mockRepository).deleteIfOlderThan("test-name-1", 100); + verify(mockRepository).deleteIfOlderThan("test-name-2", 100); + } + +} diff --git a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockTest.java b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockTest.java index 2e6f6da0..9f4eedd4 100644 --- a/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockTest.java +++ b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockTest.java @@ -14,6 +14,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.springframework.cloud.kubernetes.lock.KubernetesLock.DEFAULT_TTL; @RunWith(MockitoJUnitRunner.class) public class KubernetesLockTest { @@ -38,7 +39,7 @@ public class KubernetesLockTest { lock.lock(); - verify(repository).deleteIfExpired(NAME); + verify(repository).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository).create(NAME, HOLDER, 0); } @@ -48,7 +49,7 @@ public class KubernetesLockTest { lock.lock(); - verify(repository, times(2)).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -58,7 +59,7 @@ public class KubernetesLockTest { lock.lock(); - verify(repository, times(2)).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -68,7 +69,7 @@ public class KubernetesLockTest { lock.lockInterruptibly(); - verify(repository).deleteIfExpired(NAME); + verify(repository).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository).create(NAME, HOLDER, 0); } @@ -78,7 +79,7 @@ public class KubernetesLockTest { lock.lockInterruptibly(); - verify(repository, times(2)).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -95,7 +96,7 @@ public class KubernetesLockTest { assertThat(lock.tryLock()).isTrue(); - verify(repository).deleteIfExpired(NAME); + verify(repository).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository).create(NAME, HOLDER, 0); } @@ -105,7 +106,7 @@ public class KubernetesLockTest { assertThat(lock.tryLock()).isFalse(); - verify(repository).deleteIfExpired(NAME); + verify(repository).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository).create(NAME, HOLDER, 0); } @@ -115,7 +116,7 @@ public class KubernetesLockTest { assertThat(lock.tryLock(2, TimeUnit.SECONDS)).isTrue(); - verify(repository).deleteIfExpired(NAME); + verify(repository).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository).create(NAME, HOLDER, 0); } @@ -125,7 +126,7 @@ public class KubernetesLockTest { assertThat(lock.tryLock(2, TimeUnit.SECONDS)).isTrue(); - verify(repository, times(2)).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -135,7 +136,7 @@ public class KubernetesLockTest { assertThat(lock.tryLock(90, TimeUnit.MILLISECONDS)).isFalse(); - verify(repository).deleteIfExpired(NAME); + verify(repository).deleteIfOlderThan(NAME, DEFAULT_TTL); verify(repository).create(NAME, HOLDER, 0); }