INT-4349: Allow Number type for int headers
JIRA: https://jira.spring.io/browse/INT-4349 When headers come from the external system there is no guarantee that special headers (e.g. `sequenceNumber`, `priority` etc.) in the expected (`Integer`) type. * Widen `int` headers setting value to the `Number` type * Return primitive `int` for the `sequenceNumber` and `sequenceSize` headers since for them `IntegrationMessageHeaderAccessor` never return null for them * Remove `SequenceNumberComparator` in favor of `MessageSequenceComparator` since they are essentially duplicate each other * Fix tests to deal with primitive `int` already
This commit is contained in:
committed by
Gary Russell
parent
e6225926c4
commit
c65584a007
@@ -89,18 +89,19 @@ public class IntegrationMessageHeaderAccessor extends MessageHeaderAccessor {
|
||||
return this.getHeader(CORRELATION_ID);
|
||||
}
|
||||
|
||||
public Integer getSequenceNumber() {
|
||||
Integer sequenceNumber = this.getHeader(SEQUENCE_NUMBER, Integer.class);
|
||||
return (sequenceNumber != null ? sequenceNumber : 0);
|
||||
public int getSequenceNumber() {
|
||||
Number sequenceNumber = this.getHeader(SEQUENCE_NUMBER, Number.class);
|
||||
return (sequenceNumber != null ? sequenceNumber.intValue() : 0);
|
||||
}
|
||||
|
||||
public Integer getSequenceSize() {
|
||||
Integer sequenceSize = this.getHeader(SEQUENCE_SIZE, Integer.class);
|
||||
return (sequenceSize != null ? sequenceSize : 0);
|
||||
public int getSequenceSize() {
|
||||
Number sequenceSize = this.getHeader(SEQUENCE_SIZE, Number.class);
|
||||
return (sequenceSize != null ? sequenceSize.intValue() : 0);
|
||||
}
|
||||
|
||||
public Integer getPriority() {
|
||||
return this.getHeader(PRIORITY, Integer.class);
|
||||
Number priority = this.getHeader(PRIORITY, Number.class);
|
||||
return (priority != null ? priority.intValue() : null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -140,8 +141,8 @@ public class IntegrationMessageHeaderAccessor extends MessageHeaderAccessor {
|
||||
else if (IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER.equals(headerName)
|
||||
|| IntegrationMessageHeaderAccessor.SEQUENCE_SIZE.equals(headerName)
|
||||
|| IntegrationMessageHeaderAccessor.PRIORITY.equals(headerName)) {
|
||||
Assert.isTrue(Integer.class.isAssignableFrom(headerValue.getClass()), "The '" + headerName
|
||||
+ "' header value must be an Integer.");
|
||||
Assert.isTrue(Number.class.isAssignableFrom(headerValue.getClass()), "The '" + headerName
|
||||
+ "' header value must be a Number.");
|
||||
}
|
||||
else if (IntegrationMessageHeaderAccessor.ROUTING_SLIP.equals(headerName)) {
|
||||
Assert.isTrue(Map.class.isAssignableFrom(headerValue.getClass()), "The '" + headerName
|
||||
|
||||
@@ -90,7 +90,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
|
||||
protected final Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
private final Comparator<Message<?>> sequenceNumberComparator = new SequenceNumberComparator();
|
||||
private final Comparator<Message<?>> sequenceNumberComparator = new MessageSequenceComparator();
|
||||
|
||||
private final Map<UUID, ScheduledFuture<?>> expireGroupScheduledFutures = new HashMap<>();
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2008 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -22,23 +22,18 @@ import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.messaging.Message;
|
||||
|
||||
/**
|
||||
* A {@link Comparator} implementation based on the 'sequence number'
|
||||
* property of a {@link Message Message's} header.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Dave Syer
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
public class MessageSequenceComparator implements Comparator<Message<?>> {
|
||||
|
||||
public int compare(Message<?> message1, Message<?> message2) {
|
||||
Integer s1 = new IntegrationMessageHeaderAccessor(message1).getSequenceNumber();
|
||||
Integer s2 = new IntegrationMessageHeaderAccessor(message2).getSequenceNumber();
|
||||
if (s1 == null) {
|
||||
s1 = 0;
|
||||
}
|
||||
if (s2 == null) {
|
||||
s2 = 0;
|
||||
}
|
||||
return s1.compareTo(s2);
|
||||
@Override
|
||||
public int compare(Message<?> o1, Message<?> o2) {
|
||||
int sequenceNumber1 = new IntegrationMessageHeaderAccessor(o1).getSequenceNumber();
|
||||
int sequenceNumber2 = new IntegrationMessageHeaderAccessor(o2).getSequenceNumber();
|
||||
|
||||
return Integer.compare(sequenceNumber1, sequenceNumber2);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -36,7 +36,7 @@ import org.springframework.messaging.Message;
|
||||
*/
|
||||
public class ResequencingMessageGroupProcessor implements MessageGroupProcessor {
|
||||
|
||||
private final Comparator<Message<?>> comparator = new SequenceNumberComparator();
|
||||
private final Comparator<Message<?>> comparator = new MessageSequenceComparator();
|
||||
|
||||
public Object processMessageGroup(MessageGroup group) {
|
||||
Collection<Message<?>> messages = group.getMessages();
|
||||
|
||||
@@ -1,53 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2016 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.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.integration.aggregator;
|
||||
|
||||
import java.util.Comparator;
|
||||
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.messaging.Message;
|
||||
|
||||
/**
|
||||
* @author Dave Syer
|
||||
*
|
||||
* @since 2.0
|
||||
*
|
||||
*/
|
||||
public class SequenceNumberComparator implements Comparator<Message<?>> {
|
||||
|
||||
/**
|
||||
* If both messages have a sequence number then compare that, otherwise if one has a sequence number and the other
|
||||
* doesn't then the numbered message comes first, or finally of neither has a sequence number then they are equal in
|
||||
* rank.
|
||||
*/
|
||||
@Override
|
||||
public int compare(Message<?> o1, Message<?> o2) {
|
||||
Integer sequenceNumber1 = new IntegrationMessageHeaderAccessor(o1).getSequenceNumber();
|
||||
Integer sequenceNumber2 = new IntegrationMessageHeaderAccessor(o2).getSequenceNumber();
|
||||
if (sequenceNumber1 == sequenceNumber2) { //NOSONAR - early exit optimization
|
||||
return 0;
|
||||
}
|
||||
if (sequenceNumber1 == null) {
|
||||
return -sequenceNumber2;
|
||||
}
|
||||
if (sequenceNumber2 == null) {
|
||||
return sequenceNumber1;
|
||||
}
|
||||
return sequenceNumber1.compareTo(sequenceNumber2);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -44,7 +44,7 @@ public class SequenceSizeReleaseStrategy implements ReleaseStrategy {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(SequenceSizeReleaseStrategy.class);
|
||||
|
||||
private final Comparator<Message<?>> comparator = new SequenceNumberComparator();
|
||||
private final Comparator<Message<?>> comparator = new MessageSequenceComparator();
|
||||
|
||||
private volatile boolean releasePartialSequences;
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2008 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -18,6 +18,8 @@ package org.springframework.integration.aggregator;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
import java.util.Comparator;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
@@ -25,12 +27,13 @@ import org.springframework.integration.support.MessageBuilder;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
public class MessageSequenceComparatorTests {
|
||||
|
||||
@Test
|
||||
public void testLessThan() {
|
||||
MessageSequenceComparator comparator = new MessageSequenceComparator();
|
||||
Comparator<Message<?>> comparator = new MessageSequenceComparator();
|
||||
Message<String> message1 = MessageBuilder.withPayload("test1")
|
||||
.setSequenceNumber(1).build();
|
||||
Message<String> message2 = MessageBuilder.withPayload("test2")
|
||||
@@ -40,7 +43,7 @@ public class MessageSequenceComparatorTests {
|
||||
|
||||
@Test
|
||||
public void testEqual() {
|
||||
MessageSequenceComparator comparator = new MessageSequenceComparator();
|
||||
Comparator<Message<?>> comparator = new MessageSequenceComparator();
|
||||
Message<String> message1 = MessageBuilder.withPayload("test1")
|
||||
.setSequenceNumber(3).build();
|
||||
Message<String> message2 = MessageBuilder.withPayload("test2")
|
||||
@@ -50,7 +53,7 @@ public class MessageSequenceComparatorTests {
|
||||
|
||||
@Test
|
||||
public void testGreaterThan() {
|
||||
MessageSequenceComparator comparator = new MessageSequenceComparator();
|
||||
Comparator<Message<?>> comparator = new MessageSequenceComparator();
|
||||
Message<String> message1 = MessageBuilder.withPayload("test1")
|
||||
.setSequenceNumber(5).build();
|
||||
Message<String> message2 = MessageBuilder.withPayload("test2")
|
||||
@@ -60,7 +63,7 @@ public class MessageSequenceComparatorTests {
|
||||
|
||||
@Test
|
||||
public void testEqualWithDefaultValues() {
|
||||
MessageSequenceComparator comparator = new MessageSequenceComparator();
|
||||
Comparator<Message<?>> comparator = new MessageSequenceComparator();
|
||||
Message<String> message1 = MessageBuilder.withPayload("test1").build();
|
||||
Message<String> message2 = MessageBuilder.withPayload("test2").build();
|
||||
assertEquals(0, comparator.compare(message1, message2));
|
||||
|
||||
@@ -155,11 +155,11 @@ public class ResequencerTests {
|
||||
Message<?> reply2 = replyChannel.receive(0);
|
||||
Message<?> reply3 = replyChannel.receive(0);
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertNotNull(reply3);
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -180,18 +180,18 @@ public class ResequencerTests {
|
||||
Message<?> reply3 = replyChannel.receive(0);
|
||||
// only messages 1 and 2 should have been received by now
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertNull(reply3);
|
||||
// when sending the last message, the whole sequence must have been sent
|
||||
this.resequencer.handleMessage(message4);
|
||||
reply3 = replyChannel.receive(0);
|
||||
Message<?> reply4 = replyChannel.receive(0);
|
||||
assertNotNull(reply3);
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertNotNull(reply4);
|
||||
assertEquals(new Integer(4), new IntegrationMessageHeaderAccessor(reply4).getSequenceNumber());
|
||||
assertEquals(4, new IntegrationMessageHeaderAccessor(reply4).getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -230,18 +230,18 @@ public class ResequencerTests {
|
||||
Message<?> reply3 = replyChannel.receive(0);
|
||||
// only messages 1 and 2 should have been received by now
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertNull(reply3);
|
||||
// when sending the last message, the whole sequence must have been sent
|
||||
this.resequencer.handleMessage(message4);
|
||||
reply3 = replyChannel.receive(0);
|
||||
Message<?> reply4 = replyChannel.receive(0);
|
||||
assertNotNull(reply3);
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertNotNull(reply4);
|
||||
assertEquals(new Integer(4), new IntegrationMessageHeaderAccessor(reply4).getSequenceNumber());
|
||||
assertEquals(4, new IntegrationMessageHeaderAccessor(reply4).getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -288,7 +288,7 @@ public class ResequencerTests {
|
||||
Message<?> discard2 = discardChannel.receive(0);
|
||||
// message2 has been discarded because it came in with the wrong sequence size
|
||||
assertNotNull(discard1);
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(discard1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(discard1).getSequenceNumber());
|
||||
assertNull(discard2);
|
||||
}
|
||||
|
||||
@@ -329,13 +329,13 @@ public class ResequencerTests {
|
||||
reply3 = replyChannel.receive(0);
|
||||
Message<?> reply4 = replyChannel.receive(0);
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertNotNull(reply3);
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(reply3).getSequenceNumber());
|
||||
assertNotNull(reply4);
|
||||
assertEquals(new Integer(4), new IntegrationMessageHeaderAccessor(reply4).getSequenceNumber());
|
||||
assertEquals(4, new IntegrationMessageHeaderAccessor(reply4).getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -82,7 +82,7 @@ public class ResequencerIntegrationTests {
|
||||
inputChannel.send(message1);
|
||||
message1 = outputChannel.receive(0);
|
||||
assertNotNull(message1);
|
||||
assertEquals((Integer) 1, new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertFalse(message1.getHeaders().containsKey("foo"));
|
||||
|
||||
inputChannel.send(message2);
|
||||
@@ -90,9 +90,9 @@ public class ResequencerIntegrationTests {
|
||||
message3 = outputChannel.receive(0);
|
||||
assertNotNull(message2);
|
||||
assertNotNull(message3);
|
||||
assertEquals((Integer) 2, new IntegrationMessageHeaderAccessor(message2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(message2).getSequenceNumber());
|
||||
assertTrue(message2.getHeaders().containsKey("foo"));
|
||||
assertEquals((Integer) 3, new IntegrationMessageHeaderAccessor(message3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(message3).getSequenceNumber());
|
||||
assertFalse(message3.getHeaders().containsKey("foo"));
|
||||
|
||||
inputChannel.send(message5);
|
||||
@@ -108,11 +108,11 @@ public class ResequencerIntegrationTests {
|
||||
assertNotNull(message4);
|
||||
assertNotNull(message5);
|
||||
assertNotNull(message6);
|
||||
assertEquals((Integer) 4, new IntegrationMessageHeaderAccessor(message4).getSequenceNumber());
|
||||
assertEquals(4, new IntegrationMessageHeaderAccessor(message4).getSequenceNumber());
|
||||
assertTrue(message4.getHeaders().containsKey("foo"));
|
||||
assertEquals((Integer) 5, new IntegrationMessageHeaderAccessor(message5).getSequenceNumber());
|
||||
assertEquals(5, new IntegrationMessageHeaderAccessor(message5).getSequenceNumber());
|
||||
assertFalse(message5.getHeaders().containsKey("foo"));
|
||||
assertEquals((Integer) 6, new IntegrationMessageHeaderAccessor(message6).getSequenceNumber());
|
||||
assertEquals(6, new IntegrationMessageHeaderAccessor(message6).getSequenceNumber());
|
||||
assertFalse(message6.getHeaders().containsKey("foo"));
|
||||
|
||||
assertEquals(0, store.getMessageGroup("A").getMessages().size());
|
||||
@@ -148,7 +148,7 @@ public class ResequencerIntegrationTests {
|
||||
inputChannel.send(message1);
|
||||
message1 = outputChannel.receive(0);
|
||||
assertNotNull(message1);
|
||||
assertEquals((Integer) 1, new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -11,7 +11,7 @@
|
||||
</channel>
|
||||
|
||||
<beans:bean id="messageStore" class="org.springframework.integration.store.SimpleMessageStore" />
|
||||
|
||||
|
||||
<beans:bean id="comparator" class="org.springframework.integration.aggregator.MessageSequenceComparator" />
|
||||
|
||||
</beans:beans>
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -35,6 +35,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
/**
|
||||
* @author Dave Syer
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@ContextConfiguration
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@@ -65,11 +66,11 @@ public class ResequencerWithMessageStoreParserTests {
|
||||
Message<?> message3 = output.receive(500);
|
||||
|
||||
assertNotNull(message1);
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertNotNull(message2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(message2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(message2).getSequenceNumber());
|
||||
assertNotNull(message3);
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(message3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(message3).getSequenceNumber());
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -22,14 +22,15 @@ import java.util.List;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.aggregator.MessageSequenceComparator;
|
||||
import org.springframework.integration.annotation.Aggregator;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
public class TestAggregatorBean {
|
||||
|
||||
@@ -38,7 +39,7 @@ public class TestAggregatorBean {
|
||||
|
||||
@Aggregator
|
||||
public Message<?> createSingleMessageFromGroup(List<Message<?>> messages) {
|
||||
List<Message<?>> sortableList = new ArrayList<Message<?>>(messages);
|
||||
List<Message<?>> sortableList = new ArrayList<>(messages);
|
||||
Collections.sort(sortableList, new MessageSequenceComparator());
|
||||
StringBuffer buffer = new StringBuffer();
|
||||
Object correlationId = null;
|
||||
@@ -48,15 +49,15 @@ public class TestAggregatorBean {
|
||||
correlationId = new IntegrationMessageHeaderAccessor(message).getCorrelationId();
|
||||
}
|
||||
}
|
||||
Message<?> returnedMessage = new GenericMessage<String>(buffer.toString());
|
||||
Message<?> returnedMessage = new GenericMessage<>(buffer.toString());
|
||||
if (correlationId != null) {
|
||||
aggregatedMessages.put(correlationId, returnedMessage);
|
||||
this.aggregatedMessages.put(correlationId, returnedMessage);
|
||||
}
|
||||
return returnedMessage;
|
||||
}
|
||||
|
||||
public ConcurrentMap<Object, Message<?>> getAggregatedMessages() {
|
||||
return aggregatedMessages;
|
||||
return this.aggregatedMessages;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -22,15 +22,16 @@ import java.util.List;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.aggregator.MessageSequenceComparator;
|
||||
import org.springframework.integration.annotation.Aggregator;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@Component("endpointWithCustomizedAnnotation")
|
||||
public class TestAnnotatedEndpointWithCustomizedAggregator {
|
||||
@@ -44,7 +45,7 @@ public class TestAnnotatedEndpointWithCustomizedAggregator {
|
||||
sendPartialResultsOnExpiry = "true",
|
||||
sendTimeout = "98765432")
|
||||
public Message<?> aggregatingMethod(List<Message<?>> messages) {
|
||||
List<Message<?>> sortableList = new ArrayList<Message<?>>(messages);
|
||||
List<Message<?>> sortableList = new ArrayList<>(messages);
|
||||
Collections.sort(sortableList, new MessageSequenceComparator());
|
||||
StringBuffer buffer = new StringBuffer();
|
||||
Object correlationId = null;
|
||||
@@ -54,13 +55,13 @@ public class TestAnnotatedEndpointWithCustomizedAggregator {
|
||||
correlationId = new IntegrationMessageHeaderAccessor(message).getCorrelationId();
|
||||
}
|
||||
}
|
||||
Message<?> returnedMessage = new GenericMessage<String>(buffer.toString());
|
||||
Message<?> returnedMessage = new GenericMessage<>(buffer.toString());
|
||||
aggregatedMessages.put(correlationId, returnedMessage);
|
||||
return returnedMessage;
|
||||
}
|
||||
|
||||
public ConcurrentMap<Object, Message<?>> getAggregatedMessages() {
|
||||
return aggregatedMessages;
|
||||
return this.aggregatedMessages;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -22,15 +22,16 @@ import java.util.List;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.aggregator.MessageSequenceComparator;
|
||||
import org.springframework.integration.annotation.Aggregator;
|
||||
import org.springframework.integration.annotation.MessageEndpoint;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@MessageEndpoint("endpointWithDefaultAnnotation")
|
||||
public class TestAnnotatedEndpointWithDefaultAggregator {
|
||||
@@ -39,7 +40,7 @@ public class TestAnnotatedEndpointWithDefaultAggregator {
|
||||
|
||||
@Aggregator(inputChannel = "inputChannel")
|
||||
public Message<?> aggregatingMethod(List<Message<?>> messages) {
|
||||
List<Message<?>> sortableList = new ArrayList<Message<?>>(messages);
|
||||
List<Message<?>> sortableList = new ArrayList<>(messages);
|
||||
Collections.sort(sortableList, new MessageSequenceComparator());
|
||||
StringBuffer buffer = new StringBuffer();
|
||||
Object correlationId = null;
|
||||
@@ -49,13 +50,13 @@ public class TestAnnotatedEndpointWithDefaultAggregator {
|
||||
correlationId = new IntegrationMessageHeaderAccessor(message).getCorrelationId();
|
||||
}
|
||||
}
|
||||
Message<?> returnedMessage = new GenericMessage<String>(buffer.toString());
|
||||
aggregatedMessages.put(correlationId, returnedMessage);
|
||||
Message<?> returnedMessage = new GenericMessage<>(buffer.toString());
|
||||
this.aggregatedMessages.put(correlationId, returnedMessage);
|
||||
return returnedMessage;
|
||||
}
|
||||
|
||||
public ConcurrentMap<Object, Message<?>> getAggregatedMessages() {
|
||||
return aggregatedMessages;
|
||||
return this.aggregatedMessages;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2010 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -22,16 +22,17 @@ import java.util.List;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.aggregator.MessageSequenceComparator;
|
||||
import org.springframework.integration.annotation.Aggregator;
|
||||
import org.springframework.integration.annotation.MessageEndpoint;
|
||||
import org.springframework.integration.annotation.ReleaseStrategy;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@MessageEndpoint("endpointWithDefaultAnnotationAndCustomReleaseStrategy")
|
||||
public class TestAnnotatedEndpointWithReleaseStrategy {
|
||||
@@ -40,7 +41,7 @@ public class TestAnnotatedEndpointWithReleaseStrategy {
|
||||
|
||||
@Aggregator(inputChannel = "inputChannel")
|
||||
public Message<?> aggregatingMethod(List<Message<?>> messages) {
|
||||
List<Message<?>> sortableList = new ArrayList<Message<?>>(messages);
|
||||
List<Message<?>> sortableList = new ArrayList<>(messages);
|
||||
Collections.sort(sortableList, new MessageSequenceComparator());
|
||||
StringBuffer buffer = new StringBuffer();
|
||||
Object correlationId = null;
|
||||
@@ -50,8 +51,8 @@ public class TestAnnotatedEndpointWithReleaseStrategy {
|
||||
correlationId = new IntegrationMessageHeaderAccessor(message).getCorrelationId();
|
||||
}
|
||||
}
|
||||
Message<?> returnedMessage = new GenericMessage<String>(buffer.toString());
|
||||
aggregatedMessages.put(correlationId, returnedMessage);
|
||||
Message<?> returnedMessage = new GenericMessage<>(buffer.toString());
|
||||
this.aggregatedMessages.put(correlationId, returnedMessage);
|
||||
return returnedMessage;
|
||||
}
|
||||
|
||||
@@ -61,7 +62,7 @@ public class TestAnnotatedEndpointWithReleaseStrategy {
|
||||
}
|
||||
|
||||
public ConcurrentMap<Object, Message<?>> getAggregatedMessages() {
|
||||
return aggregatedMessages;
|
||||
return this.aggregatedMessages;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -210,10 +210,10 @@ public class BroadcastingDispatcherTests {
|
||||
dispatcher.addHandler(target2);
|
||||
dispatcher.dispatch(new GenericMessage<String>("test"));
|
||||
assertEquals(2, messages.size());
|
||||
assertEquals(0, (int) new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceNumber());
|
||||
assertEquals(0, (int) new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceSize());
|
||||
assertEquals(0, (int) new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceNumber());
|
||||
assertEquals(0, (int) new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceSize());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceSize());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceSize());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -231,14 +231,14 @@ public class BroadcastingDispatcherTests {
|
||||
Object originalId = inputMessage.getHeaders().getId();
|
||||
dispatcher.dispatch(inputMessage);
|
||||
assertEquals(3, messages.size());
|
||||
assertEquals(1, (int) new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceNumber());
|
||||
assertEquals(3, (int) new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceSize());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(messages.get(0)).getSequenceSize());
|
||||
assertEquals(originalId, new IntegrationMessageHeaderAccessor(messages.get(0)).getCorrelationId());
|
||||
assertEquals(2, (int) new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceNumber());
|
||||
assertEquals(3, (int) new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceSize());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(messages.get(1)).getSequenceSize());
|
||||
assertEquals(originalId, new IntegrationMessageHeaderAccessor(messages.get(1)).getCorrelationId());
|
||||
assertEquals(3, (int) new IntegrationMessageHeaderAccessor(messages.get(2)).getSequenceNumber());
|
||||
assertEquals(3, (int) new IntegrationMessageHeaderAccessor(messages.get(2)).getSequenceSize());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(messages.get(2)).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(messages.get(2)).getSequenceSize());
|
||||
assertEquals(originalId, new IntegrationMessageHeaderAccessor(messages.get(2)).getCorrelationId());
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2008 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -34,15 +34,15 @@ public class GenericMessageTests {
|
||||
@Test
|
||||
public void testMessageHeadersCopiedFromMap() {
|
||||
Map<String, Object> headerMap = new HashMap<String, Object>();
|
||||
headerMap.put("testAttribute", new Integer(123));
|
||||
headerMap.put("testAttribute", Integer.valueOf(123));
|
||||
headerMap.put("testProperty", "foo");
|
||||
headerMap.put(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE, 42);
|
||||
headerMap.put(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, 24);
|
||||
GenericMessage<String> message = new GenericMessage<String>("test", headerMap);
|
||||
assertEquals(new Integer(123), message.getHeaders().get("testAttribute"));
|
||||
assertEquals(123, message.getHeaders().get("testAttribute"));
|
||||
assertEquals("foo", message.getHeaders().get("testProperty", String.class));
|
||||
assertEquals(new Integer(42), new IntegrationMessageHeaderAccessor(message).getSequenceSize());
|
||||
assertEquals(new Integer(24), new IntegrationMessageHeaderAccessor(message).getSequenceNumber());
|
||||
assertEquals(42, new IntegrationMessageHeaderAccessor(message).getSequenceSize());
|
||||
assertEquals(24, new IntegrationMessageHeaderAccessor(message).getSequenceNumber());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -47,6 +47,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
@ContextConfiguration
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@@ -196,7 +197,7 @@ public class MessageBuilderTests {
|
||||
@Test
|
||||
public void testNonDestructiveSet() {
|
||||
Message<Integer> message1 = MessageBuilder.withPayload(1)
|
||||
.setPriority(42).build();
|
||||
.setPriority(42).build();
|
||||
Message<Integer> message2 = MessageBuilder.fromMessage(message1)
|
||||
.setHeaderIfAbsent(IntegrationMessageHeaderAccessor.PRIORITY, 13)
|
||||
.build();
|
||||
@@ -222,20 +223,20 @@ public class MessageBuilderTests {
|
||||
@Test
|
||||
public void testRemove() {
|
||||
Message<Integer> message1 = MessageBuilder.withPayload(1)
|
||||
.setHeader("foo", "bar").build();
|
||||
.setHeader("foo", "bar").build();
|
||||
Message<Integer> message2 = MessageBuilder.fromMessage(message1)
|
||||
.removeHeader("foo")
|
||||
.build();
|
||||
.removeHeader("foo")
|
||||
.build();
|
||||
assertFalse(message2.getHeaders().containsKey("foo"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSettingToNullRemoves() {
|
||||
Message<Integer> message1 = MessageBuilder.withPayload(1)
|
||||
.setHeader("foo", "bar").build();
|
||||
.setHeader("foo", "bar").build();
|
||||
Message<Integer> message2 = MessageBuilder.fromMessage(message1)
|
||||
.setHeader("foo", null)
|
||||
.build();
|
||||
.setHeader("foo", null)
|
||||
.build();
|
||||
assertFalse(message2.getHeaders().containsKey("foo"));
|
||||
}
|
||||
|
||||
@@ -351,4 +352,13 @@ public class MessageBuilderTests {
|
||||
assertEquals(original, result);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSequenceNumberAsLong() {
|
||||
Message<String> message = MessageBuilder.withPayload("foo")
|
||||
.setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, Long.MAX_VALUE)
|
||||
.build();
|
||||
|
||||
Integer sequenceNumber = new IntegrationMessageHeaderAccessor(message).getSequenceNumber();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -310,12 +310,12 @@ public class RecipientListRouterTests {
|
||||
assertNotNull(result1a);
|
||||
assertNotNull(result1b);
|
||||
assertEquals("test", result1a.getPayload());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1a).getSequenceNumber().intValue());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1a).getSequenceSize().intValue());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1a).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1a).getSequenceSize());
|
||||
assertNull(new IntegrationMessageHeaderAccessor(result1a).getCorrelationId());
|
||||
assertEquals("test", result1b.getPayload());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1b).getSequenceNumber().intValue());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1b).getSequenceSize().intValue());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1b).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(result1b).getSequenceSize());
|
||||
assertNull(new IntegrationMessageHeaderAccessor(result1b).getCorrelationId());
|
||||
}
|
||||
|
||||
@@ -338,12 +338,12 @@ public class RecipientListRouterTests {
|
||||
assertNotNull(result1a);
|
||||
assertNotNull(result1b);
|
||||
assertEquals("test", result1a.getPayload());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(result1a).getSequenceNumber().intValue());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(result1a).getSequenceSize().intValue());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(result1a).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(result1a).getSequenceSize());
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(result1a).getCorrelationId());
|
||||
assertEquals("test", result1b.getPayload());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(result1b).getSequenceNumber().intValue());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(result1b).getSequenceSize().intValue());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(result1b).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(result1b).getSequenceSize());
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(result1b).getCorrelationId());
|
||||
}
|
||||
|
||||
@@ -453,6 +453,7 @@ public class RecipientListRouterTests {
|
||||
public boolean accept(Message<?> message) {
|
||||
return true;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -466,6 +467,7 @@ public class RecipientListRouterTests {
|
||||
public boolean accept(Message<?> message) {
|
||||
return false;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -192,12 +192,12 @@ public class RouterParserTests {
|
||||
assertEquals(originalMessage.getHeaders().getId(), new IntegrationMessageHeaderAccessor(message1).getCorrelationId());
|
||||
assertEquals(originalMessage.getHeaders().getId(), new IntegrationMessageHeaderAccessor(message2).getCorrelationId());
|
||||
assertEquals(originalMessage.getHeaders().getId(), new IntegrationMessageHeaderAccessor(message3).getCorrelationId());
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(message1).getSequenceSize());
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(message2).getSequenceNumber());
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(message2).getSequenceSize());
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(message3).getSequenceNumber());
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(message3).getSequenceSize());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(message1).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(message1).getSequenceSize());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(message2).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(message2).getSequenceSize());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(message3).getSequenceNumber());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(message3).getSequenceSize());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -355,13 +355,13 @@ public class MethodInvokingSplitterTests {
|
||||
List<Message<?>> replies = replyChannel.clear();
|
||||
Message<?> reply1 = replies.get(0);
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply1).getSequenceSize());
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply1).getSequenceSize());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(reply1).getCorrelationId());
|
||||
Message<?> reply2 = replies.get(1);
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceSize());
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceSize());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(reply2).getCorrelationId());
|
||||
}
|
||||
|
||||
@@ -375,13 +375,13 @@ public class MethodInvokingSplitterTests {
|
||||
List<Message<?>> replies = replyChannel.clear();
|
||||
Message<?> reply1 = replies.get(0);
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply1).getSequenceSize());
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply1).getSequenceSize());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber());
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(reply1).getCorrelationId());
|
||||
Message<?> reply2 = replies.get(1);
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceSize());
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceSize());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber());
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(reply2).getCorrelationId());
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 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.
|
||||
@@ -119,17 +119,17 @@ public class SpelSplitterIntegrationTests {
|
||||
Message<?> c = output.receive(0);
|
||||
Message<?> d = output.receive(0);
|
||||
assertEquals("a", a.getPayload());
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(a).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(a).getSequenceSize());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(a).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(a).getSequenceSize());
|
||||
assertEquals("b", b.getPayload());
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(b).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(b).getSequenceSize());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(b).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(b).getSequenceSize());
|
||||
assertEquals("c", c.getPayload());
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(c).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(c).getSequenceSize());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(c).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(c).getSequenceSize());
|
||||
assertEquals("d", d.getPayload());
|
||||
assertEquals(new Integer(4), new IntegrationMessageHeaderAccessor(d).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(d).getSequenceSize());
|
||||
assertEquals(4, new IntegrationMessageHeaderAccessor(d).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(d).getSequenceSize());
|
||||
assertNull(output.receive(0));
|
||||
}
|
||||
|
||||
@@ -141,17 +141,17 @@ public class SpelSplitterIntegrationTests {
|
||||
Message<?> c = output.receive(0);
|
||||
Message<?> d = output.receive(0);
|
||||
assertEquals("a", a.getPayload());
|
||||
assertEquals(new Integer(1), new IntegrationMessageHeaderAccessor(a).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(a).getSequenceSize());
|
||||
assertEquals(1, new IntegrationMessageHeaderAccessor(a).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(a).getSequenceSize());
|
||||
assertEquals("b", b.getPayload());
|
||||
assertEquals(new Integer(2), new IntegrationMessageHeaderAccessor(b).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(b).getSequenceSize());
|
||||
assertEquals(2, new IntegrationMessageHeaderAccessor(b).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(b).getSequenceSize());
|
||||
assertEquals("c", c.getPayload());
|
||||
assertEquals(new Integer(3), new IntegrationMessageHeaderAccessor(c).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(c).getSequenceSize());
|
||||
assertEquals(3, new IntegrationMessageHeaderAccessor(c).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(c).getSequenceSize());
|
||||
assertEquals("d", d.getPayload());
|
||||
assertEquals(new Integer(4), new IntegrationMessageHeaderAccessor(d).getSequenceNumber());
|
||||
assertEquals(new Integer(0), new IntegrationMessageHeaderAccessor(d).getSequenceSize());
|
||||
assertEquals(4, new IntegrationMessageHeaderAccessor(d).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(d).getSequenceSize());
|
||||
assertNull(output.receive(0));
|
||||
}
|
||||
|
||||
@@ -177,6 +177,7 @@ public class SpelSplitterIntegrationTests {
|
||||
public Iterator<String> splitIterator(String s) {
|
||||
return Arrays.asList(s.split(",")).iterator();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -58,6 +58,8 @@ import org.springframework.util.MimeType;
|
||||
|
||||
/**
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 2.0
|
||||
*
|
||||
*/
|
||||
@@ -222,7 +224,7 @@ public class TcpMessageMapperTests {
|
||||
.getHeaders().get(IpHeaders.IP_ADDRESS));
|
||||
assertEquals(1234, message
|
||||
.getHeaders().get(IpHeaders.REMOTE_PORT));
|
||||
assertEquals(Integer.valueOf(0), new IntegrationMessageHeaderAccessor(message).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(message).getSequenceNumber());
|
||||
message = mapper.toMessage(connection);
|
||||
assertEquals(TEST_PAYLOAD, new String((byte[]) message.getPayload()));
|
||||
assertEquals("MyHost", message
|
||||
@@ -231,7 +233,7 @@ public class TcpMessageMapperTests {
|
||||
.getHeaders().get(IpHeaders.IP_ADDRESS));
|
||||
assertEquals(1234, message
|
||||
.getHeaders().get(IpHeaders.REMOTE_PORT));
|
||||
assertEquals(Integer.valueOf(0), new IntegrationMessageHeaderAccessor(message).getSequenceNumber());
|
||||
assertEquals(0, new IntegrationMessageHeaderAccessor(message).getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -306,7 +308,7 @@ public class TcpMessageMapperTests {
|
||||
assertEquals(1234, message
|
||||
.getHeaders().get(IpHeaders.REMOTE_PORT));
|
||||
IntegrationMessageHeaderAccessor headerAccessor = new IntegrationMessageHeaderAccessor(message);
|
||||
assertEquals(Integer.valueOf(1), headerAccessor.getSequenceNumber());
|
||||
assertEquals(1, headerAccessor.getSequenceNumber());
|
||||
assertEquals(message.getHeaders().get(IpHeaders.CONNECTION_ID), headerAccessor.getCorrelationId());
|
||||
message = mapper.toMessage(connection);
|
||||
headerAccessor = new IntegrationMessageHeaderAccessor(message);
|
||||
@@ -317,7 +319,7 @@ public class TcpMessageMapperTests {
|
||||
.getHeaders().get(IpHeaders.IP_ADDRESS));
|
||||
assertEquals(1234, message
|
||||
.getHeaders().get(IpHeaders.REMOTE_PORT));
|
||||
assertEquals(Integer.valueOf(2), headerAccessor.getSequenceNumber());
|
||||
assertEquals(2, headerAccessor.getSequenceNumber());
|
||||
assertEquals(message.getHeaders().get(IpHeaders.CONNECTION_ID), headerAccessor.getCorrelationId());
|
||||
assertNotNull(message.getHeaders().get("foo"));
|
||||
assertEquals("bar", message.getHeaders().get("foo"));
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2014-2016 the original author or authors.
|
||||
* Copyright 2014-2017 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.
|
||||
@@ -81,6 +81,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
/**
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 4.0
|
||||
*/
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@@ -229,12 +230,12 @@ public class ChannelSecurityInterceptorSecuredChannelAnnotationTests {
|
||||
Message<?> receive = this.securedChannelQueue.receive(10000);
|
||||
assertNotNull(receive);
|
||||
IntegrationMessageHeaderAccessor headerAccessor = new IntegrationMessageHeaderAccessor(receive);
|
||||
assertEquals(new Integer(0), headerAccessor.getSequenceNumber());
|
||||
assertEquals(0, headerAccessor.getSequenceNumber());
|
||||
|
||||
receive = this.securedChannelQueue2.receive(10000);
|
||||
assertNotNull(receive);
|
||||
headerAccessor = new IntegrationMessageHeaderAccessor(receive);
|
||||
assertEquals(new Integer(0), headerAccessor.getSequenceNumber());
|
||||
assertEquals(0, headerAccessor.getSequenceNumber());
|
||||
|
||||
this.publishSubscribeChannel.setApplySequence(true);
|
||||
|
||||
@@ -243,12 +244,12 @@ public class ChannelSecurityInterceptorSecuredChannelAnnotationTests {
|
||||
receive = this.securedChannelQueue.receive(10000);
|
||||
assertNotNull(receive);
|
||||
headerAccessor = new IntegrationMessageHeaderAccessor(receive);
|
||||
assertEquals(new Integer(1), headerAccessor.getSequenceNumber());
|
||||
assertEquals(1, headerAccessor.getSequenceNumber());
|
||||
|
||||
receive = this.securedChannelQueue2.receive(10000);
|
||||
assertNotNull(receive);
|
||||
headerAccessor = new IntegrationMessageHeaderAccessor(receive);
|
||||
assertEquals(new Integer(2), headerAccessor.getSequenceNumber());
|
||||
assertEquals(2, headerAccessor.getSequenceNumber());
|
||||
|
||||
this.publishSubscribeChannel.setApplySequence(false);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user