INT-3483: Fix AbstractCorrelatingMH deadlock

JIRA: https://jira.spring.io/browse/INT-3483

**Cherry-pick to 4.0.x & 3.0.x**
This commit is contained in:
Artem Bilan
2014-07-25 12:20:07 +03:00
committed by Gary Russell
parent 917a970b27
commit 9b2fb8a07f
2 changed files with 86 additions and 19 deletions

View File

@@ -532,10 +532,14 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
}
}
finally {
if (removeGroup) {
this.remove(group);
try {
if (removeGroup) {
this.remove(group);
}
}
finally {
lock.unlock();
}
lock.unlock();
}
}
catch (InterruptedException ie) {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2013 the original author or authors.
* Copyright 2002-2014 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.
@@ -13,16 +13,11 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.aggregator;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
import static org.junit.Assert.*;
import static org.mockito.Mockito.*;
import java.lang.reflect.Method;
import java.util.ArrayList;
@@ -30,6 +25,7 @@ import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
@@ -49,6 +45,7 @@ import org.springframework.messaging.support.GenericMessage;
/**
* @author Gary Russell
* @author Artem Bilan
* @since 2.2
*
*/
@@ -285,7 +282,8 @@ public class AbstractCorrelatingMessageHandlerTests {
QueueChannel outputChannel = new QueueChannel();
handler.setOutputChannel(outputChannel);
MessageGroupStore mgs = TestUtils.getPropertyValue(handler, "messageStore", MessageGroupStore.class);
Method forceComplete = AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
Method forceComplete =
AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
forceComplete.setAccessible(true);
mgs.addMessageToGroup("foo", new GenericMessage<String>("foo"));
GenericMessage<String> secondMessage = new GenericMessage<String>("bar");
@@ -322,9 +320,11 @@ public class AbstractCorrelatingMessageHandlerTests {
mgs.completeGroup("foo");
mgs = spy(mgs);
new DirectFieldAccessor(handler).setPropertyValue("messageStore", mgs);
Method forceComplete = AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
Method forceComplete =
AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
forceComplete.setAccessible(true);
MessageGroup group = (MessageGroup) TestUtils.getPropertyValue(mgs, "groupIdToMessageGroup", Map.class).get("foo");
MessageGroup group = (MessageGroup) TestUtils.getPropertyValue(mgs, "groupIdToMessageGroup", Map.class)
.get("foo");
assertTrue(group.isComplete());
forceComplete.invoke(handler, group);
verify(mgs, never()).getMessageGroup("foo");
@@ -354,9 +354,11 @@ public class AbstractCorrelatingMessageHandlerTests {
mgs.completeGroup("foo");
mgs = spy(mgs);
new DirectFieldAccessor(handler).setPropertyValue("messageStore", mgs);
Method forceComplete = AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
Method forceComplete =
AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
forceComplete.setAccessible(true);
MessageGroup groupInStore = (MessageGroup) TestUtils.getPropertyValue(mgs, "groupIdToMessageGroup", Map.class).get("foo");
MessageGroup groupInStore = (MessageGroup) TestUtils.getPropertyValue(mgs, "groupIdToMessageGroup", Map.class)
.get("foo");
assertTrue(groupInStore.isComplete());
assertFalse(group.isComplete());
new DirectFieldAccessor(group).setPropertyValue("lastModified", groupInStore.getLastModified());
@@ -387,9 +389,11 @@ public class AbstractCorrelatingMessageHandlerTests {
MessageGroup group = new SimpleMessageGroup(mgs.getMessageGroup("foo"));
mgs = spy(mgs);
new DirectFieldAccessor(handler).setPropertyValue("messageStore", mgs);
Method forceComplete = AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
Method forceComplete =
AbstractCorrelatingMessageHandler.class.getDeclaredMethod("forceComplete", MessageGroup.class);
forceComplete.setAccessible(true);
MessageGroup groupInStore = (MessageGroup) TestUtils.getPropertyValue(mgs, "groupIdToMessageGroup", Map.class).get("foo");
MessageGroup groupInStore = (MessageGroup) TestUtils.getPropertyValue(mgs, "groupIdToMessageGroup", Map.class)
.get("foo");
assertFalse(groupInStore.isComplete());
assertFalse(group.isComplete());
DirectFieldAccessor directFieldAccessor = new DirectFieldAccessor(group);
@@ -400,4 +404,63 @@ public class AbstractCorrelatingMessageHandlerTests {
assertNull(outputChannel.receive(0));
}
@Test
public void testInt3483DeadlockOnMessageStoreRemoveMessageGroup() throws InterruptedException {
final AggregatingMessageHandler handler =
new AggregatingMessageHandler(new DefaultAggregatingMessageGroupProcessor());
handler.setOutputChannel(new QueueChannel());
QueueChannel discardChannel = new QueueChannel();
handler.setDiscardChannel(discardChannel);
handler.setReleaseStrategy(new ReleaseStrategy() {
@Override
public boolean canRelease(MessageGroup group) {
return true;
}
});
handler.setExpireGroupsUponTimeout(false);
SimpleMessageStore messageStore = new SimpleMessageStore() {
@Override
public void removeMessageGroup(Object groupId) {
throw new RuntimeException("intentional");
}
};
handler.setMessageStore(messageStore);
handler.handleMessage(MessageBuilder.withPayload("foo")
.setCorrelationId(1)
.setSequenceNumber(1)
.setSequenceSize(2)
.build());
try {
messageStore.expireMessageGroups(0);
}
catch (Exception e) {
//suppress an intentional 'removeMessageGroup' exception
}
ExecutorService executorService = Executors.newSingleThreadExecutor();
executorService.execute(new Runnable() {
@Override
public void run() {
handler.handleMessage(MessageBuilder.withPayload("foo")
.setCorrelationId(1)
.setSequenceNumber(2)
.setSequenceSize(2)
.build());
}
});
executorService.shutdown();
/* Previously lock for the groupId hasn't been unlocked from the 'forceComplete', because it wasn't
reachable in case of exception from the BasicMessageGroupStore.removeMessageGroup
*/
assertTrue(executorService.awaitTermination(10, TimeUnit.SECONDS));
/* Since MessageGroup had been marked as 'complete', but hasn't been removed because of exception,
the second message is discarded
*/
Message<?> receive = discardChannel.receive(1000);
assertNotNull(receive);
}
}