diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index e2ef00d79a..8a49022b97 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -299,6 +299,74 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH return messageStore; } + protected Map> getExpireGroupScheduledFutures() { + return expireGroupScheduledFutures; + } + + protected MessageGroupProcessor getOutputProcessor() { + return outputProcessor; + } + + protected CorrelationStrategy getCorrelationStrategy() { + return correlationStrategy; + } + + protected ReleaseStrategy getReleaseStrategy() { + return releaseStrategy; + } + + protected MessageChannel getOutputChannel() { + return outputChannel; + } + + protected String getOutputChannelName() { + return outputChannelName; + } + + protected MessagingTemplate getMessagingTemplate() { + return messagingTemplate; + } + + protected MessageChannel getDiscardChannel() { + return discardChannel; + } + + protected String getDiscardChannelName() { + return discardChannelName; + } + + protected boolean isSendPartialResultOnExpiry() { + return sendPartialResultOnExpiry; + } + + protected boolean isSequenceAware() { + return sequenceAware; + } + + protected LockRegistry getLockRegistry() { + return lockRegistry; + } + + protected boolean isLockRegistrySet() { + return lockRegistrySet; + } + + protected long getMinimumTimeoutForEmptyGroups() { + return minimumTimeoutForEmptyGroups; + } + + protected boolean isReleasePartialSequences() { + return releasePartialSequences; + } + + protected Expression getGroupTimeoutExpression() { + return groupTimeoutExpression; + } + + protected EvaluationContext getEvaluationContext() { + return evaluationContext; + } + @Override protected void handleMessageInternal(Message message) throws Exception { Object correlationKey = correlationStrategy.getCorrelationKey(message); @@ -480,11 +548,11 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH return new IntegrationMessageHeaderAccessor(lastReleasedMessage).getSequenceNumber(); } - private MessageGroup store(Object correlationKey, Message message) { + protected MessageGroup store(Object correlationKey, Message message) { return messageStore.addMessageToGroup(correlationKey, message); } - private void expireGroup(Object correlationKey, MessageGroup group) { + protected void expireGroup(Object correlationKey, MessageGroup group) { if (logger.isInfoEnabled()) { logger.info("Expiring MessageGroup with correlationKey[" + correlationKey + "]"); } @@ -506,7 +574,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH } } - private void completeGroup(Object correlationKey, MessageGroup group) { + protected void completeGroup(Object correlationKey, MessageGroup group) { Message first = null; if (group != null) { first = group.getOne(); @@ -515,7 +583,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH } @SuppressWarnings("unchecked") - private Collection> completeGroup(Message message, Object correlationKey, MessageGroup group) { + protected Collection> completeGroup(Message message, Object correlationKey, MessageGroup group) { if (logger.isDebugEnabled()) { logger.debug("Completing group with correlationKey [" + correlationKey + "]"); } @@ -530,13 +598,13 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH return partialSequence; } - private void verifyResultCollectionConsistsOfMessages(Collection elements){ + protected void verifyResultCollectionConsistsOfMessages(Collection elements){ Class commonElementType = CollectionUtils.findCommonElementType(elements); Assert.isAssignable(Message.class, commonElementType, "The expected collection of Messages contains non-Message element: " + commonElementType); } @SuppressWarnings("rawtypes") - private void sendReplies(Object processorResult, Message message) { + protected void sendReplies(Object processorResult, Message message) { Object replyChannelHeader = null; if (message != null) { replyChannelHeader = message.getHeaders().getReplyChannel(); @@ -555,7 +623,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH } } - private void sendReplyMessage(Object reply, Object replyChannel) { + protected void sendReplyMessage(Object reply, Object replyChannel) { if (replyChannel instanceof MessageChannel) { if (reply instanceof Message) { this.messagingTemplate.send((MessageChannel) replyChannel, (Message) reply); @@ -573,7 +641,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH } } - private boolean shouldSendMultipleReplies(Iterable iter) { + protected boolean shouldSendMultipleReplies(Iterable iter) { for (Object next : iter) { if (next instanceof Message) { return true; @@ -582,7 +650,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH return false; } - private Long obtainGroupTimeout(MessageGroup group) { + protected Long obtainGroupTimeout(MessageGroup group) { return this.groupTimeoutExpression != null ? this.groupTimeoutExpression.getValue(this.evaluationContext, group, Long.class) : null; } @@ -594,7 +662,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH } } - private static class SequenceAwareMessageGroup extends SimpleMessageGroup { + protected static class SequenceAwareMessageGroup extends SimpleMessageGroup { public SequenceAwareMessageGroup(MessageGroup messageGroup) { super(messageGroup); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java index 9a2f573c11..75fb50dd11 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractMessageGroupStore.java @@ -26,6 +26,7 @@ import org.springframework.integration.support.MessageBuilderFactory; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.jmx.export.annotation.ManagedResource; +import org.springframework.messaging.Message; /** * @author Dave Syer @@ -135,6 +136,16 @@ public abstract class AbstractMessageGroupStore implements MessageGroupStore, It return count; } + @Override + public MessageGroupMetadata getGroupMetadata(Object groupId) { + throw new UnsupportedOperationException("Not yet implemented for this store"); + } + + @Override + public Message getOneMessageFromGroup(Object groupId) { + throw new UnsupportedOperationException("Not yet implemented for this store"); + } + private void expire(MessageGroup group) { RuntimeException exception = null; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java index 2b9811638e..d476cc73ea 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java @@ -28,6 +28,7 @@ import org.springframework.messaging.Message; * * @author Dave Syer * @author Oleg Zhurakousky + * @author Gary Russell */ public interface MessageGroup { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java index 67009a00b5..dd558804f6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupMetadata.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2011 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. @@ -27,43 +27,65 @@ import org.springframework.util.Assert; /** * Immutable Value Object holding metadata about a MessageGroup. - * + * * @author Oleg Zhurakousky + * @author Gary Russell * @since 2.1 */ -public class MessageGroupMetadata implements Serializable{ +public class MessageGroupMetadata implements Serializable { private static final long serialVersionUID = 1L; - + private final Object groupId; - + private final List messageIds = new LinkedList(); private final boolean complete; private final long timestamp; - + private volatile long lastModified; private final int lastReleasedMessageSequenceNumber; + private final boolean hasMessages; + + private final int size; + + private final UUID first; + public MessageGroupMetadata(MessageGroup messageGroup) { - + this(messageGroup, true, null); + } + + public MessageGroupMetadata(MessageGroup messageGroup, boolean hasMessages, UUID first) { + Assert.notNull(messageGroup, "'messageGroup' must not be null"); this.groupId = messageGroup.getGroupId(); - for (Message message : messageGroup.getMessages()) { - this.messageIds.add(message.getHeaders().getId()); + if (hasMessages) { + for (Message message : messageGroup.getMessages()) { + this.messageIds.add(message.getHeaders().getId()); + } + this.size = this.messageIds.size(); + } + else { + this.size = messageGroup.size(); } this.complete = messageGroup.isComplete(); this.timestamp = messageGroup.getTimestamp(); this.lastReleasedMessageSequenceNumber = messageGroup.getLastReleasedMessageSequenceNumber(); this.lastModified = messageGroup.getLastModified(); + this.hasMessages = hasMessages; + this.first = first; } public void remove(UUID messageId){ + if (!this.hasMessages) { + throw new IllegalStateException("Messages are not available, fetch the entire group"); + } this.messageIds.remove(messageId); } - + public void setLastModified(long lastModified) { this.lastModified = lastModified; } @@ -71,27 +93,35 @@ public class MessageGroupMetadata implements Serializable{ public Object getGroupId() { return this.groupId; } - + public Iterator messageIdIterator(){ + if (!this.hasMessages) { + throw new IllegalStateException("Messages are not available, fetch the entire group"); + } return this.messageIds.iterator(); } - + public int size(){ - return this.messageIds.size(); + return this.size; } - + public UUID firstId(){ + if (this.first != null) { + return this.first; + } + if (!this.hasMessages) { + throw new IllegalStateException("Messages are not available, fetch the entire group"); + } if (this.messageIds.size() > 0){ return this.messageIds.iterator().next(); } - return null; } public boolean isComplete() { return this.complete; } - + public long getLastModified() { return lastModified; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java index 7b216ea82e..0f5f9fe07e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroupStore.java @@ -101,6 +101,23 @@ public interface MessageGroupStore extends BasicMessageGroupStore { */ void completeGroup(Object groupId); + /** + * Obtain the group metadata without fetching any messages; must supply all other + * group properties; may include the id of the first message. + * @param groupId The group id. + * @return The metadata. + * @since 4.0 + */ + MessageGroupMetadata getGroupMetadata(Object groupId); + + /** + * Return the one {@link org.springframework.messaging.Message} from {@link org.springframework.integration.store.MessageGroup}. + * @param groupId The group identifier. + * @return the {@link org.springframework.messaging.Message}. + * @since 4.0 + */ + Message getOneMessageFromGroup(Object groupId); + /** * Invoked when a MessageGroupStore expires a group. */ diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java index 7fa5d0b599..57f47f7b12 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java @@ -68,6 +68,7 @@ public class SimpleMessageGroup implements MessageGroup { this(messageGroup.getMessages(), messageGroup.getGroupId(), messageGroup.getTimestamp(), messageGroup.isComplete()); } + @Override public long getTimestamp() { return timestamp; } @@ -76,10 +77,12 @@ public class SimpleMessageGroup implements MessageGroup { this.lastModified = lastModified; } + @Override public long getLastModified() { return lastModified; } + @Override public boolean canAdd(Message message) { return true; } @@ -92,6 +95,7 @@ public class SimpleMessageGroup implements MessageGroup { messages.remove(message); } + @Override public int getLastReleasedMessageSequenceNumber() { return lastReleasedMessageSequence; } @@ -100,6 +104,7 @@ public class SimpleMessageGroup implements MessageGroup { return this.messages.offer(message); } + @Override public Collection> getMessages() { return Collections.unmodifiableCollection(messages); } @@ -108,18 +113,22 @@ public class SimpleMessageGroup implements MessageGroup { this.lastReleasedMessageSequence = sequenceNumber; } + @Override public Object getGroupId() { return groupId; } + @Override public boolean isComplete() { return this.complete; } + @Override public void complete() { this.complete = true; } + @Override public int getSequenceSize() { if (size() == 0) { return 0; @@ -127,10 +136,12 @@ public class SimpleMessageGroup implements MessageGroup { return new IntegrationMessageHeaderAccessor(getOne()).getSequenceSize(); } + @Override public int size() { return this.messages.size(); } + @Override public Message getOne() { Message one = messages.peek(); return one; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java index 8ed914f86b..d9fef25c1e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java @@ -298,4 +298,15 @@ public class SimpleMessageStore extends AbstractMessageGroupStore public int messageGroupSize(Object groupId) { return this.getMessageGroup(groupId).size(); } + + @Override + public MessageGroupMetadata getGroupMetadata(Object groupId) { + return new MessageGroupMetadata(this.getMessageGroup(groupId)); + } + + @Override + public Message getOneMessageFromGroup(Object groupId) { + return this.getMessageGroup(groupId).getOne(); + } + } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java index 10e147bc3f..8c0e750352 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java @@ -23,9 +23,6 @@ import java.util.LinkedHashSet; import java.util.List; import java.util.UUID; -import com.mongodb.DB; -import com.mongodb.MongoException; - import org.springframework.dao.DataAccessException; import org.springframework.data.domain.Sort; import org.springframework.data.mongodb.MongoDbFactory; @@ -36,6 +33,7 @@ import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.data.mongodb.core.query.Query; import org.springframework.data.mongodb.core.query.Update; import org.springframework.integration.store.MessageGroup; +import org.springframework.integration.store.MessageGroupMetadata; import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.store.MessageStore; import org.springframework.integration.store.SimpleMessageGroup; @@ -43,6 +41,9 @@ import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.messaging.Message; import org.springframework.util.Assert; +import com.mongodb.DB; +import com.mongodb.MongoException; + /** * An alternate MongoDB {@link MessageStore} and {@link MessageGroupStore} which allows the user to * configure the instance of {@link MongoTemplate}. The mechanism of storing the messages/group of messages @@ -322,6 +323,16 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb .size(); } + @Override + public MessageGroupMetadata getGroupMetadata(Object groupId) { + throw new UnsupportedOperationException("Not yet implemented for this store"); + } + + @Override + public Message getOneMessageFromGroup(Object groupId) { + throw new UnsupportedOperationException("Not yet implemented for this store"); + } + private void expire(MessageGroup group) { RuntimeException exception = null;