diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java index bd8bb8feb9..b231753e2b 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java @@ -85,6 +85,9 @@ public class Resequencer extends AbstractMessageBarrierHandler>> barrier) { + if(barrier.getMessages().isEmpty()) { + return false; + } int sequenceSize = barrier.getMessages().first().getHeaders().getSequenceSize(); int messagesCurrentlyInBarrier = barrier.getMessages().size(); int lastReleasedSequenceNumber = barrier.getAttribute(LAST_RELEASED_SEQUENCE_NUMBER); @@ -126,13 +129,16 @@ public class Resequencer extends AbstractMessageBarrierHandler message1 = createMessage("123", "ABC", 4, 2, null); Message message2 = createMessage("456", "ABC", 4, 1, null); Message message3 = createMessage("789", "ABC", 4, 4, null); - Message message4 = createMessage("XYZ", "ABC", 4, 3, null); this.resequencer.setSendPartialResultOnTimeout(false); this.resequencer.setReleasePartialSequences(false); this.resequencer.setDiscardChannel(discardChannel); @@ -149,13 +148,47 @@ public class ResequencerTests { assertEquals(new Integer(2), reply2.getHeaders().getSequenceNumber()); assertNull(reply3); // when sending the last message, the whole sequence must have been sent - this.resequencer.handleMessage(message4); + this.resequencer.handleMessage(message3); reply3 = discardChannel.receive(0); assertNull(reply3); - Message reply4 = discardChannel.receive(0); - assertNull(reply4); } - + + + @Test + public void testResequencingWithDifferentSequenceSizes() throws InterruptedException { + QueueChannel discardChannel = new QueueChannel(); + Message message1 = createMessage("123", "ABC", 4, 2, null); + Message message2 = createMessage("456", "ABC", 5, 1, null); + this.resequencer.setSendPartialResultOnTimeout(false); + this.resequencer.setReleasePartialSequences(false); + this.resequencer.setDiscardChannel(discardChannel); + this.resequencer.setTimeout(90000); + this.resequencer.handleMessage(message1); + this.resequencer.handleMessage(message2); + this.resequencer.discardBarrier(this.resequencer.barriers.get("ABC")); + Message reply1 = discardChannel.receive(0); + Message reply2 = discardChannel.receive(0); + // only messages 1 - with sequence number 2 - should have been received by now + // the other has been discarded + assertNotNull(reply1); + assertEquals(new Integer(2), reply1.getHeaders().getSequenceNumber()); + assertNull(reply2); + } + + @Test + public void testResequencingWithWrongSequenceSizeAndNumber() throws InterruptedException { + QueueChannel discardChannel = new QueueChannel(); + Message message1 = createMessage("123", "ABC", 2, 4, null); + this.resequencer.setSendPartialResultOnTimeout(false); + this.resequencer.setReleasePartialSequences(false); + this.resequencer.setDiscardChannel(discardChannel); + this.resequencer.setTimeout(90000); + this.resequencer.handleMessage(message1); + this.resequencer.discardBarrier(this.resequencer.barriers.get("ABC")); + Message reply1 = discardChannel.receive(0); + // No message has been received - the message has been rejected. + assertNull(reply1); + } @Test public void testResequencingWithCompleteSequenceRelease() throws InterruptedException {