INT-3953: AMQP Delay Header Mapping

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

Polishing

Doc Polish
This commit is contained in:
Gary Russell
2016-02-12 13:52:41 -05:00
committed by Artem Bilan
parent ae2e1b0a8d
commit 51a16b2f45
6 changed files with 39 additions and 3 deletions

View File

@@ -70,11 +70,13 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
STANDARD_HEADER_NAMES.add(AmqpHeaders.CONTENT_LENGTH);
STANDARD_HEADER_NAMES.add(AmqpHeaders.CONTENT_TYPE);
STANDARD_HEADER_NAMES.add(AmqpHeaders.CORRELATION_ID);
STANDARD_HEADER_NAMES.add(AmqpHeaders.DELAY);
STANDARD_HEADER_NAMES.add(AmqpHeaders.DELIVERY_MODE);
STANDARD_HEADER_NAMES.add(AmqpHeaders.DELIVERY_TAG);
STANDARD_HEADER_NAMES.add(AmqpHeaders.EXPIRATION);
STANDARD_HEADER_NAMES.add(AmqpHeaders.MESSAGE_COUNT);
STANDARD_HEADER_NAMES.add(AmqpHeaders.MESSAGE_ID);
STANDARD_HEADER_NAMES.add(AmqpHeaders.RECEIVED_DELAY);
STANDARD_HEADER_NAMES.add(AmqpHeaders.RECEIVED_EXCHANGE);
STANDARD_HEADER_NAMES.add(AmqpHeaders.RECEIVED_ROUTING_KEY);
STANDARD_HEADER_NAMES.add(AmqpHeaders.REDELIVERED);
@@ -169,6 +171,10 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
if (priority != null && priority > 0) {
headers.put(IntegrationMessageHeaderAccessor.PRIORITY, priority);
}
Integer receivedDelay = amqpMessageProperties.getReceivedDelay();
if (receivedDelay != null) {
headers.put(AmqpHeaders.RECEIVED_DELAY, receivedDelay);
}
String receivedExchange = amqpMessageProperties.getReceivedExchange();
if (StringUtils.hasText(receivedExchange)) {
headers.put(AmqpHeaders.RECEIVED_EXCHANGE, receivedExchange);
@@ -253,6 +259,10 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
if (correlationId instanceof byte[]) {
amqpMessageProperties.setCorrelationId((byte[]) correlationId);
}
Integer delay = getHeaderIfAvailable(headers, AmqpHeaders.DELAY, Integer.class);
if (delay != null) {
amqpMessageProperties.setDelay(delay);
}
MessageDeliveryMode deliveryMode = getHeaderIfAvailable(headers, AmqpHeaders.DELIVERY_MODE, MessageDeliveryMode.class);
if (deliveryMode != null) {
amqpMessageProperties.setDeliveryMode(deliveryMode);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2015 the original author or authors.
* 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.
@@ -570,6 +570,16 @@ public class StubRabbitConnectionFactory implements ConnectionFactory {
throws IOException {
}
@Override
public long messageCount(String queue) throws IOException {
return 0;
}
@Override
public long consumerCount(String queue) throws IOException {
return 0;
}
}
}

View File

@@ -60,6 +60,7 @@ public class DefaultAmqpHeaderMapperTests {
headerMap.put(AmqpHeaders.CONTENT_TYPE, "test.contentType");
byte[] testCorrelationId = new byte[] {1, 2, 3};
headerMap.put(AmqpHeaders.CORRELATION_ID, testCorrelationId);
headerMap.put(AmqpHeaders.DELAY, 1234);
headerMap.put(AmqpHeaders.DELIVERY_MODE, MessageDeliveryMode.NON_PERSISTENT);
headerMap.put(AmqpHeaders.DELIVERY_TAG, 1234L);
headerMap.put(AmqpHeaders.EXPIRATION, "test.expiration");
@@ -93,6 +94,7 @@ public class DefaultAmqpHeaderMapperTests {
assertEquals(99L, amqpProperties.getContentLength());
assertEquals("test.contentType", amqpProperties.getContentType());
assertEquals(testCorrelationId, amqpProperties.getCorrelationId());
assertEquals(Integer.valueOf(1234), amqpProperties.getDelay());
assertEquals(MessageDeliveryMode.NON_PERSISTENT, amqpProperties.getDeliveryMode());
assertEquals(1234L, amqpProperties.getDeliveryTag());
assertEquals("test.expiration", amqpProperties.getExpiration());
@@ -164,6 +166,7 @@ public class DefaultAmqpHeaderMapperTests {
amqpProperties.setMessageCount(42);
amqpProperties.setMessageId("test.messageId");
amqpProperties.setPriority(22);
amqpProperties.setReceivedDelay(4567);
amqpProperties.setReceivedExchange("test.receivedExchange");
amqpProperties.setReceivedRoutingKey("test.receivedRoutingKey");
amqpProperties.setRedelivered(true);
@@ -186,6 +189,7 @@ public class DefaultAmqpHeaderMapperTests {
assertEquals("test.expiration", headerMap.get(AmqpHeaders.EXPIRATION));
assertEquals(42, headerMap.get(AmqpHeaders.MESSAGE_COUNT));
assertEquals("test.messageId", headerMap.get(AmqpHeaders.MESSAGE_ID));
assertEquals(4567, headerMap.get(AmqpHeaders.RECEIVED_DELAY));
assertEquals("test.receivedExchange", headerMap.get(AmqpHeaders.RECEIVED_EXCHANGE));
assertEquals("test.receivedRoutingKey", headerMap.get(AmqpHeaders.RECEIVED_ROUTING_KEY));
assertEquals("test.replyTo", headerMap.get(AmqpHeaders.REPLY_TO));