Moved the ID from Message to MessageHeaders.
This commit is contained in:
@@ -56,7 +56,7 @@ public class MessageStoringInterceptor extends ChannelInterceptorAdapter {
|
||||
@Override
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (message != null) {
|
||||
this.messageStore.put(message.getId(), message);
|
||||
this.messageStore.put(message.getHeaders().getId(), message);
|
||||
}
|
||||
return message;
|
||||
}
|
||||
@@ -86,7 +86,7 @@ public class MessageStoringInterceptor extends ChannelInterceptorAdapter {
|
||||
@Override
|
||||
public Message<?> postReceive(Message<?> message, MessageChannel channel) {
|
||||
if (message != null) {
|
||||
this.messageStore.remove(message.getId());
|
||||
this.messageStore.remove(message.getHeaders().getId());
|
||||
}
|
||||
return message;
|
||||
}
|
||||
|
||||
@@ -101,7 +101,7 @@ public class WireTap extends ChannelInterceptorAdapter implements Lifecycle {
|
||||
public Message<?> preSend(Message<?> message, MessageChannel channel) {
|
||||
if (this.running && this.selectorsAccept(message)) {
|
||||
Message<?> duplicate = MessageBuilder.fromMessage(message)
|
||||
.setHeader(ORIGINAL_MESSAGE_ID_KEY, message.getId())
|
||||
.setHeader(ORIGINAL_MESSAGE_ID_KEY, message.getHeaders().getId())
|
||||
.build();
|
||||
if (!this.secondaryChannel.send(duplicate, 0)) {
|
||||
if (logger.isWarnEnabled()) {
|
||||
|
||||
@@ -117,7 +117,7 @@ public class HandlerEndpoint extends AbstractEndpoint {
|
||||
Object correlationId = replyMessage.getHeaders().getCorrelationId();
|
||||
if (correlationId == null) {
|
||||
replyMessage = MessageBuilder.fromMessage(replyMessage)
|
||||
.setHeader(MessageHeaders.CORRELATION_ID, message.getId()).build();
|
||||
.setHeader(MessageHeaders.CORRELATION_ID, message.getHeaders().getId()).build();
|
||||
}
|
||||
if (replyMessage != null) {
|
||||
MessageTarget returnAddress = resolveReturnAddress(replyMessage);
|
||||
|
||||
@@ -192,8 +192,8 @@ public class SimpleMessagingGateway extends MessagingGatewaySupport implements M
|
||||
message = MessageBuilder.fromMessage(message).setReturnAddress(this.replyChannel).build();
|
||||
this.send(message);
|
||||
return (this.replyTimeout >= 0)
|
||||
? this.replyMessageCorrelator.getReply(message.getId(), this.replyTimeout)
|
||||
: this.replyMessageCorrelator.getReply(message.getId());
|
||||
? this.replyMessageCorrelator.getReply(message.getHeaders().getId(), this.replyTimeout)
|
||||
: this.replyMessageCorrelator.getReply(message.getHeaders().getId());
|
||||
}
|
||||
|
||||
private void registerReplyMessageCorrelator() {
|
||||
|
||||
@@ -130,7 +130,7 @@ public abstract class AbstractMessageHandlerAdapter extends AbstractMethodInvoki
|
||||
return null;
|
||||
}
|
||||
return MessageBuilder.fromMessage(reply).copyHeadersIfAbsent(originalMessage.getHeaders())
|
||||
.setHeaderIfAbsent(MessageHeaders.CORRELATION_ID, originalMessage.getId()).build();
|
||||
.setHeaderIfAbsent(MessageHeaders.CORRELATION_ID, originalMessage.getHeaders().getId()).build();
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -65,28 +65,6 @@ public class AsyncMessage<T> implements Future<Message<T>>, Message<T> {
|
||||
return this.future.isDone();
|
||||
}
|
||||
|
||||
public MessageHeaders getHeaders() {
|
||||
try {
|
||||
return this.future.get().getHeaders();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return null;
|
||||
} catch (ExecutionException e) {
|
||||
throw new MessagingException("failure occurred in AsyncMessage", e);
|
||||
}
|
||||
}
|
||||
|
||||
public Object getId() {
|
||||
try {
|
||||
return this.future.get().getId();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return null;
|
||||
} catch (ExecutionException e) {
|
||||
throw new MessagingException("failure occurred in AsyncMessage", e);
|
||||
}
|
||||
}
|
||||
|
||||
public T getPayload() {
|
||||
try {
|
||||
return this.future.get().getPayload();
|
||||
@@ -98,4 +76,15 @@ public class AsyncMessage<T> implements Future<Message<T>>, Message<T> {
|
||||
}
|
||||
}
|
||||
|
||||
public MessageHeaders getHeaders() {
|
||||
try {
|
||||
return this.future.get().getHeaders();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return null;
|
||||
} catch (ExecutionException e) {
|
||||
throw new MessagingException("failure occurred in AsyncMessage", e);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -21,7 +21,6 @@ import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import org.springframework.integration.util.IdGenerator;
|
||||
import org.springframework.integration.util.RandomUuidGenerator;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
@@ -31,14 +30,10 @@ import org.springframework.util.Assert;
|
||||
*/
|
||||
public class GenericMessage<T> implements Message<T>, Serializable {
|
||||
|
||||
private final static String ID_HEADER_KEY = "id";
|
||||
|
||||
private volatile T payload;
|
||||
|
||||
private final MessageHeaders headers;
|
||||
|
||||
private transient final IdGenerator defaultIdGenerator = new RandomUuidGenerator();
|
||||
|
||||
|
||||
/**
|
||||
* Create a new message with the given payload. The id will be generated by
|
||||
@@ -67,25 +62,20 @@ public class GenericMessage<T> implements Message<T>, Serializable {
|
||||
else if (headers instanceof MessageHeaders) {
|
||||
headers = new HashMap<String, Object>(headers);
|
||||
}
|
||||
headers.put(ID_HEADER_KEY, this.defaultIdGenerator.generateId());
|
||||
this.headers = new MessageHeaders(headers);
|
||||
}
|
||||
|
||||
|
||||
public Object getId() {
|
||||
return this.headers.get(ID_HEADER_KEY);
|
||||
public T getPayload() {
|
||||
return this.payload;
|
||||
}
|
||||
|
||||
public MessageHeaders getHeaders() {
|
||||
return this.headers;
|
||||
}
|
||||
|
||||
public T getPayload() {
|
||||
return this.payload;
|
||||
}
|
||||
|
||||
public String toString() {
|
||||
return "[ID=" + this.getId() + "][Headers=" + this.headers + "][Payload='" + this.payload + "']";
|
||||
return "[Payload=" + this.payload + "][Headers=" + this.headers + "]";
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -24,8 +24,6 @@ package org.springframework.integration.message;
|
||||
*/
|
||||
public interface Message<T> {
|
||||
|
||||
Object getId();
|
||||
|
||||
T getPayload();
|
||||
|
||||
MessageHeaders getHeaders();
|
||||
|
||||
@@ -24,6 +24,9 @@ import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import org.springframework.integration.util.IdGenerator;
|
||||
import org.springframework.integration.util.RandomUuidGenerator;
|
||||
|
||||
/**
|
||||
* The headers for a {@link Message}.
|
||||
*
|
||||
@@ -32,6 +35,8 @@ import java.util.Set;
|
||||
*/
|
||||
public final class MessageHeaders implements Map<String, Object>, Serializable {
|
||||
|
||||
public static final String ID = "internal.header.id";
|
||||
|
||||
public static final String TIMESTAMP = "internal.header.timestamp";
|
||||
|
||||
public static final String CORRELATION_ID = "internal.header.correlationId";
|
||||
@@ -51,13 +56,20 @@ public final class MessageHeaders implements Map<String, Object>, Serializable {
|
||||
|
||||
private final Map<String, Object> headers;
|
||||
|
||||
private transient final IdGenerator idGenerator = new RandomUuidGenerator();
|
||||
|
||||
|
||||
public MessageHeaders(Map<String, Object> headers) {
|
||||
this.headers = (headers != null ? headers : new HashMap<String, Object>());
|
||||
this.headers.put(ID, this.idGenerator.generateId());
|
||||
this.headers.put(TIMESTAMP, new Long(System.currentTimeMillis()));
|
||||
}
|
||||
|
||||
|
||||
public Object getId() {
|
||||
return this.get(ID);
|
||||
}
|
||||
|
||||
public Long getTimestamp() {
|
||||
return this.get(TIMESTAMP, Long.class);
|
||||
}
|
||||
@@ -169,7 +181,7 @@ public final class MessageHeaders implements Map<String, Object>, Serializable {
|
||||
}
|
||||
|
||||
public String toString() {
|
||||
return headers.toString();
|
||||
return this.headers.toString();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -85,7 +85,7 @@ public class SplitterMessageHandlerAdapter extends AbstractMessageHandlerAdapter
|
||||
Message<?> splitMessage = (item instanceof Message<?>) ?
|
||||
(Message<?>) item : this.createReplyMessage(item, originalMessage);
|
||||
splitMessage = MessageBuilder.fromMessage(splitMessage)
|
||||
.setCorrelationId(originalMessage.getId())
|
||||
.setCorrelationId(originalMessage.getHeaders().getId())
|
||||
.setSequenceNumber(++sequenceNumber)
|
||||
.setSequenceSize(sequenceSize).build();
|
||||
this.sendMessage(splitMessage, this.outputChannelName);
|
||||
@@ -99,7 +99,7 @@ public class SplitterMessageHandlerAdapter extends AbstractMessageHandlerAdapter
|
||||
Message<?> splitMessage = (item instanceof Message<?>) ?
|
||||
(Message<?>) item : this.createReplyMessage(item, originalMessage);
|
||||
splitMessage = MessageBuilder.fromMessage(splitMessage)
|
||||
.setCorrelationId(originalMessage.getId())
|
||||
.setCorrelationId(originalMessage.getHeaders().getId())
|
||||
.setSequenceNumber(++sequenceNumber)
|
||||
.setSequenceSize(sequenceSize).build();
|
||||
this.sendMessage(splitMessage, this.outputChannelName);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2007 the original author or authors.
|
||||
* Copyright 2002-2008 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.
|
||||
|
||||
Reference in New Issue
Block a user