From 9b2fb8a07f982fdaca753387b96bad5e13ad10cc Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 25 Jul 2014 12:20:07 +0300 Subject: [PATCH] INT-3483: Fix `AbstractCorrelatingMH` deadlock JIRA: https://jira.spring.io/browse/INT-3483 **Cherry-pick to 4.0.x & 3.0.x** --- .../AbstractCorrelatingMessageHandler.java | 10 +- ...bstractCorrelatingMessageHandlerTests.java | 95 +++++++++++++++---- 2 files changed, 86 insertions(+), 19 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index bc7ef3ae94..a78d06b432 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -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) { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandlerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandlerTests.java index 279257a6e2..991939c985 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandlerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandlerTests.java @@ -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("foo")); GenericMessage secondMessage = new GenericMessage("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); + } + }