Fixing some consistency issues in the behaviour of the Resequencer
This commit is contained in:
@@ -85,6 +85,9 @@ public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message
|
||||
}
|
||||
|
||||
private boolean hasReceivedAllMessages(MessageBarrier<SortedSet<Message<?>>> 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<SortedSet<Message
|
||||
logger.debug("A message with the same sequence number has been already received: " + message);
|
||||
return false;
|
||||
}
|
||||
// one can always add a message to the barrier if it's empty. Afterwards, assume that the complete sequence size
|
||||
//
|
||||
if (!barrier.getMessages().isEmpty() &&
|
||||
barrier.getMessages().first().getHeaders().getSequenceSize() < message.getHeaders().getSequenceNumber()) {
|
||||
if (message.getHeaders().getSequenceSize() < message.getHeaders().getSequenceNumber()) {
|
||||
logger.debug("The message has a sequence number which is larger than the sequence size: "+ message);
|
||||
return false;
|
||||
}
|
||||
if (!barrier.getMessages().isEmpty() &&
|
||||
message.getHeaders().getSequenceSize() != barrier.getMessages().first().getHeaders().getSequenceSize()) {
|
||||
logger.debug("The message has a sequence size which is different from other messages handled so far: " + message
|
||||
+ ", expected value is " + barrier.getMessages().first().getHeaders().getSequenceNumber());
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -131,7 +131,6 @@ public class ResequencerTests {
|
||||
Message<?> 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 {
|
||||
|
||||
Reference in New Issue
Block a user