From 52e82796ba24303c4b305b2dd4673bdf54fa1f47 Mon Sep 17 00:00:00 2001 From: Gytis Trikleris Date: Thu, 24 May 2018 18:45:19 +0200 Subject: [PATCH] KubernetesLock#lock() --- spring-cloud-kubernetes-lock/pom.xml | 5 + .../lock/ConfigMapLockRepository.java | 1 + .../cloud/kubernetes/lock/KubernetesLock.java | 85 ++++++++++++++ .../kubernetes/lock/KubernetesLockTest.java | 110 ++++++++++++++++++ 4 files changed, 201 insertions(+) create mode 100644 spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLock.java create mode 100644 spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockTest.java diff --git a/spring-cloud-kubernetes-lock/pom.xml b/spring-cloud-kubernetes-lock/pom.xml index 29e30e81..a1a9c411 100644 --- a/spring-cloud-kubernetes-lock/pom.xml +++ b/spring-cloud-kubernetes-lock/pom.xml @@ -61,6 +61,11 @@ assertj-core test + + org.mockito + mockito-core + test + 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 426783dc..73533ef4 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 @@ -97,6 +97,7 @@ public class ConfigMapLockRepository { public void deleteIfExpired(String name) { get(name) .filter(this::isExpired) + // TODO what if someone else deletes and creates a lock in this gap? .ifPresent(c -> delete(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 new file mode 100644 index 00000000..cadbae5f --- /dev/null +++ b/spring-cloud-kubernetes-lock/src/main/java/org/springframework/cloud/kubernetes/lock/KubernetesLock.java @@ -0,0 +1,85 @@ +/* + * 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.concurrent.TimeUnit; +import java.util.concurrent.locks.Condition; +import java.util.concurrent.locks.Lock; + +public class KubernetesLock implements Lock { + + private static final int RETRY_PERIOD = 100; + + private final ConfigMapLockRepository repository; + + private final String name; + + private final String holder; + + private final long expiration; + + public KubernetesLock(ConfigMapLockRepository repository, String name, String holder, long expiration) { + this.repository = repository; + this.name = name; + this.holder = holder; + this.expiration = expiration; + } + + @Override + public void lock() { + repository.deleteIfExpired(name); + while (true) { + try { + if (repository.create(name, holder, expiration)) { + return; + } + Thread.sleep(RETRY_PERIOD); + } catch (InterruptedException e) { + // This method cannot be interrupted + } + } + } + + @Override + public void lockInterruptibly() throws InterruptedException { + // TODO same as lock(), but thread can be interrupted + } + + @Override + public boolean tryLock() { + // TODO same as lock() but return immediately if te lock cannot be acquired + return false; + } + + @Override + public boolean tryLock(long time, TimeUnit unit) throws InterruptedException { + // TODO same as lockInterruptibly() but return after specified time if te lock cannot be acquired. + return false; + } + + @Override + public void unlock() { + // TODO delete configmap for this lock + // TODO consider only allowing lock release only from the same application and/or thread + } + + @Override + public Condition newCondition() { + throw new UnsupportedOperationException("Condition is not supported"); + } + +} 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 new file mode 100644 index 00000000..30c4e49e --- /dev/null +++ b/spring-cloud-kubernetes-lock/src/test/java/org/springframework/cloud/kubernetes/lock/KubernetesLockTest.java @@ -0,0 +1,110 @@ +package org.springframework.cloud.kubernetes.lock; + +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.junit.MockitoJUnitRunner; +import org.mockito.stubbing.Answer; + +import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +@RunWith(MockitoJUnitRunner.class) +public class KubernetesLockTest { + + private static final String NAME = "test-name"; + + private static final String HOLDER = "test-holder"; + + @Mock + private ConfigMapLockRepository repository; + + @Test + public void shouldLock() { + given(repository.create(NAME, HOLDER, 0)).willReturn(true); + + KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); + lock.lock(); + + verify(repository).deleteIfExpired(NAME); + verify(repository).create(NAME, HOLDER, 0); + } + + @Test + public void shouldWaitUntilLockCanBeAcquired() { + given(repository.create(NAME, HOLDER, 0)).will(new Answer() { + private int callsCounter = 0; + + @Override + public Boolean answer(InvocationOnMock invocationOnMock) { + return callsCounter++ > 0; + } + }); + + KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); + lock.lock(); + + verify(repository).deleteIfExpired(NAME); + verify(repository, times(2)).create(NAME, HOLDER, 0); + } + + @Test + public void lockShouldNotBeInterrupted() { + given(repository.create(NAME, HOLDER, 0)).will(new Answer() { + private int callsCounter = 0; + + @Override + public Boolean answer(InvocationOnMock invocationOnMock) throws InterruptedException { + if (callsCounter++ > 0) { + return true; + } + throw new InterruptedException("test exception"); + } + }); + + KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); + lock.lock(); + + verify(repository).deleteIfExpired(NAME); + verify(repository, times(2)).create(NAME, HOLDER, 0); + } + + @Test + public void shouldLockWithLockInterruptibly() { + + } + + @Test + public void lockInterruptiblyShouldBeInterrupted() { + + } + + @Test + public void shouldLockWithTryLock() { + + } + + @Test + public void shouldFailToLockWithTryLock() { + + } + + @Test + public void shouldFailWithTryLockTimeout() { + + } + + @Test + public void shouldUnlock() { + + } + + @Test(expected = UnsupportedOperationException.class) + public void newConditionShouldFail() { + KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); + lock.newCondition(); + } + +}