From 0d33d26b6abb1cbb0dda421feed991f7f583615f Mon Sep 17 00:00:00 2001 From: Gytis Trikleris Date: Thu, 24 May 2018 19:26:38 +0200 Subject: [PATCH] KubernetesLock#tryLock --- .../cloud/kubernetes/lock/KubernetesLock.java | 34 +++++----- .../kubernetes/lock/KubernetesLockTest.java | 68 ++++++++++++++++--- 2 files changed, 74 insertions(+), 28 deletions(-) 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 0b10396b..d6b78759 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 @@ -41,50 +41,50 @@ public class KubernetesLock implements Lock { @Override public void lock() { - repository.deleteIfExpired(name); while (true) { try { - if (repository.create(name, holder, expiration)) { + if (tryLock()) { return; } Thread.sleep(LOCK_RETRY_INTERVAL); } catch (InterruptedException e) { // This method cannot be interrupted } - } + } } @Override public void lockInterruptibly() throws InterruptedException { - repository.deleteIfExpired(name); while (true) { - try { - if (repository.create(name, holder, expiration)) { - return; - } - Thread.sleep(LOCK_RETRY_INTERVAL); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw e; + if (tryLock()) { + return; } + Thread.sleep(LOCK_RETRY_INTERVAL); } } @Override public boolean tryLock() { - // TODO same as lock() but return immediately if te lock cannot be acquired - return false; + repository.deleteIfExpired(name); + return repository.create(name, holder, expiration); } @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; + long expiration = System.currentTimeMillis() + TimeUnit.MILLISECONDS.convert(time, unit); + while (true) { + if (System.currentTimeMillis() > expiration) { + return false; + } else if (tryLock()) { + return true; + } + Thread.sleep(LOCK_RETRY_INTERVAL); + } } @Override public void unlock() { - // TODO delete configmap for this lock + // TODO delete configmap for this lock // TODO consider only allowing lock release only from the same application and/or thread } 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 8dd7ee23..7cb623ee 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 @@ -1,5 +1,8 @@ package org.springframework.cloud.kubernetes.lock; +import java.util.concurrent.TimeUnit; + +import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; @@ -7,6 +10,7 @@ import org.mockito.invocation.InvocationOnMock; import org.mockito.junit.MockitoJUnitRunner; import org.mockito.stubbing.Answer; +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; @@ -21,11 +25,17 @@ public class KubernetesLockTest { @Mock private ConfigMapLockRepository repository; + private KubernetesLock lock; + + @Before + public void before() { + lock = new KubernetesLock(repository, NAME, HOLDER, 0); + } + @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); @@ -36,10 +46,9 @@ public class KubernetesLockTest { public void lockShouldWaitUntilLockCanBeAcquired() { given(repository.create(NAME, HOLDER, 0)).will(new LockInSecondCall()); - KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); lock.lock(); - verify(repository).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfExpired(NAME); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -47,10 +56,9 @@ public class KubernetesLockTest { public void lockShouldNotBeInterrupted() { given(repository.create(NAME, HOLDER, 0)).will(new InterrupFirstCall()); - KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); lock.lock(); - verify(repository).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfExpired(NAME); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -58,7 +66,6 @@ public class KubernetesLockTest { public void shouldLockWithLockInterruptibly() throws InterruptedException { given(repository.create(NAME, HOLDER, 0)).willReturn(true); - KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); lock.lockInterruptibly(); verify(repository).deleteIfExpired(NAME); @@ -69,10 +76,9 @@ public class KubernetesLockTest { public void lockInterruptiblyShouldWaitUntilLockCanBeAcquired() throws InterruptedException { given(repository.create(NAME, HOLDER, 0)).will(new LockInSecondCall()); - KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); lock.lockInterruptibly(); - verify(repository).deleteIfExpired(NAME); + verify(repository, times(2)).deleteIfExpired(NAME); verify(repository, times(2)).create(NAME, HOLDER, 0); } @@ -80,23 +86,64 @@ public class KubernetesLockTest { public void lockInterruptiblyShouldBeInterrupted() throws InterruptedException { given(repository.create(NAME, HOLDER, 0)).will(new InterrupFirstCall()); - KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); lock.lockInterruptibly(); } @Test public void shouldLockWithTryLock() { + given(repository.create(NAME, HOLDER, 0)).willReturn(true); + assertThat(lock.tryLock()).isTrue(); + + verify(repository).deleteIfExpired(NAME); + verify(repository).create(NAME, HOLDER, 0); } @Test public void shouldFailToLockWithTryLock() { + given(repository.create(NAME, HOLDER, 0)).willReturn(false); + assertThat(lock.tryLock()).isFalse(); + + verify(repository).deleteIfExpired(NAME); + verify(repository).create(NAME, HOLDER, 0); } @Test - public void shouldFailWithTryLockTimeout() { + public void shouldLockWithTryLockTimeout() throws InterruptedException { + given(repository.create(NAME, HOLDER, 0)).willReturn(true); + assertThat(lock.tryLock(2, TimeUnit.SECONDS)).isTrue(); + + verify(repository).deleteIfExpired(NAME); + verify(repository).create(NAME, HOLDER, 0); + } + + @Test + public void tryLockShouldWaitUntilLockIsAvailable() throws InterruptedException { + given(repository.create(NAME, HOLDER, 0)).will(new LockInSecondCall()); + + assertThat(lock.tryLock(2, TimeUnit.SECONDS)).isTrue(); + + verify(repository, times(2)).deleteIfExpired(NAME); + verify(repository, times(2)).create(NAME, HOLDER, 0); + } + + @Test + public void tryLockShouldTimeout() throws InterruptedException { + given(repository.create(NAME, HOLDER, 0)).will(new LockInSecondCall()); + + assertThat(lock.tryLock(90, TimeUnit.MILLISECONDS)).isFalse(); + + verify(repository).deleteIfExpired(NAME); + verify(repository).create(NAME, HOLDER, 0); + } + + @Test(expected = InterruptedException.class) + public void tryLockShouldBeInterrupted() throws InterruptedException { + given(repository.create(NAME, HOLDER, 0)).will(new InterrupFirstCall()); + + lock.tryLock(90, TimeUnit.MILLISECONDS); } @Test @@ -106,7 +153,6 @@ public class KubernetesLockTest { @Test(expected = UnsupportedOperationException.class) public void newConditionShouldFail() { - KubernetesLock lock = new KubernetesLock(repository, NAME, HOLDER, 0); lock.newCondition(); }