KubernetesLockRegistry
This commit is contained in:
committed by
Ioannis Canellos
parent
10771acc40
commit
085d2cd5c3
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<String, KubernetesLock> 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));
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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<String, String> 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();
|
||||
}
|
||||
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user