diff --git a/ci/README.md b/ci/README.md index 660f4c6af0..c9605360a5 100644 --- a/ci/README.md +++ b/ci/README.md @@ -44,6 +44,7 @@ the payload sent to the channel. The other files are optional assertions: | `source_status` | Source status | | `source_response` | Source response payload | | `source_transformed` | Transformed source payload | +| `source_encoded` | Encoded source payload | | `destNN` | Sent payload for destination `NN` | | `destNN_transformed` | Transformed payload for destination `NN` | | `destNN_response` | Response payload from destination `NN` | diff --git a/ci/tests/101-raw-no-op/channels/01-raw-no-op/messages/01-hello-world/source_encoded b/ci/tests/101-raw-no-op/channels/01-raw-no-op/messages/01-hello-world/source_encoded new file mode 100644 index 0000000000..6769dd60bd --- /dev/null +++ b/ci/tests/101-raw-no-op/channels/01-raw-no-op/messages/01-hello-world/source_encoded @@ -0,0 +1 @@ +Hello world! \ No newline at end of file diff --git a/donkey/src/main/java/com/mirth/connect/donkey/server/data/jdbc/JdbcDao.java b/donkey/src/main/java/com/mirth/connect/donkey/server/data/jdbc/JdbcDao.java index a22f8fb1a1..999b863f0c 100644 --- a/donkey/src/main/java/com/mirth/connect/donkey/server/data/jdbc/JdbcDao.java +++ b/donkey/src/main/java/com/mirth/connect/donkey/server/data/jdbc/JdbcDao.java @@ -208,6 +208,12 @@ public void insertMessageContent(MessageContent messageContent) { insertContent(messageContent.getChannelId(), messageContent.getMessageId(), messageContent.getMetaDataId(), messageContent.getContentType(), messageContent.getContent(), messageContent.getDataType(), messageContent.isEncrypted()); } + /* + * The statement is deliberately left open here. It is the cached statement that holds the + * accumulated batch, so closing it would discard every addBatch() made so far; + * executeBatchInsertMessageContent() runs the batch and closes the statement if the + * subclass needs that. + */ @Override public void batchInsertMessageContent(MessageContent messageContent) { logger.debug(messageContent.getChannelId() + "/" + messageContent.getMessageId() + "/" + messageContent.getMetaDataId() + ": batch inserting message content (" + messageContent.getContentType().toString() + ")"); @@ -237,9 +243,10 @@ public void batchInsertMessageContent(MessageContent messageContent) { statement.addBatch(); statement.clearParameters(); } catch (SQLException e) { - throw new DonkeyDaoException(e); - } finally { + // The batch will never be executed now, so do not leave it for the next message. + clearBatchQuietly(statement); closeDatabaseObjectIfNeeded(statement); + throw new DonkeyDaoException(e); } } @@ -258,14 +265,29 @@ public void executeBatchInsertMessageContent(String channelId) { */ statement = prepareStatement("batchInsertMessageContent", channelId); statement.executeBatch(); - statement.clearBatch(); } catch (SQLException e) { throw new DonkeyDaoException(e); } finally { + clearBatchQuietly(statement); closeDatabaseObjectIfNeeded(statement); } } + /** + * Empties a cached statement's batch without letting the cleanup itself fail. The statement + * outlives the DAO in the prepared statement cache, so anything left on its batch would be + * executed along with the next message's rows. + */ + private void clearBatchQuietly(Statement statement) { + if (statement != null) { + try { + statement.clearBatch(); + } catch (SQLException e) { + logger.debug("Failed to clear batch", e); + } + } + } + @Override public void storeMessageContent(MessageContent messageContent) { logger.debug(messageContent.getChannelId() + "/" + messageContent.getMessageId() + "/" + messageContent.getMetaDataId() + ": updating message content (" + messageContent.getContentType().toString() + ")"); diff --git a/donkey/src/test/java/com/mirth/connect/donkey/server/data/jdbc/OracleJdbcDaoTest.java b/donkey/src/test/java/com/mirth/connect/donkey/server/data/jdbc/OracleJdbcDaoTest.java new file mode 100644 index 0000000000..bb0f8e0c69 --- /dev/null +++ b/donkey/src/test/java/com/mirth/connect/donkey/server/data/jdbc/OracleJdbcDaoTest.java @@ -0,0 +1,104 @@ +// SPDX-License-Identifier: MPL-2.0 +// SPDX-FileCopyrightText: Mitch Gaffigan + +package com.mirth.connect.donkey.server.data.jdbc; + +import static org.mockito.ArgumentMatchers.eq; +import static org.junit.Assert.fail; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.SQLException; + +import org.junit.Before; +import org.junit.Test; +import org.mockito.InOrder; + +import com.mirth.connect.donkey.model.message.ContentType; +import com.mirth.connect.donkey.model.message.MessageContent; +import com.mirth.connect.donkey.server.Donkey; +import com.mirth.connect.donkey.server.channel.Statistics; +import com.mirth.connect.donkey.server.data.DonkeyDaoException; +import com.mirth.connect.donkey.server.data.StatisticsUpdater; +import com.mirth.connect.donkey.util.SerializerProvider; + +/** + * Oracle is the only DAO that actually closes statements, so it is the only one where closing + * the wrong statement at the wrong time is observable. + */ +public class OracleJdbcDaoTest { + + private static final String CHANNEL_ID = "abc"; + + private OracleJdbcDao dao; + private PreparedStatement statement; + + @Before + public void before() throws SQLException { + Donkey donkey = mock(Donkey.class); + Connection connection = mock(Connection.class); + QuerySource querySource = mock(QuerySource.class); + PreparedStatementSource statementSource = mock(PreparedStatementSource.class); + SerializerProvider serializerProvider = mock(SerializerProvider.class); + StatisticsUpdater statisticsUpdater = mock(StatisticsUpdater.class); + Statistics currentStats = mock(Statistics.class); + Statistics totalStats = mock(Statistics.class); + + dao = spy(new OracleJdbcDao(donkey, connection, querySource, statementSource, serializerProvider, false, false, false, false, statisticsUpdater, currentStats, totalStats, "")); + + statement = mock(PreparedStatement.class); + doReturn(statement).when(dao).prepareStatement(eq("batchInsertMessageContent"), eq(CHANNEL_ID)); + } + + /** + * The batch lives on the cached statement, so the statement has to survive every + * batchInsertMessageContent() call and only be closed once the batch has been executed. + * Closing it earlier silently discarded the source content on Oracle. + */ + @Test + public void testBatchInsertMessageContentKeepsStatementOpenUntilExecuted() throws SQLException { + dao.batchInsertMessageContent(content(ContentType.PROCESSED_RAW, "processed raw")); + dao.batchInsertMessageContent(content(ContentType.TRANSFORMED, "transformed")); + dao.batchInsertMessageContent(content(ContentType.ENCODED, "encoded")); + + verify(statement, times(3)).addBatch(); + verify(statement, never()).close(); + + dao.executeBatchInsertMessageContent(CHANNEL_ID); + + InOrder inOrder = inOrder(statement); + inOrder.verify(statement, times(3)).addBatch(); + inOrder.verify(statement).executeBatch(); + inOrder.verify(statement).clearBatch(); + inOrder.verify(statement).close(); + } + + /** A failed batch must not be left behind for the next message to execute. */ + @Test + public void testFailedBatchInsertClearsAndClosesTheStatement() throws SQLException { + doThrow(new SQLException("no")).when(statement).addBatch(); + + try { + dao.batchInsertMessageContent(content(ContentType.ENCODED, "encoded")); + fail("Expected a DonkeyDaoException"); + } catch (DonkeyDaoException e) { + // expected + } + + InOrder inOrder = inOrder(statement); + inOrder.verify(statement).clearBatch(); + inOrder.verify(statement).close(); + } + + private static MessageContent content(ContentType contentType, String content) { + return new MessageContent(CHANNEL_ID, 1L, 0, contentType, content, "RAW", false); + } +} diff --git a/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java b/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java index f9849435ae..f2844f2ac6 100644 --- a/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java +++ b/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java @@ -55,6 +55,8 @@ static void assertFixtureFile(Message message, String fileName, String content) connector(message, SOURCE_META_DATA_ID, fileName).getStatus()); case "source_transformed" -> assertContent("source transformed", content, content(connector(message, SOURCE_META_DATA_ID, fileName).getTransformed())); + case "source_encoded" -> assertContent("source encoded", content, + content(connector(message, SOURCE_META_DATA_ID, fileName).getEncoded())); case "source_response" -> assertResponse("source response", content, connector(message, SOURCE_META_DATA_ID, fileName)); case "source_metadata.yml" -> assertMetadata("source_metadata.yml", parseYamlMap(content),