From 864ede756b0d99f061f8778ba209adf3f2600066 Mon Sep 17 00:00:00 2001 From: Mitch Gaffigan Date: Sat, 19 Sep 2026 18:48:01 -0500 Subject: [PATCH] Fix Oracle JDBC batch inserts by keeping statements open A previous change set prematurely called closeDatabaseObjectIfNeeded() after every addBatch() call. This prevented Oracle connections from accumulating and executing batch inserts properly, as the statement was closed before the batch could run. This patch removes the aggressive closure, allowing the PreparedStatement to remain open while the batch accumulates, and correctly closes it only after executeBatch() is called. It also introduces safe batch clearing on exceptions to prevent bleeding leftover batches into subsequent messages. Non-Oracle databases are not affected by this change since they do not override the closeDatabaseObjectIfNeeded method. Issue: https://github.com/OpenIntegrationEngine/engine/issues/454 Signed-off-by: Mitch Gaffigan Signed-off-by: Tony Germano --- ci/README.md | 1 + .../messages/01-hello-world/source_encoded | 1 + .../donkey/server/data/jdbc/JdbcDao.java | 28 ++++- .../server/data/jdbc/OracleJdbcDaoTest.java | 104 ++++++++++++++++++ .../smoketest/MessageAssertions.java | 2 + 5 files changed, 133 insertions(+), 3 deletions(-) create mode 100644 ci/tests/101-raw-no-op/channels/01-raw-no-op/messages/01-hello-world/source_encoded create mode 100644 donkey/src/test/java/com/mirth/connect/donkey/server/data/jdbc/OracleJdbcDaoTest.java 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),