From 6140aa7b2cfb6ae309c55a157e94b44e5d0bea4f Mon Sep 17 00:00:00 2001 From: willschipp Date: Tue, 28 May 2013 14:58:37 -0400 Subject: [PATCH] INT-3037 Fix JDBC MS Discard After Completion INT-3037 - commented out sizing component unit test to support cleanup updated commit as per Artem's comments May 29 --- .../integration/jdbc/JdbcMessageStore.java | 34 +++++++++-- .../jdbc/JdbcMessageStoreTests.java | 58 +++++++++++++++++-- 2 files changed, 83 insertions(+), 9 deletions(-) diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java index 190f5f5682..c335aa70eb 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java @@ -75,6 +75,7 @@ import org.springframework.util.StringUtils; * @author Oleg Zhurakousky * @author Matt Stine * @author Gunnar Hillert + * @author Will Schipp * * @since 2.0 */ @@ -293,6 +294,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa Assert.state(jdbcTemplate != null, "A DataSource or JdbcTemplate must be provided"); } + @Override public Message removeMessage(UUID id) { Message message = getMessage(id); if (message == null) { @@ -306,11 +308,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return null; } + @Override @ManagedAttribute public long getMessageCount() { return jdbcTemplate.queryForInt(getQuery(Query.GET_MESSAGE_COUNT), region); } + @Override public Message getMessage(UUID id) { List> list = jdbcTemplate.query(getQuery(Query.GET_MESSAGE), new Object[] { getKey(id), region }, mapper); if (list.isEmpty()) { @@ -319,6 +323,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return list.get(0); } + @Override @SuppressWarnings({ "rawtypes", "unchecked" }) public Message addMessage(final Message message) { if (message.getHeaders().containsKey(SAVED_KEY)) { @@ -342,6 +347,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa final byte[] messageBytes = serializer.convert(result); jdbcTemplate.update(getQuery(Query.CREATE_MESSAGE), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Inserting message with id key=" + messageId); @@ -355,6 +361,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return result; } + @Override public MessageGroup addMessageToGroup(Object groupId, Message message) { final String groupKey = getKey(groupId); final String messageId = getKey(message.getHeaders().getId()); @@ -381,6 +388,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa this.addMessage(message); jdbcTemplate.update(getQuery(Query.CREATE_GROUP_TO_MESSAGE), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Inserting message with id key=" + messageId + " and created date=" + createdDate); @@ -406,12 +414,14 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return jdbcTemplate.queryForInt(getQuery(Query.COUNT_ALL_MESSAGES_IN_GROUPS), region); } + @Override @ManagedAttribute public int messageGroupSize(Object groupId) { String key = getKey(groupId); return jdbcTemplate.queryForInt(getQuery(Query.COUNT_ALL_MESSAGES_IN_GROUP), key, region); } + @Override public MessageGroup getMessageGroup(Object groupId) { String key = getKey(groupId); final AtomicReference createDate = new AtomicReference(); @@ -421,12 +431,9 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa List> messages = jdbcTemplate.query(getQuery(Query.LIST_MESSAGES_BY_GROUP_KEY), new Object[] { key, region }, mapper); - if (messages.size() == 0){ - return new SimpleMessageGroup(groupId); - } - jdbcTemplate.query(getQuery(Query.GET_GROUP_INFO), new Object[] { key, region}, new RowCallbackHandler() { + @Override public void processRow(ResultSet rs) throws SQLException { updateDate.set(rs.getTimestamp("UPDATED_DATE")); @@ -459,11 +466,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return messageGroup; } + @Override public MessageGroup removeMessageFromGroup(Object groupId, Message messageToRemove) { final String groupKey = getKey(groupId); final String messageId = getKey(messageToRemove.getHeaders().getId()); jdbcTemplate.update(getQuery(Query.REMOVE_MESSAGE_FROM_GROUP), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Removing message from group with group key=" + groupKey); @@ -478,6 +487,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return getMessageGroup(groupId); } + @Override public void removeMessageGroup(Object groupId) { final String groupKey = getKey(groupId); @@ -487,6 +497,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa } jdbcTemplate.update(getQuery(Query.REMOVE_GROUP_TO_MESSAGE_JOIN), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Removing relationships for the group with group key=" + groupKey); @@ -497,6 +508,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa }); jdbcTemplate.update(getQuery(Query.DELETE_MESSAGE_GROUP), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Marking messages with group key=" + groupKey); @@ -507,11 +519,13 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa }); } + @Override public void completeGroup(Object groupId) { final long updatedDate = System.currentTimeMillis(); final String groupKey = getKey(groupId); jdbcTemplate.update(getQuery(Query.COMPLETE_GROUP), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Completing MessageGroup: " + groupKey); @@ -523,12 +537,14 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa }); } + @Override public void setLastReleasedSequenceNumberForGroup(Object groupId, final int sequenceNumber) { Assert.notNull(groupId, "'groupId' must not be null"); final long updatedDate = System.currentTimeMillis(); final String groupKey = getKey(groupId); jdbcTemplate.update(getQuery(Query.UPDATE_LAST_RELEASED_SEQUENCE), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Updating the sequence number of the last released Message in the MessageGroup: " + groupKey); @@ -542,6 +558,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa this.updateMessageGroup(groupKey); } + @Override public Message pollMessageFromGroup(Object groupId) { String key = getKey(groupId); @@ -552,6 +569,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return polledMessage; } + @Override public Iterator iterator() { final Iterator iterator = jdbcTemplate.query(getQuery(Query.LIST_GROUP_KEYS), new Object[] { region }, @@ -559,14 +577,17 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa return new Iterator() { + @Override public boolean hasNext() { return iterator.hasNext(); } + @Override public MessageGroup next() { return getMessageGroup(iterator.next()); } + @Override public void remove() { throw new UnsupportedOperationException("Cannot remove MessageGroup from this iterator."); } @@ -620,6 +641,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa private void doCreateMessageGroup(final String groupKey, final Timestamp createdDate){ jdbcTemplate.update(getQuery(Query.CREATE_MESSAGE_GROUP), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Creating message group with id key=" + groupKey + " and created date=" + createdDate); @@ -634,6 +656,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa private void doUpdateMessageGroup(final String groupKey, final Timestamp updatedDate){ jdbcTemplate.update(getQuery(Query.UPDATE_MESSAGE_GROUP), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Updating message group with id key=" + groupKey + " and updated date=" + updatedDate); @@ -647,6 +670,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa private void updateMessageGroup(final String groupId){ jdbcTemplate.update(getQuery(Query.UPDATE_GROUP), new PreparedStatementSetter() { + @Override public void setValues(PreparedStatement ps) throws SQLException { if (logger.isDebugEnabled()){ logger.debug("Updating MessageGroup: " + groupId); @@ -666,6 +690,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa jdbcTemplate.query(getQuery(Query.LIST_MESSAGEIDS_BY_GROUP_KEY), new Object[] { key, region }, new RowCallbackHandler() { + @Override public void processRow(ResultSet rs) throws SQLException { messageIds.add(UUID.fromString(rs.getString(1))); } @@ -686,6 +711,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa */ private class MessageMapper implements RowMapper> { + @Override public Message mapRow(ResultSet rs, int rowNum) throws SQLException { Message message = (Message) deserializer.convert(lobHandler.getBlobAsBytes(rs, "MESSAGE_BYTES")); return message; diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java index 56e7245139..a52566a227 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java @@ -65,6 +65,7 @@ import org.springframework.transaction.annotation.Transactional; * @author Gunnar Hillert * @author Artem Bilan * @author Gary Russell + * @author Will Schipp */ @ContextConfiguration @RunWith(SpringJUnit4ClassRunner.class) @@ -137,12 +138,14 @@ public class JdbcMessageStoreTests { public void testSerializer() throws Exception { // N.B. these serializers are not realistic (just for test purposes) messageStore.setSerializer(new Serializer>() { + @Override public void serialize(Message object, OutputStream outputStream) throws IOException { outputStream.write(((Message) object).getPayload().toString().getBytes()); outputStream.flush(); } }); messageStore.setDeserializer(new Deserializer>() { + @Override public GenericMessage deserialize(InputStream inputStream) throws IOException { BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream)); return new GenericMessage(reader.readLine()); @@ -320,6 +323,7 @@ public class JdbcMessageStoreTests { Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); messageStore.addMessageToGroup(groupId, message); messageStore.registerMessageGroupExpiryCallback(new MessageGroupCallback() { + @Override public void execute(MessageGroupStore messageGroupStore, MessageGroup group) { messageGroupStore.removeMessageGroup(group.getGroupId()); } @@ -343,6 +347,7 @@ public class JdbcMessageStoreTests { messageStore.setTimeoutOnIdle(true); messageStore.addMessageToGroup(groupId, message); messageStore.registerMessageGroupExpiryCallback(new MessageGroupCallback() { + @Override public void execute(MessageGroupStore messageGroupStore, MessageGroup group) { messageGroupStore.removeMessageGroup(group.getGroupId()); } @@ -424,8 +429,8 @@ public class JdbcMessageStoreTests { LOG.info("messageFromGroup1: " + messageFromGroup1.getHeaders().getId() + "; Sequence #: " + messageFromGroup1.getHeaders().getSequenceNumber()); LOG.info("messageFromGroup2: " + messageFromGroup2.getHeaders().getId() + "; Sequence #: " + messageFromGroup2.getHeaders().getSequenceNumber()); - assertEquals(Integer.valueOf(1), (Integer) messageFromGroup1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); - assertEquals(Integer.valueOf(2), (Integer) messageFromGroup2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(1), messageFromGroup1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(2), messageFromGroup2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); } @@ -466,8 +471,51 @@ public class JdbcMessageStoreTests { LOG.info("messageFromRegion1: " + messageFromRegion1.getHeaders().getId() + "; Sequence #: " + messageFromRegion1.getHeaders().getSequenceNumber()); LOG.info("messageFromRegion2: " + messageFromRegion2.getHeaders().getId() + "; Sequence #: " + messageFromRegion2.getHeaders().getSequenceNumber()); - assertEquals(Integer.valueOf(1), (Integer) messageFromRegion1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); - assertEquals(Integer.valueOf(2), (Integer) messageFromRegion2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(1), messageFromRegion1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(2), messageFromRegion2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); } -} + + @Test + @Transactional + public void testCompletedNotExpiredGroupINT3037() throws Exception { + /* + * based on the aggregator scenario as follows; + * + * send three messages in + * 1 of 2 + * 2 of 2 + * 2 of 2 (last again) + * + * expected behavior is that the LAST message (2 of 2 repeat) should be on the discard channel + * (discard behavior performed by the AbstractCorrelatingMessageHandler.handleMessageInternal) + */ + final JdbcMessageStore messageStore = new JdbcMessageStore(dataSource); + //init + String groupId = "group"; + //build the messages + Message oneOfTwo = MessageBuilder.withPayload("hello").setSequenceNumber(1).setSequenceSize(2).setCorrelationId(groupId).build(); + Message twoOfTwo = MessageBuilder.withPayload("world").setSequenceNumber(2).setSequenceSize(2).setCorrelationId(groupId).build(); + //add to the messageStore + messageStore.addMessageToGroup(groupId, oneOfTwo); + messageStore.addMessageToGroup(groupId, twoOfTwo); + //check that 2 messages are there + assertTrue(messageStore.getMessageGroupCount() == 1); + assertTrue(messageStore.getMessageCount() == 2); + //retrieve the group (like in the aggregator) + MessageGroup messageGroup = messageStore.getMessageGroup(groupId); + //'complete' the group + messageStore.completeGroup(messageGroup.getGroupId()); + //now clear the messages + for (Message message : messageGroup.getMessages()) { + messageStore.removeMessageFromGroup(groupId, message); + }//end for + //'add' the other message --> emulated by getting the messageGroup + messageGroup = messageStore.getMessageGroup(groupId); + //should be marked 'complete' --> old behavior it would not + assertTrue(messageGroup.isComplete()); + } + + + +} \ No newline at end of file