diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoMessageStore.java index 9e4c02717f..be36229e23 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoMessageStore.java @@ -23,16 +23,13 @@ import java.util.Map; import java.util.Set; import java.util.UUID; -import org.bson.types.ObjectId; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.BeanClassLoaderAware; import org.springframework.core.convert.converter.Converter; import org.springframework.data.mapping.context.MappingContext; import org.springframework.data.mongodb.MongoDbFactory; import org.springframework.data.mongodb.core.MongoTemplate; -import org.springframework.data.mongodb.core.SimpleMongoDbFactory; import org.springframework.data.mongodb.core.convert.MappingMongoConverter; -import org.springframework.data.mongodb.core.convert.MongoConverter; import org.springframework.data.mongodb.core.mapping.MongoMappingContext; import org.springframework.data.mongodb.core.mapping.MongoPersistentEntity; import org.springframework.data.mongodb.core.mapping.MongoPersistentProperty; @@ -45,7 +42,6 @@ import org.springframework.util.Assert; import org.springframework.util.ClassUtils; import com.mongodb.DBObject; -import com.mongodb.Mongo; /** * @author Mark Fisher @@ -55,39 +51,26 @@ public class MongoMessageStore implements MessageStore, BeanClassLoaderAware { private final static String DEFAULT_COLLECTION_NAME = "messages"; - private final MongoTemplate template; - - private volatile String collectionName = DEFAULT_COLLECTION_NAME; - //private final MongoConverter mongoConverter = new MessageReadingMongoConverter(); - private final MongoConverter mongoConverter = null; + private final MongoTemplate template; + + private final String collectionName; private volatile ClassLoader classLoader = ClassUtils.getDefaultClassLoader(); - public MongoMessageStore(Mongo mongo, String databaseName) { - Assert.notNull(mongo, "mongo must not be null"); - Assert.hasText(databaseName, "databaseName must not be empty"); - MongoDbFactory mongoFactory = new SimpleMongoDbFactory(mongo, databaseName); - MessageReadingMongoConverter converter = new MessageReadingMongoConverter(mongoFactory, new MongoMappingContext()); - this.template = new MongoTemplate(mongoFactory, converter); -// this.template.setDefaultCollectionName(DEFAULT_COLLECTION_NAME); -// this.template.createCollection(DEFAULT_COLLECTION_NAME); + public MongoMessageStore(MongoDbFactory mongoDbFactory) { + this(mongoDbFactory, null); } - - public void setCollectionName(String collectionName) { - Assert.hasText(collectionName, "collectionName must not be empty"); - this.collectionName = collectionName; + public MongoMessageStore(MongoDbFactory mongoDbFactory, String collectionName) { + Assert.notNull(mongoDbFactory, "mongoDbFactory must not be null"); + MessageReadingMongoConverter converter = new MessageReadingMongoConverter(mongoDbFactory, new MongoMappingContext()); + this.template = new MongoTemplate(mongoDbFactory, converter); + this.collectionName = (collectionName != null) ? collectionName : DEFAULT_COLLECTION_NAME; + //this.template.createCollection(collectionName); } -// public void setUsername(String username) { -// template.setUsername(username); -// } -// -// public void setPassword(String password) { -// template.setPassword(password); -// } public void setBeanClassLoader(ClassLoader classLoader) { Assert.notNull(classLoader, "classLoader must not be null"); @@ -119,11 +102,9 @@ public class MongoMessageStore implements MessageStore, BeanClassLoaderAware { private class MessageReadingMongoConverter extends MappingMongoConverter { - public MessageReadingMongoConverter( - MongoDbFactory mongoDbFactory, + public MessageReadingMongoConverter(MongoDbFactory mongoDbFactory, MappingContext, MongoPersistentProperty> mappingContext) { super(mongoDbFactory, mappingContext); - // TODO Auto-generated constructor stub } @@ -140,7 +121,7 @@ public class MongoMessageStore implements MessageStore, BeanClassLoaderAware { @Override public void write(Object source, DBObject target) { if (source instanceof Message) { - String payloadType = ((Message) source).getPayload().getClass().getName(); + //String payloadType = ((Message) source).getPayload().getClass().getName(); //target.put("_payloadType", payloadType); target.put("_id", ((Message) source).getHeaders().getId().toString()); //target.put("_id", ((Message) source).getHeaders().getId().toString()); diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoClaimCheckIntegrationTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoClaimCheckIntegrationTests.java index 43e4c339b9..58fad5b203 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoClaimCheckIntegrationTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoClaimCheckIntegrationTests.java @@ -17,6 +17,8 @@ package org.springframework.integration.mongodb.store; import org.junit.Test; +import org.springframework.data.mongodb.MongoDbFactory; +import org.springframework.data.mongodb.core.SimpleMongoDbFactory; import org.springframework.integration.Message; import org.springframework.integration.mongodb.rules.MongodbAvailable; import org.springframework.integration.mongodb.rules.MongodbAvailableTests; @@ -36,8 +38,8 @@ public class MongoClaimCheckIntegrationTests extends MongodbAvailableTests{ @Test @MongodbAvailable public void stringPayload() throws Exception { - Mongo mongo = new Mongo(); - MongoMessageStore messageStore = new MongoMessageStore(mongo, "test"); + MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test"); + MongoMessageStore messageStore = new MongoMessageStore(mongoDbFactory); ClaimCheckInTransformer checkin = new ClaimCheckInTransformer(messageStore); ClaimCheckOutTransformer checkout = new ClaimCheckOutTransformer(messageStore); Message originalMessage = MessageBuilder.withPayload("test1").build(); @@ -55,8 +57,8 @@ public class MongoClaimCheckIntegrationTests extends MongodbAvailableTests{ @Test @MongodbAvailable public void objectPayload() throws Exception { - Mongo mongo = new Mongo(); - MongoMessageStore messageStore = new MongoMessageStore(mongo, "test"); + MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test"); + MongoMessageStore messageStore = new MongoMessageStore(mongoDbFactory); ClaimCheckInTransformer checkin = new ClaimCheckInTransformer(messageStore); ClaimCheckOutTransformer checkout = new ClaimCheckOutTransformer(messageStore); Beverage payload = new Beverage(); diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoMessageStoreTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoMessageStoreTests.java index 7d3b2d74a2..4ccf3f5799 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoMessageStoreTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoMessageStoreTests.java @@ -17,6 +17,9 @@ package org.springframework.integration.mongodb.store; import org.junit.Test; + +import org.springframework.data.mongodb.MongoDbFactory; +import org.springframework.data.mongodb.core.SimpleMongoDbFactory; import org.springframework.integration.Message; import org.springframework.integration.mongodb.rules.MongodbAvailable; import org.springframework.integration.mongodb.rules.MongodbAvailableTests; @@ -35,43 +38,52 @@ public class MongoMessageStoreTests extends MongodbAvailableTests{ @Test @MongodbAvailable public void addGetWithStringPayload() throws Exception { - Mongo mongo = new Mongo(); - MongoMessageStore store = new MongoMessageStore(mongo, "test"); - Message message = MessageBuilder.withPayload("Hello").build(); - System.out.println(message); - store.addMessage(message); - assertNotNull(store.getMessage(message.getHeaders().getId())); + MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test"); + MongoMessageStore store = new MongoMessageStore(mongoDbFactory); + Message messageToStore = MessageBuilder.withPayload("Hello").build(); + //System.out.println("before: " + messageToStore); + store.addMessage(messageToStore); + Message retrievedMessage = store.getMessage(messageToStore.getHeaders().getId()); + //System.out.println("after: " + retrievedMessage); + assertNotNull(retrievedMessage); } @Test @MongodbAvailable public void addGetWithObjectDefaultConstructorPayload() throws Exception { - Mongo mongo = new Mongo(); - MongoMessageStore store = new MongoMessageStore(mongo, "test"); + MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test"); + MongoMessageStore store = new MongoMessageStore(mongoDbFactory); Person p = new Person(); p.setFname("John"); p.setLname("Doe"); - - Message message = MessageBuilder.withPayload(p).build(); - System.out.println(message); - store.addMessage(message); - Message m = store.getMessage(message.getHeaders().getId()); - assertNotNull(m); + Message messageToStore = MessageBuilder.withPayload(p).build(); + //System.out.println("before: " + messageToStore); + store.addMessage(messageToStore); + Message retrievedMessage = store.getMessage(messageToStore.getHeaders().getId()); + //System.out.println("after: " + retrievedMessage); + assertNotNull(retrievedMessage); } - - public static class Person{ + + + public static class Person { + private String fname; + private String lname; + public String getFname() { return fname; } + public void setFname(String fname) { this.fname = fname; } + public String getLname() { return lname; } + public void setLname(String lname) { this.lname = lname; }