diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java index 9261f02f14..de177247ee 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java @@ -143,8 +143,6 @@ public class JdbcChannelMessageStore implements PriorityCapableChannelMessageSto private SerializingConverter serializer; - private LobHandler lobHandler = new DefaultLobHandler(); - private MessageRowMapper messageRowMapper; private ChannelMessageStorePreparedStatementSetter preparedStatementSetter; @@ -231,10 +229,10 @@ public class JdbcChannelMessageStore implements PriorityCapableChannelMessageSto * Override the {@link LobHandler} that is used to create and unpack large objects in SQL queries. The default is * fine for almost all platforms, but some Oracle drivers require a native implementation. * @param lobHandler a {@link LobHandler} + * @deprecated since 6.4 (for removal) (with no replacement) in favor of plain JDBC driver support for byte arrays. */ + @Deprecated(forRemoval = true, since = "6.4") public void setLobHandler(LobHandler lobHandler) { - Assert.notNull(lobHandler, "The provided LobHandler must not be null."); - this.lobHandler = lobHandler; } /** @@ -396,8 +394,8 @@ public class JdbcChannelMessageStore implements PriorityCapableChannelMessageSto * and {@link ChannelMessageStorePreparedStatementSetter} was explicitly set using * {@link #setMessageRowMapper(MessageRowMapper)} and * {@link #setPreparedStatementSetter(ChannelMessageStorePreparedStatementSetter)} respectively, the default - * {@link MessageRowMapper} and {@link ChannelMessageStorePreparedStatementSetter} will be instantiate using the - * specified {@link #deserializer} and {@link #lobHandler}. + * {@link MessageRowMapper} and {@link ChannelMessageStorePreparedStatementSetter} will be instantiated using the + * specified {@link #deserializer}. * Also, if the jdbcTemplate's fetchSize property ({@link JdbcTemplate#getFetchSize()}) * is not 1, a warning will be logged. When using the {@link JdbcChannelMessageStore} * with Oracle, the fetchSize value of 1 is needed to ensure FIFO characteristics @@ -409,7 +407,7 @@ public class JdbcChannelMessageStore implements PriorityCapableChannelMessageSto Assert.notNull(this.channelMessageStoreQueryProvider, "A channelMessageStoreQueryProvider must be provided."); if (this.messageRowMapper == null) { - this.messageRowMapper = new MessageRowMapper(this.deserializer, this.lobHandler); + this.messageRowMapper = new MessageRowMapper(this.deserializer); } if (this.jdbcTemplate.getFetchSize() != 1) { @@ -417,8 +415,7 @@ public class JdbcChannelMessageStore implements PriorityCapableChannelMessageSto } if (this.preparedStatementSetter == null) { - this.preparedStatementSetter = new ChannelMessageStorePreparedStatementSetter(this.serializer, - this.lobHandler); + this.preparedStatementSetter = new ChannelMessageStorePreparedStatementSetter(this.serializer); } this.jdbcTemplate.afterPropertiesSet(); } diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java index 0bd914510d..0936816d9c 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2023 the original author or authors. + * Copyright 2002-2024 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. @@ -16,8 +16,6 @@ package org.springframework.integration.jdbc.store; -import java.sql.ResultSet; -import java.sql.SQLException; import java.sql.Timestamp; import java.util.Arrays; import java.util.Collection; @@ -38,6 +36,7 @@ import org.springframework.core.serializer.Serializer; import org.springframework.core.serializer.support.SerializingConverter; import org.springframework.dao.DataIntegrityViolationException; import org.springframework.dao.IncorrectResultSizeDataAccessException; +import org.springframework.integration.jdbc.store.channel.MessageRowMapper; import org.springframework.integration.store.AbstractMessageGroupStore; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupMetadata; @@ -49,9 +48,7 @@ import org.springframework.integration.util.FunctionIterator; import org.springframework.integration.util.UUIDConverter; import org.springframework.jdbc.core.JdbcOperations; import org.springframework.jdbc.core.JdbcTemplate; -import org.springframework.jdbc.core.RowMapper; import org.springframework.jdbc.core.SingleColumnRowMapper; -import org.springframework.jdbc.support.lob.DefaultLobHandler; import org.springframework.jdbc.support.lob.LobHandler; import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.lang.Nullable; @@ -253,8 +250,6 @@ public class JdbcMessageStore extends AbstractMessageGroupStore } } - private final MessageMapper mapper = new MessageMapper(); - private final JdbcOperations jdbcTemplate; private final Map queryCache = new ConcurrentHashMap<>(); @@ -270,9 +265,9 @@ public class JdbcMessageStore extends AbstractMessageGroupStore private boolean deserializerExplicitlySet; - private SerializingConverter serializer; + private MessageRowMapper mapper = new MessageRowMapper(this.deserializer); - private LobHandler lobHandler = new DefaultLobHandler(); + private SerializingConverter serializer; private boolean checkDatabaseOnStart = true; @@ -325,9 +320,11 @@ public class JdbcMessageStore extends AbstractMessageGroupStore * Override the {@link LobHandler} that is used to create and unpack large objects in SQL queries. The default is * fine for almost all platforms, but some Oracle drivers require a native implementation. * @param lobHandler a {@link LobHandler} + * @deprecated since 6.4 (for removal) (with no replacement) in favor of plain JDBC driver support for byte arrays. */ + @Deprecated(forRemoval = true, since = "6.4") public void setLobHandler(LobHandler lobHandler) { - this.lobHandler = lobHandler; + } /** @@ -347,6 +344,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore public void setDeserializer(Deserializer> deserializer) { this.deserializer = new AllowListDeserializingConverter((Deserializer) deserializer); this.deserializerExplicitlySet = true; + this.mapper = new MessageRowMapper(this.deserializer); } /** @@ -459,8 +457,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore ps.setString(1, messageId); // NOSONAR - magic number ps.setString(2, this.region); // NOSONAR - magic number ps.setTimestamp(3, new Timestamp(System.currentTimeMillis())); // NOSONAR - magic number - - this.lobHandler.getLobCreator().setBlobAsBytes(ps, 4, messageBytes); // NOSONAR - magic number + ps.setBytes(4, messageBytes); // NOSONAR - magic number }); } catch (DataIntegrityViolationException ex) { @@ -780,23 +777,4 @@ public class JdbcMessageStore extends AbstractMessageGroupStore return input == null ? null : UUIDConverter.getUUID(input).toString(); } - /** - * Convenience class to be used to unpack a message from a result set row. Uses column named in the result set to - * extract the required data, so that select clause ordering is unimportant. - */ - private final class MessageMapper implements RowMapper> { - - @Override - public Message mapRow(ResultSet rs, int rowNum) throws SQLException { - byte[] messageBytes = JdbcMessageStore.this.lobHandler.getBlobAsBytes(rs, "MESSAGE_BYTES"); - if (messageBytes == null) { - return null; - } - else { - return (Message) JdbcMessageStore.this.deserializer.convert(messageBytes); - } - } - - } - } diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/ChannelMessageStorePreparedStatementSetter.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/ChannelMessageStorePreparedStatementSetter.java index d98bc13f4f..acb7944ccd 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/ChannelMessageStorePreparedStatementSetter.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/ChannelMessageStorePreparedStatementSetter.java @@ -1,5 +1,5 @@ /* - * Copyright 2017-2019 the original author or authors. + * Copyright 2017-2024 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. @@ -39,7 +39,7 @@ import org.springframework.util.Assert; *

* This class can be extended for any custom data structure or columns types. * For this purpose the {@code protected} constructor is provided for inheritors. - * In this case the {@link #serializer} and {@link #lobHandler} are null to avoid + * In this case the {@link #serializer} is {@code null} to avoid * extra serialization actions if the target custom behavior doesn't imply them. * * @author Meherzad Lahewala @@ -53,8 +53,6 @@ public class ChannelMessageStorePreparedStatementSetter { private final SerializingConverter serializer; - private final LobHandler lobHandler; - /** * Instantiate a {@link ChannelMessageStorePreparedStatementSetter} with the provided * serializer and lobHandler, which both must not be null. @@ -62,24 +60,35 @@ public class ChannelMessageStorePreparedStatementSetter { * the request message * @param lobHandler the {@link LobHandler} to store {@code byte[]} of the request * message to prepared statement + * @deprecated since 6.4 (for removal) (if favor of {@link #ChannelMessageStorePreparedStatementSetter(SerializingConverter)}) + * with a plain JDBC driver support for byte arrays. */ + @Deprecated(forRemoval = true, since = "6.4") public ChannelMessageStorePreparedStatementSetter(SerializingConverter serializer, LobHandler lobHandler) { + this(serializer); + } + + /** + * Instantiate a {@link ChannelMessageStorePreparedStatementSetter} with the provided + * serializer and lobHandler, which both must not be null. + * @param serializer the {@link SerializingConverter} to build {@code byte[]} from + * the request message + * @since 6.4 + */ + public ChannelMessageStorePreparedStatementSetter(SerializingConverter serializer) { Assert.notNull(serializer, "'serializer' must not be null"); - Assert.notNull(lobHandler, "'lobHandler' must not be null"); this.serializer = serializer; - this.lobHandler = lobHandler; } /** * The default constructor for inheritors who are not interested in the message * serialization to {@code byte[]}. - * The {@link #serializer} and {@link #lobHandler} are null from this constructor, + * The {@link #serializer} is {@code null} from this constructor, * therefore any serialization isn't happened in the default {@link #setValues} implementation. * A target implementor must ensure the proper custom logic for storing message. */ protected ChannelMessageStorePreparedStatementSetter() { this.serializer = null; - this.lobHandler = null; } /** @@ -91,7 +100,7 @@ public class ChannelMessageStorePreparedStatementSetter { *

  • 3 - region *
  • 4 - createdDate *
  • 5 - priority if enabled, otherwise null - *
  • 6 - serialized message if {@link #serializer} and {@link #lobHandler} are provided. + *
  • 6 - serialized message if {@link #serializer} is provided. * * An inheritor may consider to call this method for population common properties and perform * custom message serialization logic for the parameter #6. @@ -99,7 +108,7 @@ public class ChannelMessageStorePreparedStatementSetter { * @param preparedStatement the {@link PreparedStatement} to populate columns based on the provided arguments * @param requestMessage the {@link Message} to store * @param groupId the group id for the message to store - * @param region the region in the target table to distinguish different data base clients + * @param region the region in the target table to distinguish different database clients * @param priorityEnabled the flag to indicate if priority has to be stored * @throws SQLException the exception throws during data population */ @@ -126,7 +135,7 @@ public class ChannelMessageStorePreparedStatementSetter { if (this.serializer != null) { byte[] messageBytes = this.serializer.convert(requestMessage); - this.lobHandler.getLobCreator().setBlobAsBytes(preparedStatement, 6, messageBytes); // NOSONAR magic number + preparedStatement.setBytes(6, messageBytes); // NOSONAR magic number } } diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/MessageRowMapper.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/MessageRowMapper.java index 631d2ce4e5..c054bebb8f 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/MessageRowMapper.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/channel/MessageRowMapper.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2024 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. @@ -23,6 +23,7 @@ import org.springframework.integration.support.converter.AllowListDeserializingC import org.springframework.jdbc.core.RowMapper; import org.springframework.jdbc.support.lob.LobHandler; import org.springframework.messaging.Message; +import org.springframework.util.Assert; /** * Convenience class to be used to unpack a {@link Message} from a result set @@ -40,22 +41,32 @@ public class MessageRowMapper implements RowMapper> { private final AllowListDeserializingConverter deserializer; - private final LobHandler lobHandler; - /** * Construct an instance based on the provided {@link AllowListDeserializingConverter} * and {@link LobHandler}. * @param deserializer the {@link AllowListDeserializingConverter} to use. * @param lobHandler the {@link LobHandler} to use. + * @deprecated since 6.4 (for removal) (if favor of {@link #MessageRowMapper(AllowListDeserializingConverter)}) + * with a plain JDBC driver support for byte arrays. */ + @Deprecated(forRemoval = true, since = "6.4") public MessageRowMapper(AllowListDeserializingConverter deserializer, LobHandler lobHandler) { + this(deserializer); + } + + /** + * Construct an instance based on the provided {@link AllowListDeserializingConverter}. + * @param deserializer the {@link AllowListDeserializingConverter} to use. + * @since 6.4 + */ + public MessageRowMapper(AllowListDeserializingConverter deserializer) { + Assert.notNull(deserializer, "'deserializer' must not be null"); this.deserializer = deserializer; - this.lobHandler = lobHandler; } @Override public Message mapRow(ResultSet rs, int rowNum) throws SQLException { - byte[] blobAsBytes = this.lobHandler.getBlobAsBytes(rs, "MESSAGE_BYTES"); + byte[] blobAsBytes = rs.getBytes("MESSAGE_BYTES"); if (blobAsBytes == null) { return null; } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java index c37cbc1eec..aa94a7a918 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java @@ -32,7 +32,6 @@ import org.springframework.integration.jdbc.store.JdbcMessageStore; import org.springframework.integration.store.MessageStore; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.test.util.TestUtils; -import org.springframework.jdbc.support.lob.LobHandler; import org.springframework.messaging.Message; import org.springframework.test.util.ReflectionTestUtils; @@ -52,14 +51,14 @@ public class JdbcMessageStoreParserTests { public void testSimpleMessageStoreWithDataSource() { setUp("defaultJdbcMessageStore.xml", getClass()); MessageStore store = context.getBean("messageStore", MessageStore.class); - assertThat(store instanceof JdbcMessageStore).isTrue(); + assertThat(store).isInstanceOf(JdbcMessageStore.class); } @Test public void testSimpleMessageStoreWithTemplate() { setUp("jdbcOperationsJdbcMessageStore.xml", getClass()); MessageStore store = context.getBean("messageStore", MessageStore.class); - assertThat(store instanceof JdbcMessageStore).isTrue(); + assertThat(store).isInstanceOf(JdbcMessageStore.class); } @Test @@ -67,9 +66,9 @@ public class JdbcMessageStoreParserTests { setUp("serializerJdbcMessageStore.xml", getClass()); MessageStore store = context.getBean("messageStore", MessageStore.class); Object serializer = TestUtils.getPropertyValue(store, "serializer.serializer"); - assertThat(serializer instanceof EnhancedSerializer).isTrue(); + assertThat(serializer).isInstanceOf(EnhancedSerializer.class); Object deserializer = TestUtils.getPropertyValue(store, "deserializer.deserializer"); - assertThat(deserializer instanceof EnhancedSerializer).isTrue(); + assertThat(deserializer).isInstanceOf(EnhancedSerializer.class); } @Test @@ -78,7 +77,6 @@ public class JdbcMessageStoreParserTests { MessageStore store = context.getBean("messageStore", MessageStore.class); assertThat(ReflectionTestUtils.getField(store, "region")).isEqualTo("FOO"); assertThat(ReflectionTestUtils.getField(store, "tablePrefix")).isEqualTo("BAR_"); - assertThat(ReflectionTestUtils.getField(store, "lobHandler")).isEqualTo(context.getBean(LobHandler.class)); } @AfterEach @@ -99,8 +97,7 @@ public class JdbcMessageStoreParserTests { private final Deserializer targetDeserializer = new DefaultDeserializer(); public Object deserialize(InputStream inputStream) throws IOException { - Message message = (Message) targetDeserializer.deserialize(inputStream); - return message; + return targetDeserializer.deserialize(inputStream); } public void serialize(Object object, OutputStream outputStream) throws IOException { diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/soupedUpJdbcMessageStore.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/soupedUpJdbcMessageStore.xml index d5e6a89002..fe41120a82 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/soupedUpJdbcMessageStore.xml +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/soupedUpJdbcMessageStore.xml @@ -9,12 +9,9 @@ - - - diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/oracle/OracleJdbcMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/oracle/OracleJdbcMessageStoreTests.java new file mode 100644 index 0000000000..5f2e3ab09a --- /dev/null +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/oracle/OracleJdbcMessageStoreTests.java @@ -0,0 +1,542 @@ +/* + * Copyright 2024 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. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.jdbc.oracle; + +import java.io.BufferedReader; +import java.io.InputStreamReader; +import java.util.Properties; +import java.util.UUID; + +import javax.sql.DataSource; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.condition.DisabledIfEnvironmentVariable; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.IntegrationMessageHeaderAccessor; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.history.MessageHistory; +import org.springframework.integration.jdbc.store.JdbcMessageStore; +import org.springframework.integration.store.MessageGroup; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.integration.test.predicate.MessagePredicate; +import org.springframework.integration.util.UUIDConverter; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.DataSourceTransactionManager; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.annotation.Repeat; +import org.springframework.test.annotation.Rollback; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionTemplate; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Artem Bilan + * + * @since 6.4 + */ +@SpringJUnitConfig +@DirtiesContext +public class OracleJdbcMessageStoreTests implements OracleContainerTest { + + private static final Log LOG = LogFactory.getLog(OracleJdbcMessageStoreTests.class); + + @Autowired + private DataSource dataSource; + + private JdbcMessageStore messageStore; + + @Autowired + private PlatformTransactionManager transactionManager; + + @BeforeEach + public void init() { + messageStore = new JdbcMessageStore(dataSource); + messageStore.setRegion("JdbcMessageStoreTests"); + } + + @AfterEach + public void afterTest() { + final JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource); + new TransactionTemplate(this.transactionManager).execute(status -> { + final int deletedGroupToMessageRows = jdbcTemplate.update("delete from INT_GROUP_TO_MESSAGE"); + final int deletedMessages = jdbcTemplate.update("delete from INT_MESSAGE"); + final int deletedMessageGroups = jdbcTemplate.update("delete from INT_MESSAGE_GROUP"); + + LOG.info(String.format("Cleaning Database - Deleted Messages: %s, " + + "Deleted GroupToMessage Rows: %s, Deleted Message Groups: %s", + deletedMessages, deletedGroupToMessageRows, deletedMessageGroups)); + return null; + }); + } + + @Test + @Transactional + public void testGetNonExistent() { + Message result = messageStore.getMessage(UUID.randomUUID()); + assertThat(result).isNull(); + } + + @Test + @Transactional + public void testAddAndGet() { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + Message result = messageStore.getMessage(saved.getHeaders().getId()); + assertThat(result).isNotNull(); + assertThat(saved).matches(new MessagePredicate(result)); + } + + @Test + @Transactional + public void testWithMessageHistory() { + Message message = new GenericMessage<>("Hello"); + DirectChannel fooChannel = new DirectChannel(); + fooChannel.setBeanName("fooChannel"); + DirectChannel barChannel = new DirectChannel(); + barChannel.setBeanName("barChannel"); + + message = MessageHistory.write(message, fooChannel); + message = MessageHistory.write(message, barChannel); + messageStore.addMessage(message); + message = messageStore.getMessage(message.getHeaders().getId()); + MessageHistory messageHistory = MessageHistory.read(message); + assertThat(messageHistory).isNotNull(); + assertThat(messageHistory.size()).isEqualTo(2); + Properties fooChannelHistory = messageHistory.get(0); + assertThat(fooChannelHistory.get("name")).isEqualTo("fooChannel"); + assertThat(fooChannelHistory.get("type")).isEqualTo("channel"); + } + + @Test + @Transactional + public void testSize() { + Message message = MessageBuilder.withPayload("foo").build(); + messageStore.addMessage(message); + assertThat(messageStore.getMessageCount()).isEqualTo(1); + } + + @Test + @Transactional + public void testSerializer() { + // N.B. these serializers are not realistic (just for test purposes) + messageStore.setSerializer((object, outputStream) -> { + outputStream.write(object.getPayload().toString().getBytes()); + outputStream.flush(); + }); + messageStore.setDeserializer(inputStream -> { + BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream)); + return new GenericMessage<>(reader.readLine()); + }); + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + assertThat(messageStore.getMessage(message.getHeaders().getId())).isNotNull(); + Message result = messageStore.getMessage(saved.getHeaders().getId()); + assertThat(result).isNotNull(); + assertThat(result.getPayload()).isEqualTo("foo"); + } + + @Test + @Transactional + public void testAddAndGetWithDifferentRegion() { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + messageStore.setRegion("FOO"); + Message result = messageStore.getMessage(saved.getHeaders().getId()); + assertThat(result).isNull(); + } + + @Test + @Transactional + public void testAddAndUpdate() { + Message message = MessageBuilder.withPayload("foo").setCorrelationId("X").build(); + message = messageStore.addMessage(message); + message = MessageBuilder.fromMessage(message).setCorrelationId("Y").build(); + message = messageStore.addMessage(message); + assertThat(new IntegrationMessageHeaderAccessor(messageStore.getMessage(message.getHeaders().getId())) + .getCorrelationId()).isEqualTo("Y"); + } + + @Test + @Transactional + public void testAddAndUpdateAlreadySaved() { + Message message = MessageBuilder.withPayload("foo").build(); + message = messageStore.addMessage(message); + Message result = messageStore.addMessage(message); + assertThat(result).isEqualTo(message); + } + + @Test + @Transactional + public void testAddAndUpdateAlreadySavedAndCopied() { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + Message copy = MessageBuilder.fromMessage(saved).build(); + Message result = messageStore.addMessage(copy); + assertThat(result).isEqualTo(copy); + assertThat(result).isEqualTo(saved); + assertThat(messageStore.getMessage(saved.getHeaders().getId())).isNotNull(); + } + + @Test + @Transactional + public void testAddAndUpdateWithChange() { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + Message copy = MessageBuilder.fromMessage(saved).setHeader("newHeader", 1).build(); + Message result = messageStore.addMessage(copy); + assertThat(result).isNotSameAs(saved); + assertThat(saved).matches(new MessagePredicate(result, "newHeader")); + assertThat(messageStore.getMessage(saved.getHeaders().getId())).isNotNull(); + } + + @Test + @Transactional + public void testAddAndRemoveMessageGroup() { + Message message = MessageBuilder.withPayload("foo").build(); + message = messageStore.addMessage(message); + assertThat(messageStore.removeMessage(message.getHeaders().getId())).isNotNull(); + } + + @Test + @Transactional + public void testAddAndGetMessageGroup() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + long now = System.currentTimeMillis(); + messageStore.addMessageToGroup(groupId, message); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(1); + assertThat(group.getTimestamp() >= now).as("Timestamp too early: " + group.getTimestamp() + "<" + now).isTrue(); + } + + @Test + @Transactional + public void testAddAndRemoveMessageFromMessageGroup() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.removeMessagesFromGroup(groupId, message); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(0); + } + + @Test + @Transactional + public void testRemoveMessageGroup() { + JdbcTemplate template = new JdbcTemplate(dataSource); + template.afterPropertiesSet(); + String groupId = "X"; + + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.removeMessageGroup(groupId); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(0); + + String uuidGroupId = UUIDConverter.getUUID(groupId).toString(); + assertThat(template.queryForList( + "SELECT * from INT_GROUP_TO_MESSAGE where GROUP_KEY = '" + uuidGroupId + "'").size() == 0).isTrue(); + } + + @Test + @Transactional + public void testCompleteMessageGroup() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.completeGroup(groupId); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.isComplete()).isTrue(); + assertThat(group.size()).isEqualTo(1); + } + + @Test + @Transactional + public void testUpdateLastReleasedSequence() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.setLastReleasedSequenceNumberForGroup(groupId, 5); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.getLastReleasedMessageSequenceNumber()).isEqualTo(5); + } + + @Test + @Transactional + public void testMessageGroupCount() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").build(); + messageStore.addMessageToGroup(groupId, message); + assertThat(messageStore.getMessageGroupCount()).isEqualTo(1); + } + + @Test + @Transactional + public void testMessageGroupSizes() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").build(); + messageStore.addMessageToGroup(groupId, message); + assertThat(messageStore.getMessageCountForAllMessageGroups()).isEqualTo(1); + } + + @Test + @Transactional + public void testOrderInMessageGroup() throws Exception { + String groupId = "X"; + + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("foo").setCorrelationId(groupId).build()); + Thread.sleep(1); + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("bar").setCorrelationId(groupId).build()); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(2); + assertThat(messageStore.pollMessageFromGroup(groupId).getPayload()).isEqualTo("foo"); + assertThat(messageStore.pollMessageFromGroup(groupId).getPayload()).isEqualTo("bar"); + } + + @Test + @Transactional + @DisabledIfEnvironmentVariable(named = "bamboo_buildKey", matches = ".*?", + disabledReason = "Timing is too short for CI") + public void testExpireMessageGroupOnCreateOnly() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.registerMessageGroupExpiryCallback( + (messageGroupStore, group) -> messageGroupStore.removeMessageGroup(group.getGroupId())); + Thread.sleep(1000); + messageStore.expireMessageGroups(2000); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(1); + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("bar").setCorrelationId(groupId).build()); + Thread.sleep(2001); + messageStore.expireMessageGroups(2000); + group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(0); + } + + @Test + @Transactional + @DisabledIfEnvironmentVariable(named = "bamboo_buildKey", matches = ".*?", + disabledReason = "Timing is too short for CI") + public void testExpireMessageGroupOnIdleOnly() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.setTimeoutOnIdle(true); + messageStore.addMessageToGroup(groupId, message); + messageStore.registerMessageGroupExpiryCallback( + (messageGroupStore, group) -> messageGroupStore.removeMessageGroup(group.getGroupId())); + Thread.sleep(1000); + messageStore.expireMessageGroups(2000); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(1); + Thread.sleep(2000); + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("bar").setCorrelationId(groupId).build()); + group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(2); + Thread.sleep(2000); + messageStore.expireMessageGroups(2000); + group = messageStore.getMessageGroup(groupId); + assertThat(group.size()).isEqualTo(0); + } + + @Test + @Transactional + public void testMessagePollingFromTheGroup() throws Exception { + + final String groupX = "X"; + + messageStore.addMessageToGroup(groupX, MessageBuilder.withPayload("foo").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup(groupX, MessageBuilder.withPayload("bar").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup(groupX, MessageBuilder.withPayload("baz").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup("Y", MessageBuilder.withPayload("barA").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup("Y", MessageBuilder.withPayload("bazA").setCorrelationId(groupX).build()); + Thread.sleep(100); + MessageGroup group = messageStore.getMessageGroup(groupX); + assertThat(group.size()).isEqualTo(3); + + Message message1 = messageStore.pollMessageFromGroup(groupX); + assertThat(message1).isNotNull(); + assertThat(message1.getPayload()).isEqualTo("foo"); + + group = messageStore.getMessageGroup(groupX); + assertThat(group.size()).isEqualTo(2); + + Message message2 = messageStore.pollMessageFromGroup(groupX); + assertThat(message2).isNotNull(); + assertThat(message2.getPayload()).isEqualTo("bar"); + + group = messageStore.getMessageGroup(groupX); + assertThat(group.size()).isEqualTo(1); + } + + @Test + @Transactional + @Rollback(false) + @Repeat(20) + public void testSameMessageToMultipleGroups() { + final String group1Id = "group1"; + final String group2Id = "group2"; + + final Message message = MessageBuilder.withPayload("foo").build(); + + final MessageBuilder builder1 = MessageBuilder.fromMessage(message); + final MessageBuilder builder2 = MessageBuilder.fromMessage(message); + + builder1.setSequenceNumber(1); + builder2.setSequenceNumber(2); + + final Message message1 = builder1.build(); + final Message message2 = builder2.build(); + + messageStore.addMessageToGroup(group1Id, message1); + messageStore.addMessageToGroup(group2Id, message2); + + final Message messageFromGroup1 = messageStore.pollMessageFromGroup(group1Id); + final Message messageFromGroup2 = messageStore.pollMessageFromGroup(group2Id); + + assertThat(messageFromGroup1).isNotNull(); + assertThat(messageFromGroup2).isNotNull(); + + LOG.info("messageFromGroup1: " + messageFromGroup1.getHeaders().getId() + + "; Sequence #: " + new IntegrationMessageHeaderAccessor(messageFromGroup1).getSequenceNumber()); + LOG.info("messageFromGroup2: " + messageFromGroup2.getHeaders().getId() + + "; Sequence #: " + new IntegrationMessageHeaderAccessor(messageFromGroup2).getSequenceNumber()); + + assertThat(messageFromGroup1.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER)) + .isEqualTo(1); + assertThat(messageFromGroup2.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER)) + .isEqualTo(2); + + } + + @Test + @Transactional + @Rollback(false) + @Repeat(20) + public void testSameMessageAndGroupToMultipleRegions() { + final String groupId = "myGroup"; + final String region1 = "region1"; + final String region2 = "region2"; + + final JdbcMessageStore messageStore1 = new JdbcMessageStore(dataSource); + messageStore1.setRegion(region1); + + final JdbcMessageStore messageStore2 = new JdbcMessageStore(dataSource); + messageStore1.setRegion(region2); + + final Message message = MessageBuilder.withPayload("foo").build(); + + final MessageBuilder builder1 = MessageBuilder.fromMessage(message); + final MessageBuilder builder2 = MessageBuilder.fromMessage(message); + + builder1.setSequenceNumber(1); + builder2.setSequenceNumber(2); + + final Message message1 = builder1.build(); + final Message message2 = builder2.build(); + + messageStore1.addMessageToGroup(groupId, message1); + messageStore2.addMessageToGroup(groupId, message2); + + final Message messageFromRegion1 = messageStore1.pollMessageFromGroup(groupId); + final Message messageFromRegion2 = messageStore2.pollMessageFromGroup(groupId); + + assertThat(messageFromRegion1).isNotNull(); + assertThat(messageFromRegion2).isNotNull(); + + LOG.info("messageFromRegion1: " + messageFromRegion1.getHeaders().getId() + + "; Sequence #: " + new IntegrationMessageHeaderAccessor(messageFromRegion1).getSequenceNumber()); + LOG.info("messageFromRegion2: " + messageFromRegion2.getHeaders().getId() + + "; Sequence #: " + new IntegrationMessageHeaderAccessor(messageFromRegion2).getSequenceNumber()); + + assertThat(messageFromRegion1.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER)) + .isEqualTo(1); + assertThat(messageFromRegion2.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER)) + .isEqualTo(2); + } + + @Test + public void testMessageGroupCondition() { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").build(); + this.messageStore.addMessagesToGroup(groupId, message); + this.messageStore.setGroupCondition(groupId, "testCondition"); + assertThat(this.messageStore.getMessageGroup(groupId).getCondition()).isEqualTo("testCondition"); + } + + @Test + public void sameMessageInTwoGroupsNotRemovedByFirstGroup() { + GenericMessage testMessage = new GenericMessage<>("test data"); + + messageStore.addMessageToGroup("1", testMessage); + messageStore.addMessageToGroup("2", testMessage); + + messageStore.removeMessageGroup("1"); + + assertThat(messageStore.getMessageCount()).isEqualTo(1); + + messageStore.removeMessageGroup("2"); + + assertThat(messageStore.getMessageCount()).isEqualTo(0); + } + + @Test + public void removeMessagesFromGroupDontRemoveSameMessageInOtherGroup() { + GenericMessage testMessage = new GenericMessage<>("test data"); + + messageStore.addMessageToGroup("1", testMessage); + messageStore.addMessageToGroup("2", testMessage); + + messageStore.removeMessagesFromGroup("1", testMessage); + + assertThat(messageStore.getMessageCount()).isEqualTo(1); + assertThat(messageStore.messageGroupSize("1")).isEqualTo(0); + assertThat(messageStore.messageGroupSize("2")).isEqualTo(1); + } + + @Configuration(proxyBeanMethods = false) + public static class Config { + + @Bean + DataSource dataSource() { + return OracleContainerTest.dataSource(); + } + + @Bean + PlatformTransactionManager transactionManager(DataSource dataSource) { + return new DataSourceTransactionManager(dataSource); + } + + } + +} diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractJdbcChannelMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractJdbcChannelMessageStoreTests.java index ea8e3db81b..ea134b20b5 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractJdbcChannelMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractJdbcChannelMessageStoreTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2023 the original author or authors. + * Copyright 2002-2024 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. @@ -28,8 +28,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.core.serializer.support.SerializingConverter; import org.springframework.integration.jdbc.store.JdbcChannelMessageStore; import org.springframework.integration.support.MessageBuilder; -import org.springframework.jdbc.support.lob.DefaultLobHandler; -import org.springframework.jdbc.support.lob.LobHandler; import org.springframework.messaging.Message; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; @@ -135,14 +133,12 @@ public abstract class AbstractJdbcChannelMessageStoreTests { private SerializingConverter serializer = new SerializingConverter(); - private LobHandler lobHandler = new DefaultLobHandler(); - @Override public void setValues(PreparedStatement preparedStatement, Message requestMessage, Object groupId, String region, boolean priorityEnabled) throws SQLException { super.setValues(preparedStatement, requestMessage, groupId, region, priorityEnabled); byte[] messageBytes = this.serializer.convert(requestMessage); - this.lobHandler.getLobCreator().setBlobAsBytes(preparedStatement, 6, messageBytes); + preparedStatement.setBytes(6, messageBytes); } }; diff --git a/src/reference/antora/modules/ROOT/pages/jdbc/message-store.adoc b/src/reference/antora/modules/ROOT/pages/jdbc/message-store.adoc index d52974c508..c8b5ef0488 100644 --- a/src/reference/antora/modules/ROOT/pages/jdbc/message-store.adoc +++ b/src/reference/antora/modules/ROOT/pages/jdbc/message-store.adoc @@ -46,11 +46,10 @@ The following example shows some other optional attributes: [source,xml] ---- - + ---- -In the preceding example, we have specified a `LobHandler` for dealing with messages as large objects (which is often necessary for Oracle) and a prefix for the table names in the queries generated by the store. +In the preceding example, we have specified a prefix for the table names in the queries generated by the store. The table name prefix defaults to `INT_`. [[jdbc-message-store-channels]] diff --git a/src/reference/antora/modules/ROOT/pages/whats-new.adoc b/src/reference/antora/modules/ROOT/pages/whats-new.adoc index a08aec7fda..51e1035129 100644 --- a/src/reference/antora/modules/ROOT/pages/whats-new.adoc +++ b/src/reference/antora/modules/ROOT/pages/whats-new.adoc @@ -23,4 +23,11 @@ In general the project has been moved to the latest dependency versions. === Remote File Adapters Changes The `AbstractRemoteFileStreamingMessageSource` has now a convenient `clearFetchedCache()` API to remove references from cache for not processed remote files. -The references stay in cache because polling configuration does not allow to process all the fetched in one cycle, and the target `SessionFactory` might be changed between polling cycles, e.g. via `RotatingServerAdvice`. \ No newline at end of file +The references stay in cache because polling configuration does not allow to process all the fetched in one cycle, and the target `SessionFactory` might be changed between polling cycles, e.g. via `RotatingServerAdvice`. + +[[x6.4-jdbc-changes]] +=== JDBC Changes + +The `LobHandler` (and respective API) has been deprecated for removal in Spring Framework `6.2`. +Respective option on `JdbcMessageStore` (and similar) have been deprecated as well. +The byte array handling for serialized message is fully deferred to JDBC driver. \ No newline at end of file