INT-3337: MongoDbMessageStore Refactoring
JIRA: https://jira.spring.io/browse/INT-3337 INT-3337: Addressing PR comments INT-3337: Add support to store/read Messages * Add support to store/read `ErrorMessage`. As a trick for `Throwable` is selected (de)serializing converter INT-3337: Add converters for `Message<?>` impls INT-3337: Make `MDbMS.MessageWrapper` AuditAware Add `_id` persistence field Addressing PR comments Polishing - Docs - Check for null ApplicationContext - Add afterPropertiesSet() to tests
This commit is contained in:
committed by
Gary Russell
parent
6d6bef58b2
commit
c11e3ba3b5
@@ -23,6 +23,7 @@ import org.springframework.integration.store.SimpleMessageStore;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
|
||||
/**
|
||||
@@ -34,6 +35,7 @@ import org.springframework.util.ObjectUtils;
|
||||
* a reference to the message and changes will be reflected there too.
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
@@ -52,9 +54,10 @@ public class MutableMessage<T> implements Message<T>, Serializable {
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public MutableMessage(T payload, MessageHeaders headers) {
|
||||
this.payload = payload;
|
||||
public MutableMessage(T payload, Map<String, Object> headers) {
|
||||
Assert.notNull(payload, "payload must not be null");
|
||||
this.headers = new MessageHeaders(headers);
|
||||
this.payload = payload;
|
||||
// Needs SPR-11468 to avoid DFA and header manipulation
|
||||
rawHeaders = (Map<String, Object>) new DirectFieldAccessor(this.headers)
|
||||
.getPropertyValue("headers");
|
||||
@@ -64,6 +67,7 @@ public class MutableMessage<T> implements Message<T>, Serializable {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public MessageHeaders getHeaders() {
|
||||
return this.headers;
|
||||
@@ -75,6 +79,7 @@ public class MutableMessage<T> implements Message<T>, Serializable {
|
||||
}
|
||||
|
||||
public void setPayload(T payload) {
|
||||
Assert.notNull(payload, "'payload' must not be null");
|
||||
this.payload = payload;
|
||||
}
|
||||
|
||||
|
||||
@@ -362,6 +362,7 @@ public class ConfigurableMongoDbMessageStore extends AbstractMessageGroupStore
|
||||
* when the application context is configured with auditing. The document is not
|
||||
* currently Auditable.
|
||||
*/
|
||||
@SuppressWarnings("unused")
|
||||
@Id
|
||||
private String _id;
|
||||
|
||||
|
||||
@@ -29,10 +29,18 @@ import java.util.Map;
|
||||
import java.util.Properties;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.springframework.beans.BeansException;
|
||||
import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.beans.factory.BeanClassLoaderAware;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.context.ApplicationContext;
|
||||
import org.springframework.context.ApplicationContextAware;
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.core.serializer.support.DeserializingConverter;
|
||||
import org.springframework.core.serializer.support.SerializingConverter;
|
||||
import org.springframework.data.annotation.Id;
|
||||
import org.springframework.data.annotation.Transient;
|
||||
import org.springframework.data.convert.WritingConverter;
|
||||
import org.springframework.data.domain.Sort;
|
||||
import org.springframework.data.domain.Sort.Direction;
|
||||
import org.springframework.data.mapping.context.MappingContext;
|
||||
@@ -46,6 +54,8 @@ import org.springframework.data.mongodb.core.mapping.MongoPersistentProperty;
|
||||
import org.springframework.data.mongodb.core.query.Query;
|
||||
import org.springframework.data.mongodb.core.query.Update;
|
||||
import org.springframework.integration.history.MessageHistory;
|
||||
import org.springframework.integration.message.AdviceMessage;
|
||||
import org.springframework.integration.message.MutableMessage;
|
||||
import org.springframework.integration.store.AbstractMessageGroupStore;
|
||||
import org.springframework.integration.store.MessageGroup;
|
||||
import org.springframework.integration.store.MessageGroupStore;
|
||||
@@ -54,6 +64,7 @@ import org.springframework.integration.store.SimpleMessageGroup;
|
||||
import org.springframework.jmx.export.annotation.ManagedAttribute;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.support.ErrorMessage;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ClassUtils;
|
||||
@@ -76,7 +87,8 @@ import com.mongodb.DBObject;
|
||||
* @author Artem Bilan
|
||||
* @since 2.1
|
||||
*/
|
||||
public class MongoDbMessageStore extends AbstractMessageGroupStore implements MessageStore, BeanClassLoaderAware {
|
||||
public class MongoDbMessageStore extends AbstractMessageGroupStore
|
||||
implements MessageStore, BeanClassLoaderAware, ApplicationContextAware, InitializingBean {
|
||||
|
||||
private final static String DEFAULT_COLLECTION_NAME = "messages";
|
||||
|
||||
@@ -90,17 +102,19 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
|
||||
private final static String GROUP_UPDATE_TIMESTAMP_KEY = "_group_update_timestamp";
|
||||
|
||||
private final static String PAYLOAD_TYPE_KEY = "_payloadType";
|
||||
|
||||
private final static String CREATED_DATE = "_createdDate";
|
||||
|
||||
|
||||
private final MongoTemplate template;
|
||||
|
||||
private final MessageReadingMongoConverter converter;
|
||||
|
||||
private final String collectionName;
|
||||
|
||||
private volatile ClassLoader classLoader = ClassUtils.getDefaultClassLoader();
|
||||
|
||||
private ApplicationContext applicationContext;
|
||||
|
||||
|
||||
/**
|
||||
* Create a MongoDbMessageStore using the provided {@link MongoDbFactory}.and the default collection name.
|
||||
@@ -119,9 +133,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
*/
|
||||
public MongoDbMessageStore(MongoDbFactory mongoDbFactory, String collectionName) {
|
||||
Assert.notNull(mongoDbFactory, "mongoDbFactory must not be null");
|
||||
MessageReadingMongoConverter converter = new MessageReadingMongoConverter(mongoDbFactory, new MongoMappingContext());
|
||||
converter.afterPropertiesSet();
|
||||
this.template = new MongoTemplate(mongoDbFactory, converter);
|
||||
this.converter = new MessageReadingMongoConverter(mongoDbFactory, new MongoMappingContext());
|
||||
this.template = new MongoTemplate(mongoDbFactory, this.converter);
|
||||
this.collectionName = (StringUtils.hasText(collectionName)) ? collectionName : DEFAULT_COLLECTION_NAME;
|
||||
}
|
||||
|
||||
@@ -132,6 +145,20 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
this.classLoader = classLoader;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
|
||||
this.applicationContext = applicationContext;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() throws Exception {
|
||||
if (this.applicationContext != null) {
|
||||
this.template.setApplicationContext(this.applicationContext);
|
||||
this.converter.setApplicationContext(this.applicationContext);
|
||||
}
|
||||
this.converter.afterPropertiesSet();
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Message<T> addMessage(Message<T> message) {
|
||||
Assert.notNull(message, "'message' must not be null");
|
||||
@@ -335,6 +362,10 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
customConverters.add(new DBObjectToUUIDConverter());
|
||||
customConverters.add(new MessageHistoryToDBObjectConverter());
|
||||
customConverters.add(new DBObjectToGenericMessageConverter());
|
||||
customConverters.add(new DBObjectToMutableMessageConverter());
|
||||
customConverters.add(new DBObjectToErrorMessageConverter());
|
||||
customConverters.add(new DBObjectToAdviceMessageConverter());
|
||||
customConverters.add(new ThrowableToBytesConverter());
|
||||
this.setCustomConversions(new CustomConversions(customConverters));
|
||||
super.afterPropertiesSet();
|
||||
}
|
||||
@@ -355,24 +386,18 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
return super.read(clazz, source);
|
||||
}
|
||||
if (source != null) {
|
||||
Map<String, Object> headers = this.normalizeHeaders((Map<String, Object>) source.get("headers"));
|
||||
|
||||
Object payload = source.get("payload");
|
||||
Object payloadType = source.get(PAYLOAD_TYPE_KEY);
|
||||
if (payloadType != null && payload instanceof DBObject) {
|
||||
try {
|
||||
Class<?> payloadClass = ClassUtils.forName(payloadType.toString(), classLoader);
|
||||
payload = this.read(payloadClass, (DBObject) payload);
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new IllegalStateException("failed to load class: " + payloadType, e);
|
||||
}
|
||||
Message<?> message = null;
|
||||
Object messageType = source.get("_messageType");
|
||||
if (messageType == null) {
|
||||
messageType = GenericMessage.class.getName();
|
||||
}
|
||||
GenericMessage message = new GenericMessage(payload, headers);
|
||||
Map innerMap = (Map) new DirectFieldAccessor(message.getHeaders()).getPropertyValue("headers");
|
||||
// using reflection to set ID and TIMESTAMP since they are immutable through MessageHeaders
|
||||
innerMap.put(MessageHeaders.ID, headers.get(MessageHeaders.ID));
|
||||
innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP));
|
||||
try {
|
||||
message = (Message<?>) this.read(ClassUtils.forName(messageType.toString(), classLoader), source);
|
||||
}
|
||||
catch (ClassNotFoundException e) {
|
||||
throw new IllegalStateException("failed to load class: " + messageType, e);
|
||||
}
|
||||
|
||||
Long groupTimestamp = (Long)source.get(GROUP_TIMESTAMP_KEY);
|
||||
Long lastModified = (Long)source.get(GROUP_UPDATE_TIMESTAMP_KEY);
|
||||
Integer lastReleasedSequenceNumber = (Integer)source.get(LAST_RELEASED_SEQUENCE_NUMBER);
|
||||
@@ -432,8 +457,34 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
}
|
||||
return normalizedHeaders;
|
||||
}
|
||||
|
||||
private Object extractPayload(DBObject source) {
|
||||
Object payload = source.get("payload");
|
||||
if (payload instanceof DBObject) {
|
||||
DBObject payloadObject = (DBObject) payload;
|
||||
Object payloadType = payloadObject.get("_class");
|
||||
try {
|
||||
Class<?> payloadClass = ClassUtils.forName(payloadType.toString(), classLoader);
|
||||
payload = this.read(payloadClass, payloadObject);
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new IllegalStateException("failed to load class: " + payloadType, e);
|
||||
}
|
||||
}
|
||||
return payload;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static void enhanceHeaders(MessageHeaders messageHeaders, Map<String, Object> headers) {
|
||||
Map<String, Object> innerMap = (Map<String, Object>) new DirectFieldAccessor(messageHeaders).getPropertyValue("headers");
|
||||
// using reflection to set ID and TIMESTAMP since they are immutable through MessageHeaders
|
||||
innerMap.put(MessageHeaders.ID, headers.get(MessageHeaders.ID));
|
||||
innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP));
|
||||
}
|
||||
|
||||
|
||||
private static class UuidToDBObjectConverter implements Converter<UUID, DBObject> {
|
||||
@Override
|
||||
public DBObject convert(UUID source) {
|
||||
@@ -474,53 +525,118 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
private class DBObjectToGenericMessageConverter implements Converter<DBObject, GenericMessage<?>> {
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("unchecked")
|
||||
public GenericMessage<?> convert(DBObject source) {
|
||||
MessageReadingMongoConverter converter = (MessageReadingMongoConverter) MongoDbMessageStore.this.template
|
||||
.getConverter();
|
||||
Map<String, Object> headers = converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
|
||||
|
||||
Object payload = source.get("payload");
|
||||
Object payloadType = source.get(PAYLOAD_TYPE_KEY);
|
||||
if (payloadType != null && payload instanceof DBObject) {
|
||||
public GenericMessage<?> convert(DBObject source) {
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
|
||||
|
||||
GenericMessage<?> message = new GenericMessage<Object>(MongoDbMessageStore.this.converter.extractPayload(source), headers);
|
||||
enhanceHeaders(message.getHeaders(), headers);
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private class DBObjectToMutableMessageConverter implements Converter<DBObject, MutableMessage<?>> {
|
||||
|
||||
@Override
|
||||
public MutableMessage<?> convert(DBObject source) {
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
|
||||
|
||||
return new MutableMessage<Object>(MongoDbMessageStore.this.converter.extractPayload(source), headers);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private class DBObjectToAdviceMessageConverter implements Converter<DBObject, AdviceMessage> {
|
||||
|
||||
@Override
|
||||
public AdviceMessage convert(DBObject source) {
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
|
||||
|
||||
Message<?> inputMessage = null;
|
||||
|
||||
if (source.get("inputMessage") != null) {
|
||||
DBObject inputMessageObject = (DBObject) source.get("inputMessage");
|
||||
Object inputMessageType = inputMessageObject.get("_class");
|
||||
try {
|
||||
Class<?> payloadClass = ClassUtils.forName(payloadType.toString(), classLoader);
|
||||
payload = converter.read(payloadClass, (DBObject) payload);
|
||||
Class<?> messageClass = ClassUtils.forName(inputMessageType.toString(), classLoader);
|
||||
inputMessage = (Message<?>) MongoDbMessageStore.this.converter.read(messageClass, inputMessageObject);
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new IllegalStateException("failed to load class: " + payloadType, e);
|
||||
throw new IllegalStateException("failed to load class: " + inputMessageType, e);
|
||||
}
|
||||
}
|
||||
|
||||
@SuppressWarnings("rawtypes")
|
||||
GenericMessage<Object> message = new GenericMessage(payload, headers);
|
||||
Map<String, Object> innerMap = (Map<String, Object>) new DirectFieldAccessor(message.getHeaders()).getPropertyValue("headers");
|
||||
// using reflection to set ID and TIMESTAMP since they are immutable through MessageHeaders
|
||||
innerMap.put(MessageHeaders.ID, headers.get(MessageHeaders.ID));
|
||||
innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP));
|
||||
AdviceMessage message = new AdviceMessage(MongoDbMessageStore.this.converter.extractPayload(source), headers, inputMessage);
|
||||
enhanceHeaders(message.getHeaders(), headers);
|
||||
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private class DBObjectToErrorMessageConverter implements Converter<DBObject, ErrorMessage> {
|
||||
|
||||
private final Converter<byte[], Object> deserializingConverter = new DeserializingConverter();
|
||||
|
||||
@Override
|
||||
public ErrorMessage convert(DBObject source) {
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
|
||||
|
||||
Object payload = this.deserializingConverter.convert((byte[]) source.get("payload"));
|
||||
ErrorMessage message = new ErrorMessage((Throwable) payload, headers);
|
||||
enhanceHeaders(message.getHeaders(), headers);
|
||||
|
||||
return message;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@WritingConverter
|
||||
private class ThrowableToBytesConverter implements Converter<Throwable, byte[]> {
|
||||
|
||||
private final Converter<Object, byte[]> serializingConverter = new SerializingConverter();
|
||||
|
||||
@Override
|
||||
public byte[] convert(Throwable source) {
|
||||
return serializingConverter.convert(source);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* Wrapper class used for storing Messages in MongoDB along with their "group" metadata.
|
||||
*/
|
||||
private static final class MessageWrapper {
|
||||
|
||||
/*
|
||||
* Needed as a persistence property to suppress 'Cannot determine IsNewStrategy' MappingException
|
||||
* when the application context is configured with auditing. The document is not
|
||||
* currently Auditable.
|
||||
*/
|
||||
@SuppressWarnings("unused")
|
||||
@Id
|
||||
private String _id;
|
||||
|
||||
private volatile Object _groupId;
|
||||
|
||||
@Transient
|
||||
private final Message<?> message;
|
||||
|
||||
@SuppressWarnings("unused")
|
||||
private final String _messageType;
|
||||
|
||||
private final Object payload;
|
||||
|
||||
@SuppressWarnings("unused")
|
||||
private final Map<String, ?> headers;
|
||||
|
||||
@SuppressWarnings("unused")
|
||||
private final String _payloadType;
|
||||
private final Message<?> inputMessage;
|
||||
|
||||
private volatile long _group_timestamp;
|
||||
|
||||
@@ -533,9 +649,15 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
public MessageWrapper(Message<?> message) {
|
||||
Assert.notNull(message, "'message' must not be null");
|
||||
this.message = message;
|
||||
this._messageType = message.getClass().getName();
|
||||
this.payload = message.getPayload();
|
||||
this.headers = message.getHeaders();
|
||||
this._payloadType = this.payload.getClass().getName();
|
||||
if (message instanceof AdviceMessage) {
|
||||
this.inputMessage = ((AdviceMessage) message).getInputMessage();
|
||||
}
|
||||
else {
|
||||
this.inputMessage = null;
|
||||
}
|
||||
}
|
||||
|
||||
public int get_LastReleasedSequenceNumber() {
|
||||
|
||||
@@ -131,6 +131,7 @@ public abstract class AbstractMongoDbMessageGroupStoreTests extends MongoDbAvail
|
||||
Message<?> messageA = new GenericMessage<String>("A");
|
||||
Message<?> messageB = new GenericMessage<String>("B");
|
||||
store.addMessageToGroup(1, messageA);
|
||||
Thread.sleep(10);
|
||||
store.addMessageToGroup(1, messageB);
|
||||
assertEquals(2, store.messageGroupSize(1));
|
||||
Message<?> out = store.pollMessageFromGroup(1);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2013 the original author or authors.
|
||||
* Copyright 2002-2014 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,22 +18,29 @@ package org.springframework.integration.mongodb.store;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.io.Serializable;
|
||||
import java.util.Properties;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.hamcrest.Matchers;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.integration.history.MessageHistory;
|
||||
import org.springframework.integration.message.AdviceMessage;
|
||||
import org.springframework.integration.message.MutableMessage;
|
||||
import org.springframework.integration.mongodb.rules.MongoDbAvailable;
|
||||
import org.springframework.integration.mongodb.rules.MongoDbAvailableTests;
|
||||
import org.springframework.integration.store.MessageStore;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.messaging.support.ErrorMessage;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
import com.mongodb.Mongo;
|
||||
|
||||
@@ -151,6 +158,108 @@ public abstract class AbstractMongoDbMessageStoreTests extends MongoDbAvailableT
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testInt3076MessageAsPayload() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<?> messageToStore = new GenericMessage<Message<?>>(MessageBuilder.withPayload(p).build());
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage.getPayload() instanceof GenericMessage);
|
||||
assertEquals(messageToStore.getPayload(), retrievedMessage.getPayload());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
assertEquals(((Message<?>) messageToStore.getPayload()).getPayload(), p);
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testInt3076AdviceMessage() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<Person> inputMessage = MessageBuilder.withPayload(p).build();
|
||||
Message<?> messageToStore = new AdviceMessage("foo", inputMessage);
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage instanceof AdviceMessage);
|
||||
assertEquals(messageToStore.getPayload(), retrievedMessage.getPayload());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
assertEquals(inputMessage, ((AdviceMessage) retrievedMessage).getInputMessage());
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testAdviceMessageAsPayload() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<Person> inputMessage = MessageBuilder.withPayload(p).build();
|
||||
Message<?> messageToStore = new GenericMessage<Message<?>>(new AdviceMessage("foo", inputMessage));
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage.getPayload() instanceof AdviceMessage);
|
||||
AdviceMessage adviceMessage = (AdviceMessage) retrievedMessage.getPayload();
|
||||
assertEquals("foo", adviceMessage.getPayload());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
assertEquals(inputMessage, adviceMessage.getInputMessage());
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testMutableMessageAsPayload() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<?> messageToStore = new GenericMessage<Message<?>>(new MutableMessage<Object>(p));
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage.getPayload() instanceof MutableMessage);
|
||||
assertEquals(messageToStore.getPayload(), retrievedMessage.getPayload());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
assertEquals(((Message<?>) messageToStore.getPayload()).getPayload(), p);
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testInt3076ErrorMessage() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<Person> failedMessage = MessageBuilder.withPayload(p).build();
|
||||
MessagingException messagingException;
|
||||
try {
|
||||
throw new RuntimeException("intentional");
|
||||
}
|
||||
catch (Exception e) {
|
||||
messagingException = new MessagingException(failedMessage, "intentional MessagingException", e);
|
||||
}
|
||||
Message<?> messageToStore = new ErrorMessage(messagingException);
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage instanceof ErrorMessage);
|
||||
assertThat(retrievedMessage.getPayload(), Matchers.instanceOf(MessagingException.class));
|
||||
assertEquals("intentional MessagingException", ((MessagingException) retrievedMessage.getPayload()).getMessage());
|
||||
assertEquals(failedMessage, ((MessagingException) retrievedMessage.getPayload()).getFailedMessage());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
}
|
||||
|
||||
|
||||
public static class Foo implements Serializable {
|
||||
/**
|
||||
*
|
||||
|
||||
@@ -16,26 +16,11 @@
|
||||
package org.springframework.integration.mongodb.store;
|
||||
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import org.hamcrest.Matchers;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.context.support.GenericApplicationContext;
|
||||
import org.springframework.data.mongodb.MongoDbFactory;
|
||||
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
|
||||
import org.springframework.integration.message.AdviceMessage;
|
||||
import org.springframework.integration.mongodb.rules.MongoDbAvailable;
|
||||
import org.springframework.integration.store.MessageStore;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.messaging.support.ErrorMessage;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
import com.mongodb.Mongo;
|
||||
|
||||
@@ -56,68 +41,4 @@ public class ConfigurableMongoDbMessageStoreTests extends AbstractMongoDbMessage
|
||||
return mongoDbMessageStore;
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testInt3076MessageAsPayload() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<?> messageToStore = new GenericMessage<Message<?>>(MessageBuilder.withPayload(p).build());
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage.getPayload() instanceof GenericMessage);
|
||||
assertEquals(messageToStore.getPayload(), retrievedMessage.getPayload());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
assertEquals(((Message<?>) messageToStore.getPayload()).getPayload(), p);
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testInt3076AdviceMessage() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<Person> inputMessage = MessageBuilder.withPayload(p).build();
|
||||
Message<?> messageToStore = new AdviceMessage("foo", inputMessage);
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage instanceof AdviceMessage);
|
||||
assertEquals(messageToStore.getPayload(), retrievedMessage.getPayload());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
assertEquals(inputMessage, ((AdviceMessage) retrievedMessage).getInputMessage());
|
||||
assertEquals(messageToStore, retrievedMessage);
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testInt3076ErrorMessage() throws Exception{
|
||||
MessageStore store = this.getMessageStore();
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
Message<Person> failedMessage = MessageBuilder.withPayload(p).build();
|
||||
MessagingException messagingException;
|
||||
try {
|
||||
throw new RuntimeException("intentional");
|
||||
}
|
||||
catch (Exception e) {
|
||||
messagingException = new MessagingException(failedMessage, "intentional MessagingException", e);
|
||||
}
|
||||
Message<?> messageToStore = new ErrorMessage(messagingException);
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertTrue(retrievedMessage instanceof ErrorMessage);
|
||||
assertThat(retrievedMessage.getPayload(), Matchers.instanceOf(MessagingException.class));
|
||||
assertEquals("intentional MessagingException", ((MessagingException) retrievedMessage.getPayload()).getMessage());
|
||||
assertEquals(failedMessage, ((MessagingException) retrievedMessage.getPayload()).getFailedMessage());
|
||||
assertEquals(messageToStore.getHeaders(), retrievedMessage.getHeaders());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2007-2013 the original author or authors
|
||||
* Copyright 2007-2014 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.
|
||||
@@ -33,7 +33,9 @@ public class MongoDbMessageGroupStoreTests extends AbstractMongoDbMessageGroupSt
|
||||
|
||||
@Override
|
||||
protected MongoDbMessageStore getMessageGroupStore() throws Exception {
|
||||
return new MongoDbMessageStore( new SimpleMongoDbFactory(new Mongo(), "test"));
|
||||
MongoDbMessageStore mongoDbMessageStore = new MongoDbMessageStore( new SimpleMongoDbFactory(new Mongo(), "test"));
|
||||
mongoDbMessageStore.afterPropertiesSet();
|
||||
return mongoDbMessageStore;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2013 the original author or authors.
|
||||
* Copyright 2002-2014 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.
|
||||
@@ -24,12 +24,12 @@ import org.junit.Test;
|
||||
|
||||
import org.springframework.data.mongodb.MongoDbFactory;
|
||||
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.integration.mongodb.rules.MongoDbAvailable;
|
||||
import org.springframework.integration.mongodb.rules.MongoDbAvailableTests;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.transformer.ClaimCheckInTransformer;
|
||||
import org.springframework.integration.transformer.ClaimCheckOutTransformer;
|
||||
import org.springframework.messaging.Message;
|
||||
|
||||
import com.mongodb.Mongo;
|
||||
|
||||
@@ -44,6 +44,7 @@ public class MongoDbMessageStoreClaimCheckIntegrationTests extends MongoDbAvaila
|
||||
public void stringPayload() throws Exception {
|
||||
MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test");
|
||||
MongoDbMessageStore messageStore = new MongoDbMessageStore(mongoDbFactory);
|
||||
messageStore.afterPropertiesSet();
|
||||
ClaimCheckInTransformer checkin = new ClaimCheckInTransformer(messageStore);
|
||||
ClaimCheckOutTransformer checkout = new ClaimCheckOutTransformer(messageStore);
|
||||
Message<?> originalMessage = MessageBuilder.withPayload("test1").build();
|
||||
@@ -60,6 +61,7 @@ public class MongoDbMessageStoreClaimCheckIntegrationTests extends MongoDbAvaila
|
||||
public void objectPayload() throws Exception {
|
||||
MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test");
|
||||
MongoDbMessageStore messageStore = new MongoDbMessageStore(mongoDbFactory);
|
||||
messageStore.afterPropertiesSet();
|
||||
ClaimCheckInTransformer checkin = new ClaimCheckInTransformer(messageStore);
|
||||
ClaimCheckOutTransformer checkout = new ClaimCheckOutTransformer(messageStore);
|
||||
Beverage payload = new Beverage();
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2013 the original author or authors.
|
||||
* Copyright 2002-2014 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.
|
||||
@@ -29,7 +29,9 @@ public class MongoDbMessageStoreTests extends AbstractMongoDbMessageStoreTests {
|
||||
|
||||
@Override
|
||||
protected MessageStore getMessageStore() throws Exception {
|
||||
return new MongoDbMessageStore(new SimpleMongoDbFactory(new Mongo(), "test"));
|
||||
MongoDbMessageStore mongoDbMessageStore = new MongoDbMessageStore(new SimpleMongoDbFactory(new Mongo(), "test"));
|
||||
mongoDbMessageStore.afterPropertiesSet();
|
||||
return mongoDbMessageStore;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -6,15 +6,15 @@
|
||||
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd">
|
||||
|
||||
<int:aggregator input-channel="inputChannel" output-channel="outputChannel" message-store="mongoStore"/>
|
||||
|
||||
|
||||
<int:channel id="outputChannel">
|
||||
<int:queue/>
|
||||
</int:channel>
|
||||
|
||||
|
||||
<bean id="mongoStore" class="org.springframework.integration.mongodb.store.MongoDbMessageStore">
|
||||
<constructor-arg ref="mongoConnectionFactory"/>
|
||||
</bean>
|
||||
|
||||
|
||||
<bean id="mongoConnectionFactory" class="org.springframework.data.mongodb.core.SimpleMongoDbFactory">
|
||||
<constructor-arg>
|
||||
<bean class="com.mongodb.Mongo"/>
|
||||
|
||||
@@ -88,8 +88,8 @@
|
||||
|
||||
<para>
|
||||
Spring Integration's MongoDB module provides the <classname>MongoDbMessageStore</classname> which is an implementation of both
|
||||
the <classname>MessageStore</classname> strategy (mainly used by the <emphasis>QueueChannel</emphasis> and <emphasis>ClaimCheck</emphasis>
|
||||
patterns) and the <classname>MessageGroupStore</classname> strategy (mainly used by the <emphasis>Aggregator</emphasis> and
|
||||
the <classname>MessageStore</classname> strategy (mainly used by the <emphasis>ClaimCheck</emphasis>pattern)
|
||||
and the <classname>MessageGroupStore</classname> strategy (mainly used by the <emphasis>Aggregator</emphasis> and
|
||||
<emphasis>Resequencer</emphasis> patterns).
|
||||
</para>
|
||||
|
||||
@@ -111,15 +111,18 @@
|
||||
and an <emphasis>Aggregator</emphasis>. As you can see it is a simple bean configuration, and it expects a
|
||||
<classname>MongoDbFactory</classname> as a constructor argument.
|
||||
</para>
|
||||
<para>
|
||||
The <classname>MongoDbMessageStore</classname> expands the <interfacename>Message</interfacename> as a Mongo document
|
||||
with all nested properties using the Spring Data Mongo Mapping mechanism. It is useful when you need to have access to
|
||||
the <code>payload</code> or <code>headers</code> for auditing or analytics, for example, against stored messages.
|
||||
</para>
|
||||
<important>
|
||||
<para>
|
||||
The <classname>MongoDbMessageStore</classname> uses a custom <classname>MappingMongoConverter</classname> implementation
|
||||
to store <interfacename>Message</interfacename>s as MongoDB documents and there are some limitations
|
||||
for the properties (<code>payload</code> and <code>header</code> values) of the <interfacename>Message</interfacename>.
|
||||
For example an <classname>ErrorMessage</classname> can't be converted to the MongoDB document, because it has an
|
||||
<classname>Exception</classname> property, where the cause property is infintely recursed. Also, there is no ability to
|
||||
configure
|
||||
custom converters for complex domain <code>payload</code>s or <code>header</code> values.
|
||||
For example, there is no ability to configure custom converters for complex domain <code>payload</code>s or <code>header</code> values.
|
||||
Or to provide a custom <classname>MongoTemplate</classname> (or <classname>MappingMongoConverter</classname>).
|
||||
To achieve these capabilities, an alternative MongoDB <interfacename>MessageStore</interfacename> implementation has been
|
||||
introduced; see next paragraph.
|
||||
</para>
|
||||
@@ -128,8 +131,7 @@
|
||||
<emphasis>Spring Integration 3.0</emphasis> introduced the <classname>ConfigurableMongoDbMessageStore</classname> -
|
||||
<interfacename>MessageStore</interfacename> and <interfacename>MessageGroupStore</interfacename> implementation.
|
||||
This class can receive, as a constructor argument, a <classname>MongoTemplate</classname>, with which you can
|
||||
configure with a
|
||||
custom <classname>WriteConcern</classname>, for example. Another constructor requires a
|
||||
configure with a custom <classname>WriteConcern</classname>, for example. Another constructor requires a
|
||||
<classname>MappingMongoConverter</classname>, and a <interfacename>MongoDbFactory</interfacename>,
|
||||
which allows you to provide some custom conversions for <interfacename>Message</interfacename>s and their properties.
|
||||
Note, by default, the <classname>ConfigurableMongoDbMessageStore</classname> uses standard Java serialization
|
||||
@@ -137,8 +139,8 @@
|
||||
properties from <classname>MongoTemplate</classname>, which is built from the provided
|
||||
<interfacename>MongoDbFactory</interfacename> and <classname>MappingMongoConverter</classname>.
|
||||
The default name for the collection stored by the <classname>ConfigurableMongoDbMessageStore</classname> is
|
||||
<code>configurableStoreMessages</code>. It is recommended to use this implementation for robust and flexible solutions.
|
||||
The <classname>MongoDbMessageStore</classname> remains for backward compatibility and may be removed in future releases.
|
||||
<code>configurableStoreMessages</code>. It is recommended to use this implementation for robust and flexible solutions
|
||||
when messages contain complex data types.
|
||||
</para>
|
||||
</section>
|
||||
|
||||
|
||||
Reference in New Issue
Block a user