INT-4105: RedisLock: Unlock Local in obtain()
JIRA: https://jira.spring.io/browse/INT-4105 Fixes GH-1888 The lock in Redis can be expired in between `obtain()` calls. So, even if we return a new lock instance, the old one must clear properly. * Add `lock.unlock()` to the `obtain()` if the case of expiration in Redis. * Add `warn` for the exception on the `lock.unlock()`. We don't care about error here and just proceed to a new instance. **Chery-pick to 4.3.x** Remove the lock reference from `weakThreadLocks` as well Add `registry.setUseWeakReferences(true);` to test-case * Add assert for `lock.unlock()` in case of expiration * Add `UUID.randomUUID()` for the registry key to avoid cross-talking during concurrent builds, e.g. on the CI server https://build.spring.io/browse/INT-MASTER-352 This is actually a fix for the JIRA: https://jira.spring.io/browse/INT-4083
This commit is contained in:
committed by
Gary Russell
parent
763547562c
commit
450ac7f18d
@@ -259,7 +259,20 @@ public final class RedisLockRegistry implements LockRegistry {
|
||||
if (lock != null && lock.thread != null) {
|
||||
RedisLock lockInStore = this.redisTemplate.boundValueOps(this.registryKey + ":" + lockKey).get();
|
||||
if (lockInStore == null || !lock.equals(lockInStore)) {
|
||||
getHardThreadLocks().remove(lock);
|
||||
try {
|
||||
lock.unlock();
|
||||
}
|
||||
catch (Exception e) {
|
||||
if (logger.isWarnEnabled()) {
|
||||
logger.warn("Lock was released due to expiration. A new one will be obtained...", e);
|
||||
}
|
||||
}
|
||||
if (this.hardThreadLocks.get() != null) {
|
||||
this.hardThreadLocks.get().remove(lock);
|
||||
}
|
||||
if (this.weakThreadLocks.get() != null) {
|
||||
this.weakThreadLocks.get().remove(lock);
|
||||
}
|
||||
lock = null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -30,6 +30,7 @@ import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Executors;
|
||||
@@ -57,6 +58,7 @@ import org.springframework.integration.test.util.TestUtils;
|
||||
/**
|
||||
* @author Gary Russell
|
||||
* @author Konstantin Yakimov
|
||||
* @author Artem Bilan
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
@@ -64,6 +66,10 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
|
||||
private final Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
private final String registryKey = UUID.randomUUID().toString();
|
||||
|
||||
private final String registryKey2 = UUID.randomUUID().toString();
|
||||
|
||||
@Rule
|
||||
public Log4jLevelAdjuster adjuster = new Log4jLevelAdjuster(Level.TRACE, "org.springframework.integration.redis");
|
||||
|
||||
@@ -71,8 +77,8 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@After
|
||||
public void setupShutDown() {
|
||||
RedisTemplate<String, ?> template = this.createTemplate();
|
||||
template.delete("rlrTests");
|
||||
template.delete("rlrTests2");
|
||||
template.delete(this.registryKey + ":*");
|
||||
template.delete(this.registryKey2 + ":*");
|
||||
}
|
||||
|
||||
private RedisTemplate<String, ?> createTemplate() {
|
||||
@@ -86,7 +92,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testLock() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Lock lock = registry.obtain("foo");
|
||||
lock.lock();
|
||||
@@ -103,7 +109,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testLockInterruptibly() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Lock lock = registry.obtain("foo");
|
||||
lock.lockInterruptibly();
|
||||
@@ -120,7 +126,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testReentrantLock() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Lock lock1 = registry.obtain("foo");
|
||||
lock1.lock();
|
||||
@@ -145,7 +151,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testReentrantLockInterruptibly() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Lock lock1 = registry.obtain("foo");
|
||||
lock1.lockInterruptibly();
|
||||
@@ -170,7 +176,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testTwoLocks() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Lock lock1 = registry.obtain("foo");
|
||||
lock1.lockInterruptibly();
|
||||
@@ -195,7 +201,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testTwoThreadsSecondFailsToGetLock() throws Exception {
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
final Lock lock1 = registry.obtain("foo");
|
||||
lock1.lockInterruptibly();
|
||||
final AtomicBoolean locked = new AtomicBoolean();
|
||||
@@ -228,7 +234,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testTwoThreads() throws Exception {
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
final Lock lock1 = registry.obtain("foo");
|
||||
final AtomicBoolean locked = new AtomicBoolean();
|
||||
final CountDownLatch latch1 = new CountDownLatch(1);
|
||||
@@ -269,8 +275,8 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testTwoThreadsDifferentRegistries() throws Exception {
|
||||
final RedisLockRegistry registry1 = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
final RedisLockRegistry registry2 = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
final RedisLockRegistry registry1 = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
final RedisLockRegistry registry2 = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
final Lock lock1 = registry1.obtain("foo");
|
||||
final AtomicBoolean locked = new AtomicBoolean();
|
||||
final CountDownLatch latch1 = new CountDownLatch(1);
|
||||
@@ -319,7 +325,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testTwoThreadsWrongOneUnlocks() throws Exception {
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
final Lock lock = registry.obtain("foo");
|
||||
lock.lockInterruptibly();
|
||||
final AtomicBoolean locked = new AtomicBoolean();
|
||||
@@ -350,7 +356,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testList() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests");
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey);
|
||||
Lock foo = registry.obtain("foo");
|
||||
foo.lockInterruptibly();
|
||||
Lock bar = registry.obtain("bar");
|
||||
@@ -368,7 +374,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testExpireNoLockInStore() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests", 1000);
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey, 100);
|
||||
Lock foo = registry.obtain("foo");
|
||||
foo.lockInterruptibly();
|
||||
waitForExpire("foo");
|
||||
@@ -382,10 +388,30 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
assertNull(TestUtils.getPropertyValue(registry, "hardThreadLocks", ThreadLocal.class).get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testExpireDuringSecondObtain() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey, 100);
|
||||
registry.setUseWeakReferences(true);
|
||||
Lock foo = registry.obtain("foo");
|
||||
foo.lockInterruptibly();
|
||||
waitForExpire("foo");
|
||||
Lock foo1 = registry.obtain("foo");
|
||||
assertNotSame(foo, foo1);
|
||||
|
||||
try {
|
||||
foo.unlock();
|
||||
fail("IllegalStateException");
|
||||
}
|
||||
catch (IllegalStateException e) {
|
||||
assertThat(e.getMessage(), containsString("Lock is not locked"));
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testExpireNewLockInStore() throws Exception {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests", 1000);
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey, 100);
|
||||
Lock foo1 = registry.obtain("foo");
|
||||
foo1.lockInterruptibly();
|
||||
waitForExpire("foo");
|
||||
@@ -397,8 +423,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
fail("Expected exception");
|
||||
}
|
||||
catch (IllegalStateException e) {
|
||||
assertThat(e.getMessage(), containsString("Lock was released due to expiration"));
|
||||
assertThat(e.getMessage(), containsString("lock in store:"));
|
||||
assertThat(e.getMessage(), containsString("Lock is not locked"));
|
||||
}
|
||||
foo2.unlock();
|
||||
assertNull(TestUtils.getPropertyValue(registry, "hardThreadLocks", ThreadLocal.class).get());
|
||||
@@ -408,10 +433,10 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@RedisAvailable
|
||||
public void testEquals() throws Exception {
|
||||
RedisConnectionFactory connectionFactory = this.getConnectionFactoryForTest();
|
||||
RedisLockRegistry registry1 = new RedisLockRegistry(connectionFactory, "rlrTests");
|
||||
RedisLockRegistry registry1 = new RedisLockRegistry(connectionFactory, this.registryKey);
|
||||
registry1.setUseWeakReferences(true);
|
||||
RedisLockRegistry registry2 = new RedisLockRegistry(connectionFactory, "rlrTests");
|
||||
RedisLockRegistry registry3 = new RedisLockRegistry(connectionFactory, "rlrTests2");
|
||||
RedisLockRegistry registry2 = new RedisLockRegistry(connectionFactory, this.registryKey);
|
||||
RedisLockRegistry registry3 = new RedisLockRegistry(connectionFactory, this.registryKey2);
|
||||
Lock lock1 = registry1.obtain("foo");
|
||||
Lock lock2 = registry1.obtain("foo");
|
||||
assertEquals(lock1, lock2);
|
||||
@@ -441,7 +466,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testThreadLocalListLeaks() {
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), "rlrTests", 1000);
|
||||
RedisLockRegistry registry = new RedisLockRegistry(this.getConnectionFactoryForTest(), this.registryKey, 100);
|
||||
registry.setUseWeakReferences(true);
|
||||
|
||||
for (int i = 0; i < 10; i++) {
|
||||
@@ -469,7 +494,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
@RedisAvailable
|
||||
public void testExpireNotChanged() throws Exception {
|
||||
RedisConnectionFactory connectionFactory = this.getConnectionFactoryForTest();
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, "rlrTests", 10000);
|
||||
final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000);
|
||||
Lock lock = registry.obtain("foo");
|
||||
lock.lock();
|
||||
|
||||
@@ -499,7 +524,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests {
|
||||
private void waitForExpire(String key) throws Exception {
|
||||
RedisTemplate<String, ?> template = this.createTemplate();
|
||||
int n = 0;
|
||||
while (n++ < 100 && template.keys("rlrTests:" + key).size() > 0) {
|
||||
while (n++ < 100 && template.keys(this.registryKey + ":" + key).size() > 0) {
|
||||
Thread.sleep(100);
|
||||
}
|
||||
assertTrue(key + " key did not expire", n < 100);
|
||||
|
||||
Reference in New Issue
Block a user