From fb7e4d68691adf869fc5884ec925a99646b77906 Mon Sep 17 00:00:00 2001 From: Vedran Pavic Date: Thu, 23 Mar 2017 00:18:33 +0100 Subject: [PATCH] INT-4246: Revoke leader when lock isn't acquired JIRA: https://jira.spring.io/browse/INT-4246 The `LockRegistryLeaderInitiator` currently does not revoke leadership when leading instance is unable to acquire lock from the underlying lock. This can result in multiple `LockRegistryLeaderInitiator` instances becoming leaders in the exceptional situations such as lock timeouts. * Handle leadership revoking when leading `LockRegistryLeaderInitiator` is unable to acquire lock Polish javadoc **Cherry-pick to master** Conflicts: spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java spring-integration-core/src/test/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiatorTests.java * Stop initiators in the end of test --- .../leader/LockRegistryLeaderInitiator.java | 57 +++++++-------- .../LockRegistryLeaderInitiatorTests.java | 72 ++++++++++++++++++- .../util/CallerBlocksPolicyTests.java | 5 +- .../integration/util/SimplePoolTests.java | 7 +- .../jdbc/lock/JdbcLockRegistry.java | 2 +- 5 files changed, 104 insertions(+), 39 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java b/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java index 806a57f0fd..f88312451d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiator.java @@ -55,7 +55,7 @@ import org.springframework.util.Assert; * * @author Dave Syer * @author Artem Bilan - * + * @author Vedran Pavic * @since 4.3.1 */ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBean, ApplicationEventPublisherAware { @@ -313,16 +313,7 @@ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBe if (acquired) { // Success: we are now leader this.locked = true; - LockRegistryLeaderInitiator.this.candidate.onGranted(this.context); - if (LockRegistryLeaderInitiator.this.leaderEventPublisher != null) { - try { - LockRegistryLeaderInitiator.this.leaderEventPublisher.publishOnGranted( - LockRegistryLeaderInitiator.this, this.context, this.lockKey); - } - catch (Exception e) { - logger.warn("Error publishing OnGranted event.", e); - } - } + handleGranted(); } } else if (acquired) { @@ -333,6 +324,9 @@ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBe Thread.sleep(LockRegistryLeaderInitiator.this.heartBeatMillis); } else { + this.locked = false; + // We were not able to acquire it, therefore not leading any more + handleRevoked(); // Try again quickly in case the lock holder dropped it Thread.sleep(LockRegistryLeaderInitiator.this.busyWaitMillis); } @@ -342,17 +336,7 @@ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBe this.lock.unlock(); this.locked = false; // The lock was broken and we are no longer leader - LockRegistryLeaderInitiator.this.candidate.onRevoked(this.context); - if (LockRegistryLeaderInitiator.this.leaderEventPublisher != null) { - try { - LockRegistryLeaderInitiator.this.leaderEventPublisher.publishOnRevoked( - LockRegistryLeaderInitiator.this, this.context, - LockRegistryLeaderInitiator.this.candidate.getRole()); - } - catch (Exception e1) { - logger.warn("Error publishing OnRevoked event.", e); - } - } + handleRevoked(); // Give it a chance to elect some other leader. Thread.sleep(LockRegistryLeaderInitiator.this.busyWaitMillis); Thread.currentThread().interrupt(); @@ -365,17 +349,7 @@ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBe this.lock.unlock(); if (this.locked) { // We are stopping, therefore not leading any more - LockRegistryLeaderInitiator.this.candidate.onRevoked(this.context); - if (LockRegistryLeaderInitiator.this.leaderEventPublisher != null) { - try { - LockRegistryLeaderInitiator.this.leaderEventPublisher.publishOnRevoked( - LockRegistryLeaderInitiator.this, this.context, - LockRegistryLeaderInitiator.this.candidate.getRole()); - } - catch (Exception e) { - logger.warn("Error publishing OnRevoked event.", e); - } - } + handleRevoked(); } this.locked = false; } @@ -386,6 +360,23 @@ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBe return this.locked; } + private void handleGranted() throws InterruptedException { + LockRegistryLeaderInitiator.this.candidate.onGranted(this.context); + if (LockRegistryLeaderInitiator.this.leaderEventPublisher != null) { + LockRegistryLeaderInitiator.this.leaderEventPublisher.publishOnGranted( + LockRegistryLeaderInitiator.this, this.context, this.lockKey); + } + } + + private void handleRevoked() { + LockRegistryLeaderInitiator.this.candidate.onRevoked(this.context); + if (LockRegistryLeaderInitiator.this.leaderEventPublisher != null) { + LockRegistryLeaderInitiator.this.leaderEventPublisher.publishOnRevoked( + LockRegistryLeaderInitiator.this, this.context, + LockRegistryLeaderInitiator.this.candidate.getRole()); + } + } + } /** diff --git a/spring-integration-core/src/test/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiatorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiatorTests.java index 1849dfeacd..983e11887d 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiatorTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/support/leader/LockRegistryLeaderInitiatorTests.java @@ -17,11 +17,19 @@ package org.springframework.integration.support.leader; import static org.hamcrest.CoreMatchers.is; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.mock; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.locks.Lock; import org.junit.Before; import org.junit.Test; @@ -36,6 +44,7 @@ import org.springframework.integration.support.locks.LockRegistry; /** * @author Dave Syer * @author Artem Bilan + * @author Vedran Pavic * * @since 4.3.1 */ @@ -120,11 +129,72 @@ public class LockRegistryLeaderInitiatorTests { initiator.start(); assertTrue(onGranted.await(10, TimeUnit.SECONDS)); - assertTrue(initiator.getContext().isLeader()); + assertFalse(initiator.getContext().isLeader()); initiator.stop(); } + @Test + public void competingWithLock() throws Exception { + // switch used to toggle which registry obtains lock + AtomicBoolean firstLocked = new AtomicBoolean(true); + + // set up first registry instance - this one will be able to obtain lock initially + LockRegistry firstRegistry = mock(LockRegistry.class); + Lock firstLock = mock(Lock.class); + given(firstRegistry.obtain(anyString())).willReturn(firstLock); + given(firstLock.tryLock(anyLong(), any(TimeUnit.class))).willAnswer(i -> firstLocked.get()); + + // set up first initiator instance using first LockRegistry + LockRegistryLeaderInitiator first = + new LockRegistryLeaderInitiator(firstRegistry, new DefaultCandidate()); + CountDownLatch firstGranted = new CountDownLatch(1); + CountDownLatch firstRevoked = new CountDownLatch(1); + first.setHeartBeatMillis(10); + first.setBusyWaitMillis(1); + first.setLeaderEventPublisher(new CountingPublisher(firstGranted, firstRevoked)); + + // set up second registry instance - this one will NOT be able to obtain lock initially + LockRegistry secondRegistry = mock(LockRegistry.class); + Lock secondLock = mock(Lock.class); + given(secondRegistry.obtain(anyString())).willReturn(secondLock); + given(secondLock.tryLock(anyLong(), any(TimeUnit.class))).willAnswer(i -> !firstLocked.get()); + + // set up second initiator instance using second LockRegistry + LockRegistryLeaderInitiator second = + new LockRegistryLeaderInitiator(secondRegistry, new DefaultCandidate()); + CountDownLatch secondGranted = new CountDownLatch(1); + CountDownLatch secondRevoked = new CountDownLatch(1); + second.setHeartBeatMillis(10); + second.setBusyWaitMillis(1); + second.setLeaderEventPublisher(new CountingPublisher(secondGranted, secondRevoked)); + + // start initiators + first.start(); + second.start(); + + Thread.sleep(100); + + // first initiator should lead and publish granted event + assertThat(first.getContext().isLeader(), is(true)); + assertThat(second.getContext().isLeader(), is(false)); + assertThat(firstGranted.await(10, TimeUnit.SECONDS), is(true)); + + // simulate first registry instance unable to obtain lock, for example due to lock timeout + firstLocked.set(false); + + Thread.sleep(100); + + // second initiator should take lead and publish granted event, first initiator should publish revoked event + assertThat(second.getContext().isLeader(), is(true)); + assertThat(first.getContext().isLeader(), is(false)); + assertThat(secondGranted.await(10, TimeUnit.SECONDS), is(true)); + assertThat(firstRevoked.await(10, TimeUnit.SECONDS), is(true)); + + first.stop(); + second.stop(); + } + private static class CountingPublisher implements LeaderEventPublisher { private final CountDownLatch granted; diff --git a/spring-integration-core/src/test/java/org/springframework/integration/util/CallerBlocksPolicyTests.java b/spring-integration-core/src/test/java/org/springframework/integration/util/CallerBlocksPolicyTests.java index 4d4db12413..582811688c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/util/CallerBlocksPolicyTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/util/CallerBlocksPolicyTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2016 the original author or authors. + * Copyright 2014-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -34,6 +34,7 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; /** * @author Gary Russell + * @author Artem Bilan * @since 3.0.3 * */ @@ -66,6 +67,7 @@ public class CallerBlocksPolicyTests { assertTrue(latch.await(10, TimeUnit.SECONDS)); assertThat(e.get(), instanceOf(RejectedExecutionException.class)); assertEquals("Max wait time expired to queue task", e.get().getMessage()); + te.destroy(); } @Test @@ -104,6 +106,7 @@ public class CallerBlocksPolicyTests { }); assertTrue(latch.await(10, TimeUnit.SECONDS)); assertNull(e.get()); + te.destroy(); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/util/SimplePoolTests.java b/spring-integration-core/src/test/java/org/springframework/integration/util/SimplePoolTests.java index 1d598ee129..34a22ea8d5 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/util/SimplePoolTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/util/SimplePoolTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -59,7 +59,7 @@ public class SimplePoolTests { } @Test - public void testOverCommitandResize() { + public void testOverCommitAndResize() { final Set strings = new HashSet(); final AtomicBoolean stale = new AtomicBoolean(); SimplePool pool = stringPool(2, strings, stale); @@ -155,7 +155,7 @@ public class SimplePoolTests { SimplePool pool = new SimplePool(size, new SimplePool.PoolItemCallback() { private int i; public String createForPool() { - String string = new String("String" + i++); + String string = "String" + i++; strings.add(string); return string; } @@ -171,4 +171,5 @@ public class SimplePoolTests { }); return pool; } + } diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/lock/JdbcLockRegistry.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/lock/JdbcLockRegistry.java index 339faf240b..59dc96e9cd 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/lock/JdbcLockRegistry.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/lock/JdbcLockRegistry.java @@ -44,7 +44,7 @@ import org.springframework.util.Assert; * * @author Dave Syer * @author Artem Bilan - * @author Vedran Pavić + * @author Vedran Pavic * * @since 4.3 */