From 081d0d160bff7285545c199cdd94a6716c63a38d Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 8 May 2018 19:02:14 -0400 Subject: [PATCH] Fix one more race condition in RedisLeaderTests https://build.spring.io/browse/INT-SI50X-JOB1-56 We can't wait for the latch in the interruptable code flow; we can't have a round-robing election guarantees. * Add `Thread.sleep(LockRegistryLeaderInitiator.this.busyWaitMillis)` to the `LockRegistryLeaderInitiator` when we restart the main task * Remove latches waiting and thread shifting from the `RedisLockRegistryLeaderInitiatorTests` * Use long `busyWaitMillis` for yielding initiator to let the second candidate to be elected **Cherry-pick to 5.0.x** --- .../leader/LockRegistryLeaderInitiator.java | 7 ++- ...RedisLockRegistryLeaderInitiatorTests.java | 56 +++++-------------- 2 files changed, 20 insertions(+), 43 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 fb2e503110..4e2a6cc27b 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 @@ -408,7 +408,12 @@ public class LockRegistryLeaderInitiator implements SmartLifecycle, DisposableBe if (isRunning()) { logger.warn("Restarting LeaderSelector for " + this.context + " because of error.", e); LockRegistryLeaderInitiator.this.future = - LockRegistryLeaderInitiator.this.executorService.submit(this); + LockRegistryLeaderInitiator.this.executorService.submit( + () -> { + // Give it a chance to elect some other leader. + Thread.sleep(LockRegistryLeaderInitiator.this.busyWaitMillis); + return call(); + }); } return null; } diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/leader/RedisLockRegistryLeaderInitiatorTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/leader/RedisLockRegistryLeaderInitiatorTests.java index ba581a7ccd..7011f8d86f 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/leader/RedisLockRegistryLeaderInitiatorTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/leader/RedisLockRegistryLeaderInitiatorTests.java @@ -23,12 +23,9 @@ import static org.junit.Assert.assertThat; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.Executor; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; import org.junit.Rule; import org.junit.Test; @@ -40,8 +37,8 @@ import org.springframework.integration.redis.rules.RedisAvailableTests; import org.springframework.integration.redis.util.RedisLockRegistry; import org.springframework.integration.support.leader.LockRegistryLeaderInitiator; import org.springframework.integration.test.rule.Log4j2LevelAdjuster; +import org.springframework.integration.test.support.LongRunningIntegrationTest; import org.springframework.scheduling.concurrent.CustomizableThreadFactory; -import org.springframework.util.ReflectionUtils; /** * @author Artem Bilan @@ -52,7 +49,8 @@ import org.springframework.util.ReflectionUtils; */ public class RedisLockRegistryLeaderInitiatorTests extends RedisAvailableTests { - private static final Log logger = LogFactory.getLog(RedisLockRegistryLeaderInitiatorTests.class); + @Rule + public LongRunningIntegrationTest longTests = new LongRunningIntegrationTest(); @Rule public Log4j2LevelAdjuster adjuster = @@ -105,55 +103,28 @@ public class RedisLockRegistryLeaderInitiatorTests extends RedisAvailableTests { CountDownLatch acquireLockFailed1 = new CountDownLatch(1); CountDownLatch acquireLockFailed2 = new CountDownLatch(1); - Executor latchesExecutor = Executors.newCachedThreadPool(); + initiator1.setLeaderEventPublisher(new CountingPublisher(granted1, revoked1, acquireLockFailed1)); - initiator1.setLeaderEventPublisher(new CountingPublisher(granted1, revoked1, acquireLockFailed1) { + initiator2.setLeaderEventPublisher(new CountingPublisher(granted2, revoked2, acquireLockFailed2)); - @Override - public void publishOnRevoked(Object source, Context context, String role) { - latchesExecutor.execute(() -> { - // It's difficult to see round-robin election, so block one initiator until the second is elected. - try { - assertThat(granted2.await(10, TimeUnit.SECONDS), is(true)); - } - catch (InterruptedException e) { - ReflectionUtils.rethrowRuntimeException(e); - } - super.publishOnRevoked(source, context, role); - }); - - } - - }); - - initiator2.setLeaderEventPublisher(new CountingPublisher(granted2, revoked2, acquireLockFailed2) { - - @Override - public void publishOnRevoked(Object source, Context context, String role) { - latchesExecutor.execute(() -> { - try { - // It's difficult to see round-robin election, so block one initiator until the second is elected. - assertThat(granted1.await(10, TimeUnit.SECONDS), is(true)); - } - catch (InterruptedException e) { - ReflectionUtils.rethrowRuntimeException(e); - } - super.publishOnRevoked(source, context, role); - }); - } - - }); + // It's hard to see round-robin election, so let's make the yielding initiator to sleep long before restarting + initiator1.setBusyWaitMillis(5000); initiator1.getContext().yield(); assertThat(revoked1.await(10, TimeUnit.SECONDS), is(true)); + assertThat(granted2.await(10, TimeUnit.SECONDS), is(true)); assertThat(initiator2.getContext().isLeader(), is(true)); assertThat(initiator1.getContext().isLeader(), is(false)); + initiator1.setBusyWaitMillis(LockRegistryLeaderInitiator.DEFAULT_BUSY_WAIT_TIME); + initiator2.setBusyWaitMillis(5000); + initiator2.getContext().yield(); assertThat(revoked2.await(10, TimeUnit.SECONDS), is(true)); + assertThat(granted1.await(10, TimeUnit.SECONDS), is(true)); assertThat(initiator1.getContext().isLeader(), is(true)); assertThat(initiator2.getContext().isLeader(), is(false)); @@ -161,7 +132,8 @@ public class RedisLockRegistryLeaderInitiatorTests extends RedisAvailableTests { initiator2.stop(); CountDownLatch revoked11 = new CountDownLatch(1); - initiator1.setLeaderEventPublisher(new CountingPublisher(new CountDownLatch(1), revoked11, new CountDownLatch(1))); + initiator1.setLeaderEventPublisher(new CountingPublisher(new CountDownLatch(1), revoked11, + new CountDownLatch(1))); initiator1.getContext().yield();