From 456848e6ed6aa65f0ffa0bd4032ef2447a817cb2 Mon Sep 17 00:00:00 2001 From: Mitch Gaffigan Date: Sat, 19 Sep 2026 11:01:42 -0500 Subject: [PATCH 1/4] Add server endpoint for submitting a batch message processMessage returns a single message id, taken from the response handler's selected result. For a batch that is only the first or last message, leaving a caller no way to learn the ids of the others. The smoke test harness needs them to assert per-message content. POST /channels/{channelId}/batchMessagesWithObj returns every id. The response handler is already threaded through dispatchBatchMessage, so CollectingResponseHandler only has to record each id as it is set. It keeps the ids and not the DispatchResults, so a long batch does not retain the processed messages, maps and content of the messages that have already finished. Processing is unchanged. Whether a batch is split is still governed by the channel's Process Batch setting; on a channel with batch processing disabled the dispatch takes the single-message path and the endpoint returns a one-element list. EngineController.dispatchRawMessage gains a five-argument overload. The existing four-argument form delegates to it with a SimpleResponseHandler, so its callers are unaffected, but the interface gains an abstract method and out-of-tree implementations will need updating. Signed-off-by: Mitch Gaffigan --- .../batch/CollectingResponseHandler.java | 32 +++++++++++++++++++ .../com/mirth/connect/client/core/Client.java | 10 ++++++ .../api/servlets/MessageServletInterface.java | 13 ++++++++ .../server/api/servlets/MessageServlet.java | 22 +++++++++++++ .../controllers/DonkeyEngineController.java | 7 +++- .../server/controllers/EngineController.java | 4 +++ 6 files changed, 87 insertions(+), 1 deletion(-) create mode 100644 donkey/src/main/java/com/mirth/connect/donkey/server/message/batch/CollectingResponseHandler.java diff --git a/donkey/src/main/java/com/mirth/connect/donkey/server/message/batch/CollectingResponseHandler.java b/donkey/src/main/java/com/mirth/connect/donkey/server/message/batch/CollectingResponseHandler.java new file mode 100644 index 0000000000..3f43aeeddd --- /dev/null +++ b/donkey/src/main/java/com/mirth/connect/donkey/server/message/batch/CollectingResponseHandler.java @@ -0,0 +1,32 @@ +// SPDX-License-Identifier: MPL-2.0 +// SPDX-FileCopyrightText: Mitch Gaffigan + +package com.mirth.connect.donkey.server.message.batch; + +import java.util.ArrayList; +import java.util.List; + +import com.mirth.connect.donkey.server.channel.DispatchResult; + +/** + * Reports every message of a batch, where {@link ResponseHandler#getResultForResponse()} reports + * only the one selected for the response. + */ +public class CollectingResponseHandler extends SimpleResponseHandler { + + private final List messageIds = new ArrayList(); + + @Override + public void setDispatchResult(DispatchResult dispatchResult) { + super.setDispatchResult(dispatchResult); + + if (dispatchResult != null) { + messageIds.add(dispatchResult.getMessageId()); + } + } + + public List getMessageIds() { + return messageIds; + } + +} diff --git a/server/src/main/java/com/mirth/connect/client/core/Client.java b/server/src/main/java/com/mirth/connect/client/core/Client.java index a927508d67..eaa261fc7e 100644 --- a/server/src/main/java/com/mirth/connect/client/core/Client.java +++ b/server/src/main/java/com/mirth/connect/client/core/Client.java @@ -1785,6 +1785,16 @@ public Long processMessage(String channelId, RawMessage rawMessage) throws Clien return getServlet(MessageServletInterface.class).processMessage(channelId, rawMessage); } + /** + * Processes a new message through a channel, returning the ID of every message it produced. + * + * @see MessageServletInterface#processBatchMessage + */ + @Override + public List processBatchMessage(String channelId, RawMessage rawMessage) throws ClientException { + return getServlet(MessageServletInterface.class).processBatchMessage(channelId, rawMessage); + } + /** * Processes a new message through a channel, using the RawMessage object. * diff --git a/server/src/main/java/com/mirth/connect/client/core/api/servlets/MessageServletInterface.java b/server/src/main/java/com/mirth/connect/client/core/api/servlets/MessageServletInterface.java index 4879f0c600..f8f88a552f 100644 --- a/server/src/main/java/com/mirth/connect/client/core/api/servlets/MessageServletInterface.java +++ b/server/src/main/java/com/mirth/connect/client/core/api/servlets/MessageServletInterface.java @@ -97,6 +97,19 @@ public Long processMessage(// @formatter:off @ExampleObject(name = "rawMessage", ref = "../apiexamples/raw_message_json") }) }) RawMessage rawMessage) throws ClientException; // @formatter:on + @POST + @Path("/{channelId}/batchMessagesWithObj") + @Operation(summary = "Processes a new message through a channel, returning the ID of every message it produced.") + @MirthOperation(name = "processMessages", display = "Process messages", permission = Permissions.MESSAGES_PROCESS, type = ExecuteType.ASYNC) + public List processBatchMessage(// @formatter:off + @Param("channelId") @Parameter(description = "The ID of the channel.", required = true) @PathParam("channelId") String channelId, + @Param("rawMessage") @RequestBody(description = "The RawMessage object to process.", required = true, content = { + @Content(mediaType = MediaType.APPLICATION_XML, examples = { + @ExampleObject(name = "rawMessage", ref = "../apiexamples/raw_message_xml") }), + @Content(mediaType = MediaType.APPLICATION_JSON, examples = { + @ExampleObject(name = "rawMessage", ref = "../apiexamples/raw_message_json") }) }) RawMessage rawMessage) throws ClientException; + // @formatter:on + @GET @Path("/{channelId}/messages/{messageId}") @Operation(summary = "Retrieve a message by ID.") diff --git a/server/src/main/java/com/mirth/connect/server/api/servlets/MessageServlet.java b/server/src/main/java/com/mirth/connect/server/api/servlets/MessageServlet.java index 9cecece5e9..0e3e778c67 100644 --- a/server/src/main/java/com/mirth/connect/server/api/servlets/MessageServlet.java +++ b/server/src/main/java/com/mirth/connect/server/api/servlets/MessageServlet.java @@ -45,6 +45,7 @@ import com.mirth.connect.donkey.server.channel.ChannelException; import com.mirth.connect.donkey.server.channel.DispatchResult; import com.mirth.connect.donkey.server.message.batch.BatchMessageException; +import com.mirth.connect.donkey.server.message.batch.CollectingResponseHandler; import com.mirth.connect.model.MessageImportResult; import com.mirth.connect.model.ServerEvent; import com.mirth.connect.model.ServerEvent.Level; @@ -133,6 +134,27 @@ public Long processMessage(final String channelId, final RawMessage rawMessage) return null; } + @Override + @CheckAuthorizedChannelId + public List processBatchMessage(final String channelId, final RawMessage rawMessage) { + CollectingResponseHandler responseHandler = new CollectingResponseHandler(); + + try { + engineController.dispatchRawMessage(channelId, rawMessage, true, true, responseHandler); + + containerRequestContext.setProperty(ResponseCodeFilter.RESPONSE_CODE_PROPERTY, Response.Status.CREATED.getStatusCode()); + return responseHandler.getMessageIds(); + } catch (ChannelException e) { + // Do nothing. An error should have been logged. + } catch (BatchMessageException e) { + logger.error("Error processing batch message for channel " + channelId, e); + } + + containerRequestContext.setProperty(ResponseCodeFilter.RESPONSE_CODE_PROPERTY, Response.Status.INTERNAL_SERVER_ERROR.getStatusCode()); + + return null; + } + @Override @CheckAuthorizedChannelId public Message getMessageContent(String channelId, Long messageId, List metaDataIds) { diff --git a/server/src/main/java/com/mirth/connect/server/controllers/DonkeyEngineController.java b/server/src/main/java/com/mirth/connect/server/controllers/DonkeyEngineController.java index 1dbeb66b03..d661850330 100644 --- a/server/src/main/java/com/mirth/connect/server/controllers/DonkeyEngineController.java +++ b/server/src/main/java/com/mirth/connect/server/controllers/DonkeyEngineController.java @@ -1119,6 +1119,11 @@ public Channel getDeployedChannel(String channelId) { @Override public DispatchResult dispatchRawMessage(String channelId, RawMessage rawMessage, boolean force, boolean canBatch) throws ChannelException, BatchMessageException { + return dispatchRawMessage(channelId, rawMessage, force, canBatch, new SimpleResponseHandler()); + } + + @Override + public DispatchResult dispatchRawMessage(String channelId, RawMessage rawMessage, boolean force, boolean canBatch, ResponseHandler responseHandler) throws ChannelException, BatchMessageException { if (!isDeployed(channelId)) { ChannelException e = new ChannelException(true); logger.error("Could not find channel to route to: " + channelId, e); @@ -1133,7 +1138,6 @@ public DispatchResult dispatchRawMessage(String channelId, RawMessage rawMessage } else { BatchRawMessage batchRawMessage = new BatchRawMessage(new BatchMessageReader(rawMessage.getRawData()), rawMessage.getSourceMap()); - ResponseHandler responseHandler = new SimpleResponseHandler(); sourceConnector.dispatchBatchMessage(batchRawMessage, responseHandler, rawMessage.getDestinationMetaDataIds()); return responseHandler.getResultForResponse(); @@ -1144,6 +1148,7 @@ public DispatchResult dispatchRawMessage(String channelId, RawMessage rawMessage try { dispatchResult = sourceConnector.dispatchRawMessage(rawMessage, force); dispatchResult.setAttemptedResponse(true); + responseHandler.setDispatchResult(dispatchResult); } finally { sourceConnector.finishDispatch(dispatchResult); } diff --git a/server/src/main/java/com/mirth/connect/server/controllers/EngineController.java b/server/src/main/java/com/mirth/connect/server/controllers/EngineController.java index f088643e21..8f32383e4b 100644 --- a/server/src/main/java/com/mirth/connect/server/controllers/EngineController.java +++ b/server/src/main/java/com/mirth/connect/server/controllers/EngineController.java @@ -22,6 +22,7 @@ import com.mirth.connect.donkey.model.channel.DebugOptions; import com.mirth.connect.donkey.server.channel.DispatchResult; import com.mirth.connect.donkey.server.message.batch.BatchMessageException; +import com.mirth.connect.donkey.server.message.batch.ResponseHandler; import com.mirth.connect.model.ChannelStatistics; import com.mirth.connect.model.DashboardStatus; import com.mirth.connect.model.ServerEventContext; @@ -73,6 +74,9 @@ public interface EngineController { public DispatchResult dispatchRawMessage(String channelId, RawMessage rawMessage, boolean force, boolean canBatch) throws ChannelException, BatchMessageException; + /** Dispatches a raw message, reporting each message it produced to {@code responseHandler}. */ + public DispatchResult dispatchRawMessage(String channelId, RawMessage rawMessage, boolean force, boolean canBatch, ResponseHandler responseHandler) throws ChannelException, BatchMessageException; + /** * Returns a list of DashboardStatus objects representing the running channels. * From df1e9368e974403dc37534faa6c760b1dbad7d83 Mon Sep 17 00:00:00 2001 From: Mitch Gaffigan Date: Sat, 19 Sep 2026 11:39:47 -0500 Subject: [PATCH 2/4] Support multi-message fixtures in the smoke test harness A batch fixture produces several messages, and the harness could only submit a payload and assert against one of them. Assertion files for a fixture that expects N messages now live in numbered subdirectories, 01 through NN; a fixture that expects one message keeps its flat layout and reads exactly as before. The generator rejects numbering that does not start at 01 or that leaves gaps, and rejects a fixture that mixes loose assertion files with numbered ones. An empty source_rejected file declares that the server must refuse the submission. It cannot be combined with assertion files, since a refused submission produces no message to assert against. source_raw, destNN_raw and destNN_encoded are now assertable. OieServer.submitMessage calls the new batchMessagesWithObj endpoint so the harness learns every id the payload produced, and runMessage fails if the count does not match what the fixture expects. Failure output renders every message rather than one, and now includes the source map. ci/tests/120-delimited-batch covers the new layout with a CSV batch channel, independent of any XML change. Note that routing every submission through batchMessagesWithObj leaves processMessage without end-to-end coverage. Signed-off-by: Mitch Gaffigan --- ci/README.md | 34 ++- .../channels/01-csv/channel.xml | 226 ++++++++++++++++++ .../01-csv/messages/01-batch/01/source_raw | 1 + .../messages/01-batch/01/source_transformed | 1 + .../01-csv/messages/01-batch/02/source_raw | 1 + .../messages/01-batch/02/source_transformed | 1 + .../channels/01-csv/messages/01-batch/source | 2 + smoketest/generate-smoke-tests.gradle | 46 +++- .../smoketest/Harness.java | 138 +++++++++-- .../smoketest/MessageAssertions.java | 6 +- .../smoketest/OieServer.java | 12 +- 11 files changed, 427 insertions(+), 41 deletions(-) create mode 100644 ci/tests/120-delimited-batch/channels/01-csv/channel.xml create mode 100644 ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_raw create mode 100644 ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_transformed create mode 100644 ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_raw create mode 100644 ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_transformed create mode 100644 ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/source diff --git a/ci/README.md b/ci/README.md index c9605360a5..1fd3f24e94 100644 --- a/ci/README.md +++ b/ci/README.md @@ -35,16 +35,26 @@ ci/tests/ ``` `channel.xml` is an exported OIE channel. Every message fixture requires `source`, -the payload sent to the channel. The other files are optional assertions: +the payload sent to the channel. These files describe the submission: -| File | Assertion | +| File | Meaning | | --- | --- | +| `source` | The payload sent to the channel. Required. | | `source_sourcemap.yml` | Source map supplied with `source` | +| `source_rejected` | The server must refuse the submission. | + +The rest are optional assertions against the message the channel produced: + +| File | Assertion | +| --- | --- | | `source_metadata.yml` | Selected source message metadata | +| `source_raw` | Raw source payload | +| `source_encoded` | Encoded source payload | | `source_status` | Source status | | `source_response` | Source response payload | | `source_transformed` | Transformed source payload | -| `source_encoded` | Encoded source payload | +| `destNN_raw` | Raw payload for destination `NN` | +| `destNN_encoded` | Encoded payload for destination `NN` | | `destNN` | Sent payload for destination `NN` | | `destNN_transformed` | Transformed payload for destination `NN` | | `destNN_response` | Response payload from destination `NN` | @@ -64,6 +74,24 @@ alpine-temurin21-postgres ubuntu-temurin21-postgres ``` +### Parse Batch + +If you are using "Parse Batch" mode, put each expected message's assertion +files in a numbered subdirectory instead. There must be one directory per +expected message. + +```text +messages/ + 03-two-messages/ + source + 01/ + source_raw + source_status + 02/ + source_raw + source_status +``` + ## Add a Configuration Add a Compose file at `ci/configurations/.compose.yml`. Configuration names diff --git a/ci/tests/120-delimited-batch/channels/01-csv/channel.xml b/ci/tests/120-delimited-batch/channels/01-csv/channel.xml new file mode 100644 index 0000000000..9e6f3cf7da --- /dev/null +++ b/ci/tests/120-delimited-batch/channels/01-csv/channel.xml @@ -0,0 +1,226 @@ + + 732e3b46-e151-47e8-8699-d10d088e99db + 2 + CSV batch + + 1 + + 0 + sourceConnector + + + + None + true + true + false + 1 + + + Default Resource + [Default Resource] + + + 1000 + + + + + DELIMITED + XML + + + , + \n + " + true + \ + + Name + Description + + false + true + + + , + \n + " + true + \ + + + Record + 0 + + false + + + + + + + false + + + Element_Name + + 1 + + + + + + + + + Channel Reader + SOURCE + true + true + + + + 1 + Destination 1 + + + + false + false + 10000 + false + 0 + false + false + 1 + + false + + + Default Resource + [Default Resource] + + + 1000 + true + + none + ${message.encodedData} + + + + + XML + XML + + + false + + + Element_Name + + 1 + + + + + + + false + + + Element_Name + + 1 + + + + + + + + RAW + RAW + + + JavaScript + + + + + + JavaScript + + + + + + + + Channel Writer + DESTINATION + true + true + + + // Modify the message variable below to pre process data +return message; + // This script executes once after a message has been processed +// Responses returned from here will be stored as "Postprocessor" in the response map +return; + // This script executes once when the channel is deployed +// You only have access to the globalMap and globalChannelMap here to persist data +return; + // This script executes once when the channel is undeployed +// You only have access to the globalMap and globalChannelMap here to persist data +return; + + true + DEVELOPMENT + false + false + false + false + false + false + STARTED + true + + + SOURCE + STRING + mirth_source + + + TYPE + STRING + mirth_type + + + + None + + + + + Default Resource + [Default Resource] + + + + + + true + + + America/Chicago + + + true + false + + 1 + + + \ No newline at end of file diff --git a/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_raw b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_raw new file mode 100644 index 0000000000..7e278783dd --- /dev/null +++ b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_raw @@ -0,0 +1 @@ +OIE,Awesome diff --git a/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_transformed b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_transformed new file mode 100644 index 0000000000..4095ebbce6 --- /dev/null +++ b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/01/source_transformed @@ -0,0 +1 @@ +OIEAwesome \ No newline at end of file diff --git a/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_raw b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_raw new file mode 100644 index 0000000000..2813f675f2 --- /dev/null +++ b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_raw @@ -0,0 +1 @@ +Other,Not as awesome diff --git a/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_transformed b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_transformed new file mode 100644 index 0000000000..2cbea13928 --- /dev/null +++ b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/02/source_transformed @@ -0,0 +1 @@ +OtherNot as awesome \ No newline at end of file diff --git a/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/source b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/source new file mode 100644 index 0000000000..25fc94e553 --- /dev/null +++ b/ci/tests/120-delimited-batch/channels/01-csv/messages/01-batch/source @@ -0,0 +1,2 @@ +OIE,Awesome +Other,Not as awesome diff --git a/smoketest/generate-smoke-tests.gradle b/smoketest/generate-smoke-tests.gradle index 7ee6a2879f..7184ddfdca 100644 --- a/smoketest/generate-smoke-tests.gradle +++ b/smoketest/generate-smoke-tests.gradle @@ -117,33 +117,67 @@ abstract class GenerateSmokeTests extends DefaultTask { out << " }\n" } + /** Files that describe the submission rather than the messages it produces. */ + private static final List INPUT_FILES = ['source', 'source_sourcemap.yml', 'source_rejected'] + void writeMessage(StringBuilder out, File root, File resOut, File messageDir, int order) { File source = new File(messageDir, 'source') if (!source.isFile()) { throw new GradleException("Message fixture is missing its source payload: ${messageDir}") } boolean hasSourceMap = new File(messageDir, 'source_sourcemap.yml').isFile() + boolean rejected = new File(messageDir, 'source_rejected').isFile() String base = 'fixtures/' + rel(root, messageDir) + // Naming files relative to the message directory yields "source_status" for a fixture that + // produces one message and "02/source_status" for one that produces several. List assertionFiles = [] - messageDir.listFiles().findAll { it.isFile() }.sort().each { File file -> + messageDir.eachFileRecurse(groovy.io.FileType.FILES) { File file -> + if (file.name.startsWith('.')) return copyResource(file, new File(resOut, 'fixtures/' + rel(root, file))) - if (file.name != 'source' && file.name != 'source_sourcemap.yml') { - assertionFiles << file.name + String path = rel(messageDir, file) + if (!INPUT_FILES.contains(path)) { + assertionFiles << path } } + assertionFiles.sort() - StringBuilder args = new StringBuilder("channelId, ${quote(base)}, ${hasSourceMap}") - assertionFiles.each { args << ", ${quote(it)}" } + int messageCount = expectedMessageCount(messageDir, assertionFiles, rejected) out << "\n @Test\n" out << " @Order(${order})\n" out << " @DisplayName(${quote(messageDir.name)})\n" out << " void message_${ident(messageDir.name)}() throws Exception {\n" - out << " Harness.runMessage(${args});\n" + if (rejected) { + out << " Harness.runRejectedMessage(channelId, ${quote(base)}, ${hasSourceMap});\n" + } else { + StringBuilder args = new StringBuilder("channelId, ${quote(base)}, ${hasSourceMap}, ${messageCount}") + assertionFiles.each { args << ", ${quote(it)}" } + out << " Harness.runMessage(${args});\n" + } out << " }\n" } + /** One message, unless the fixture has a numbered directory per message it expects. */ + static int expectedMessageCount(File messageDir, List assertionFiles, boolean rejected) { + if (rejected && !assertionFiles.isEmpty()) { + throw new GradleException("source_rejected means no message is created, so it cannot be combined with ${assertionFiles} in ${messageDir}") + } + + List slots = messageDir.listFiles().findAll { it.isDirectory() }*.name.sort() + if (slots.any { !(it ==~ /\d+/) }) { + throw new GradleException("Subdirectories of a message fixture must be message numbers, found ${slots} in ${messageDir}") + } + if (slots.indexed().any { index, slot -> Integer.parseInt(slot) != index + 1 }) { + throw new GradleException("Numbered message directories must run from 01 with no gaps, found ${slots} in ${messageDir}") + } + if (!slots.isEmpty() && assertionFiles.any { !it.contains('/') }) { + throw new GradleException("Move the loose assertion files into a numbered message directory: ${messageDir}") + } + + return slots.isEmpty() ? 1 : slots.size() + } + static List readConfigurations(File caseDir) { File file = new File(caseDir, 'configurations') return file.isFile() ? file.readLines('UTF-8')*.trim().findAll { it && !it.startsWith('#') } : [] diff --git a/smoketest/src/test/java/org/openintegrationengine/smoketest/Harness.java b/smoketest/src/test/java/org/openintegrationengine/smoketest/Harness.java index 5a2b702819..386c392299 100644 --- a/smoketest/src/test/java/org/openintegrationengine/smoketest/Harness.java +++ b/smoketest/src/test/java/org/openintegrationengine/smoketest/Harness.java @@ -7,6 +7,7 @@ import java.io.InputStream; import java.io.UncheckedIOException; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.Arrays; import java.util.LinkedHashMap; import java.util.List; @@ -14,6 +15,7 @@ import org.junit.jupiter.api.Assumptions; +import com.mirth.connect.client.core.ClientException; import com.mirth.connect.donkey.model.message.ConnectorMessage; import com.mirth.connect.donkey.model.message.Message; import com.mirth.connect.donkey.model.message.MessageContent; @@ -63,40 +65,48 @@ public static void undeploy(String channelId) { /** * Submits {@code /source} (with {@code /source_sourcemap.yml} when * {@code hasSourceMap}) into the channel, then retries the named assertion files until - * they all hold or the message reaches a terminal state. Because the message is written + * they all hold or every message reaches a terminal state. Because messages are written * asynchronously, an early poll can legitimately fail; only a failure that persists once - * the message is terminal is a real failure. + * the messages are terminal is a real failure. + * + * @param assertionFiles a bare file name for a payload that produces one message, or one + * prefixed with the message's 1-based number when it produces several, + * e.g. {@code "02/source_status"}. */ - public static void runMessage(String channelId, String base, boolean hasSourceMap, String... assertionFiles) - throws Exception { + public static void runMessage(String channelId, String base, boolean hasSourceMap, int expectedMessageCount, + String... assertionFiles) throws Exception { String source = resource(base + "/source"); - Map sourceMap = hasSourceMap - ? MessageAssertions.parseSourceMap(resource(base + "/source_sourcemap.yml")) - : new LinkedHashMap<>(); + Map sourceMap = sourceMap(base, hasSourceMap); // Load the fixtures once; the poll loop below may check them many times. - Map assertions = new LinkedHashMap<>(); - for (String fileName : assertionFiles) { - assertions.put(fileName, resource(base + "/" + fileName)); + List> assertions = loadAssertions(base, expectedMessageCount, assertionFiles); + + List messageIds; + try { + messageIds = server().submitMessage(channelId, source, sourceMap); + } catch (ClientException e) { + throw new AssertionError(base + " failed: the server refused the payload. Add source_rejected" + + " to the fixture if that is expected. " + e.getMessage(), e); } - long messageId = server().submitMessage(channelId, source, sourceMap); + if (messageIds.size() != expectedMessageCount) { + throw new AssertionError(base + " failed: expected the payload to produce " + expectedMessageCount + + " message(s), found " + messageIds.size() + " " + messageIds); + } long deadline = System.nanoTime() + HarnessConfig.TIMEOUT.toNanos(); AssertionError lastFailure = null; - Message lastMessage = null; + List lastMessages = List.of(); while (System.nanoTime() < deadline) { - Message message = server().fetchMessage(channelId, messageId); - if (message != null) { - lastMessage = message; + List messages = fetchMessages(channelId, messageIds); + lastMessages = messages; + if (!messages.contains(null)) { try { - for (Map.Entry assertion : assertions.entrySet()) { - MessageAssertions.assertFixtureFile(message, assertion.getKey(), assertion.getValue()); - } + assertAll(messages, assertions); return; } catch (AssertionError e) { lastFailure = e; - if (isTerminal(message)) { + if (messages.stream().allMatch(Harness::isTerminal)) { break; } } @@ -106,10 +116,30 @@ public static void runMessage(String channelId, String base, boolean hasSourceMa if (lastFailure != null) { throw new AssertionError(base + " failed: " + lastFailure.getMessage() - + "\n\n" + describe(lastMessage), lastFailure); + + "\n\n" + describe(lastMessages), lastFailure); } - throw new AssertionError("Timed out after " + HarnessConfig.TIMEOUT.toSeconds() + "s waiting for message " - + messageId + " for fixture " + base + "\n\n" + describe(lastMessage)); + throw new AssertionError("Timed out after " + HarnessConfig.TIMEOUT.toSeconds() + "s waiting for message(s) " + + messageIds + " for fixture " + base + "\n\n" + describe(lastMessages)); + } + + /** + * Submits {@code /source} and requires the server to refuse it. A refusal reaches the + * client only as a {@link ClientException} carrying the status line as text, so the refusal + * itself is the assertion. + */ + public static void runRejectedMessage(String channelId, String base, boolean hasSourceMap) throws Exception { + String source = resource(base + "/source"); + Map sourceMap = sourceMap(base, hasSourceMap); + + List messageIds; + try { + messageIds = server().submitMessage(channelId, source, sourceMap); + } catch (ClientException refused) { + return; + } + + throw new AssertionError(base + " failed: expected the server to reject the payload, but it was accepted" + + " and produced message(s) " + messageIds + "\n\n" + describe(fetchMessages(channelId, messageIds))); } /** Reads a staged fixture from the classpath. */ @@ -124,6 +154,45 @@ static String resource(String path) { } } + private static Map sourceMap(String base, boolean hasSourceMap) { + return hasSourceMap + ? MessageAssertions.parseSourceMap(resource(base + "/source_sourcemap.yml")) + : new LinkedHashMap<>(); + } + + /** Groups the assertion files by the message they describe, keyed by their bare file name. */ + private static List> loadAssertions(String base, int expectedMessageCount, + String[] assertionFiles) { + List> assertions = new ArrayList<>(); + for (int index = 0; index < expectedMessageCount; index++) { + assertions.add(new LinkedHashMap<>()); + } + + for (String path : assertionFiles) { + int separator = path.lastIndexOf('/'); + int index = separator < 0 ? 0 : Integer.parseInt(path.substring(0, separator)) - 1; + assertions.get(index).put(path.substring(separator + 1), resource(base + "/" + path)); + } + return assertions; + } + + private static void assertAll(List messages, List> assertions) { + for (int index = 0; index < messages.size(); index++) { + for (Map.Entry assertion : assertions.get(index).entrySet()) { + MessageAssertions.assertFixtureFile(messages.get(index), assertion.getKey(), assertion.getValue()); + } + } + } + + /** Reads each message back, leaving a null in place of one the server has not stored yet. */ + private static List fetchMessages(String channelId, List messageIds) throws ClientException { + List messages = new ArrayList<>(messageIds.size()); + for (Long messageId : messageIds) { + messages.add(server().fetchMessage(channelId, messageId)); + } + return messages; + } + /** True once the server has finished processing and no connector is still pending. */ private static boolean isTerminal(Message message) { if (!message.isProcessed()) { @@ -137,8 +206,26 @@ private static boolean isTerminal(Message message) { .noneMatch(connectorMessage -> PENDING_STATUSES.contains(connectorMessage.getStatus())); } - /** Renders the message the way a fixture author needs to see it to fix a mismatch. */ - private static String describe(Message message) { + /** Renders the messages the way a fixture author needs to see them to fix a mismatch. */ + private static String describe(List messages) { + if (messages.isEmpty()) { + return "No messages were retrieved from the server."; + } + + StringBuilder detail = new StringBuilder(); + for (int index = 0; index < messages.size(); index++) { + if (index > 0) { + detail.append("\n\n"); + } + if (messages.size() > 1) { + detail.append("--- message ").append(index + 1).append(" of ").append(messages.size()).append(" ---\n"); + } + detail.append(describeMessage(messages.get(index))); + } + return detail.toString(); + } + + private static String describeMessage(Message message) { if (message == null) { return "No message was retrieved from the server."; } @@ -159,7 +246,8 @@ private static String describe(Message message) { appendContent(detail, "encoded", connectorMessage.getEncoded()); appendContent(detail, "sent", connectorMessage.getSent()); appendContent(detail, "response", connectorMessage.getResponse()); - detail.append("\n connectorMap=").append(connectorMessage.getConnectorMap()) + detail.append("\n sourceMap=").append(connectorMessage.getSourceMap()) + .append("\n connectorMap=").append(connectorMessage.getConnectorMap()) .append("\n metaDataMap=").append(connectorMessage.getMetaDataMap()); if (connectorMessage.getProcessingError() != null) { detail.append("\n processingError=").append(connectorMessage.getProcessingError()); diff --git a/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java b/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java index f2844f2ac6..b99ca9a283 100644 --- a/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java +++ b/smoketest/src/test/java/org/openintegrationengine/smoketest/MessageAssertions.java @@ -35,7 +35,7 @@ final class MessageAssertions { private static final Pattern RESPONSE_ENVELOPE = Pattern.compile("^\\s*].*", Pattern.DOTALL); /** Destination assertion files are {@code dest} plus an optional suffix. */ - private static final Pattern DEST_NAME = Pattern.compile("dest(\\d+)(_transformed|_response|_status|_metadata\\.yml)?"); + private static final Pattern DEST_NAME = Pattern.compile("dest(\\d+)(_raw|_transformed|_encoded|_response|_status|_metadata\\.yml)?"); /** Source connector metadata id; destination N is metadata id N. */ private static final int SOURCE_META_DATA_ID = 0; @@ -53,6 +53,8 @@ static void assertFixtureFile(Message message, String fileName, String content) switch (fileName) { case "source_status" -> assertStatus("source status", content, connector(message, SOURCE_META_DATA_ID, fileName).getStatus()); + case "source_raw" -> assertContent("source raw", content, + content(connector(message, SOURCE_META_DATA_ID, fileName).getRaw())); case "source_transformed" -> assertContent("source transformed", content, content(connector(message, SOURCE_META_DATA_ID, fileName).getTransformed())); case "source_encoded" -> assertContent("source encoded", content, @@ -77,7 +79,9 @@ private static void assertDestination(Message message, String fileName, String c switch (suffix) { case "" -> assertContent(fileName, content, content(destination.getSent())); + case "_raw" -> assertContent(fileName, content, content(destination.getRaw())); case "_transformed" -> assertContent(fileName, content, content(destination.getTransformed())); + case "_encoded" -> assertContent(fileName, content, content(destination.getEncoded())); case "_response" -> assertResponse(fileName, content, destination); case "_status" -> assertStatus(fileName, content, destination.getStatus()); case "_metadata.yml" -> assertMetadata(fileName, parseYamlMap(content), destination); diff --git a/smoketest/src/test/java/org/openintegrationengine/smoketest/OieServer.java b/smoketest/src/test/java/org/openintegrationengine/smoketest/OieServer.java index 9df99d794b..7fe6d79868 100644 --- a/smoketest/src/test/java/org/openintegrationengine/smoketest/OieServer.java +++ b/smoketest/src/test/java/org/openintegrationengine/smoketest/OieServer.java @@ -130,14 +130,14 @@ private void awaitStarted(String channelId, String label) throws Exception { + HarnessConfig.TIMEOUT.toSeconds() + "s; last state was " + lastState); } - /** Submits a source payload and returns the new message id. */ - long submitMessage(String channelId, String rawData, Map sourceMap) throws ClientException { + /** Submits a source payload and returns the id of every message it produced, in dispatch order. */ + List submitMessage(String channelId, String rawData, Map sourceMap) throws ClientException { RawMessage rawMessage = new RawMessage(rawData, null, sourceMap); - Long messageId = client.processMessage(channelId, rawMessage); - if (messageId == null) { - throw new AssertionError("Server returned no message id for channel " + channelId); + List messageIds = client.processBatchMessage(channelId, rawMessage); + if (messageIds == null) { + throw new AssertionError("Server returned no message ids for channel " + channelId); } - return messageId; + return messageIds; } /** Reads one message back, with content, so assertions can inspect every connector. */ From 26fbdc5151a09e2ebd34d55c771e4be7e05c636f Mon Sep 17 00:00:00 2001 From: Mitch Gaffigan Date: Sat, 19 Sep 2026 11:40:50 -0500 Subject: [PATCH 3/4] Add XXE test during xml parse batch (CVE-2026-82578) Signed-off-by: Mitch Gaffigan --- .../channels/01-xml-batch-xxe/channel.xml | 220 ++++++++++++++++++ .../01-xml-batch-xxe/messages/01-xxe/source | 5 + .../messages/01-xxe/source_rejected | 0 .../messages/02-two-messages/01/source_raw | 1 + .../messages/02-two-messages/02/source_raw | 1 + .../messages/02-two-messages/source | 5 + 6 files changed, 232 insertions(+) create mode 100644 ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/channel.xml create mode 100644 ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source create mode 100644 ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source_rejected create mode 100644 ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/01/source_raw create mode 100644 ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/02/source_raw create mode 100644 ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/source diff --git a/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/channel.xml b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/channel.xml new file mode 100644 index 0000000000..4987ee9fa5 --- /dev/null +++ b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/channel.xml @@ -0,0 +1,220 @@ + + da8f77be-aee0-41a6-bbc7-ad12a45ec293 + 2 + XML Batch XXE + + 3 + + 0 + sourceConnector + + + + None + true + true + false + 1 + + + Default Resource + [Default Resource] + + + 1000 + + + + + XML + XML + + + false + + + XPath_Query + + 1 + /Batch/Message + + + + + + false + + + Element_Name + + 1 + + + + + + + + + Channel Reader + SOURCE + true + true + + + + 1 + Destination 1 + + + + false + false + 10000 + false + 0 + false + false + 1 + + false + + + Default Resource + [Default Resource] + + + 1000 + true + + none + ${message.encodedData} + + + + + XML + XML + + + false + + + Element_Name + + 1 + + + + + + + false + + + Element_Name + + 1 + + + + + + + + XML + XML + + + false + + + Element_Name + + 1 + + + + + + + false + + + Element_Name + + 1 + + + + + + + + + Channel Writer + DESTINATION + true + true + + + // Modify the message variable below to pre process data +return message; + // This script executes once after a message has been processed +// Responses returned from here will be stored as "Postprocessor" in the response map +return; + // This script executes once when the channel is deployed +// You only have access to the globalMap and globalChannelMap here to persist data +return; + // This script executes once when the channel is undeployed +// You only have access to the globalMap and globalChannelMap here to persist data +return; + + true + DEVELOPMENT + false + false + false + false + false + false + STARTED + true + + + SOURCE + STRING + mirth_source + + + TYPE + STRING + mirth_type + + + + None + + + + + Default Resource + [Default Resource] + + + + + + true + + + America/Chicago + + + true + false + + 1 + + + \ No newline at end of file diff --git a/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source new file mode 100644 index 0000000000..6c79e522ab --- /dev/null +++ b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source @@ -0,0 +1,5 @@ + + ]> + + &xxe; + \ No newline at end of file diff --git a/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source_rejected b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/01-xxe/source_rejected new file mode 100644 index 0000000000..e69de29bb2 diff --git a/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/01/source_raw b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/01/source_raw new file mode 100644 index 0000000000..9df09439db --- /dev/null +++ b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/01/source_raw @@ -0,0 +1 @@ +alpha diff --git a/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/02/source_raw b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/02/source_raw new file mode 100644 index 0000000000..182fdd49b7 --- /dev/null +++ b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/02/source_raw @@ -0,0 +1 @@ +bravo diff --git a/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/source b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/source new file mode 100644 index 0000000000..4be910e0ff --- /dev/null +++ b/ci/tests/210-xml-batch-xxe/channels/01-xml-batch-xxe/messages/02-two-messages/source @@ -0,0 +1,5 @@ + + + alpha + bravo + From 77e86a5e2006c8ab8245636a41fb90475b519a1a Mon Sep 17 00:00:00 2001 From: Mitch Gaffigan Date: Sat, 19 Sep 2026 12:16:17 -0500 Subject: [PATCH 4/4] Fix XXE in XML batch parsing CVE-2026-82578 found an XXE vulnerability in XML batch parsing. This closes that vulnerability and adds regression tests. setNamespaceAware(true) is not a behavior change. XPath.evaluate(String, InputSource, QName) built its DOM with a namespace aware DocumentBuilder, so the batch splitter has always been namespace aware. Without it, split messages lose the xmlns declarations they inherit from the batch element. Signed-off-by: Mitch Gaffigan --- .../datatypes/xml/XMLBatchAdaptor.java | 10 ++- .../datatypes/xml/XMLBatchAdaptorTest.java | 69 +++++++++++++++++++ 2 files changed, 78 insertions(+), 1 deletion(-) create mode 100644 server/src/test/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptorTest.java diff --git a/server/src/main/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptor.java b/server/src/main/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptor.java index bb6994e339..5e4ac94874 100644 --- a/server/src/main/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptor.java +++ b/server/src/main/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptor.java @@ -18,6 +18,7 @@ import java.util.Map; import javax.xml.XMLConstants; +import javax.xml.parsers.DocumentBuilderFactory; import javax.xml.transform.OutputKeys; import javax.xml.transform.Transformer; import javax.xml.transform.TransformerFactory; @@ -33,6 +34,7 @@ import org.mozilla.javascript.Context; import org.mozilla.javascript.Script; import org.mozilla.javascript.Scriptable; +import org.w3c.dom.Document; import org.w3c.dom.Node; import org.w3c.dom.NodeList; import org.xml.sax.InputSource; @@ -42,6 +44,7 @@ import com.mirth.connect.donkey.server.message.batch.BatchMessageException; import com.mirth.connect.donkey.server.message.batch.BatchMessageReader; import com.mirth.connect.donkey.server.message.batch.BatchMessageReceiver; +import com.mirth.connect.model.converters.DocumentSerializer; import com.mirth.connect.plugins.datatypes.xml.XMLBatchProperties.SplitType; import com.mirth.connect.server.controllers.ContextFactoryController; import com.mirth.connect.server.controllers.ControllerFactory; @@ -127,7 +130,12 @@ private String getMessageFromReader() throws Exception { XPath xpath = xPathFactory.newXPath(); - nodeList = (NodeList) xpath.evaluate(query.toString(), new InputSource(bufferedReader), XPathConstants.NODESET); + // Parse the XML securely to prevent XXE + DocumentBuilderFactory documentBuilderFactory = DocumentSerializer.getSecureDocumentBuilderFactory(); + documentBuilderFactory.setNamespaceAware(true); + Document document = documentBuilderFactory.newDocumentBuilder().parse(new InputSource(bufferedReader)); + + nodeList = (NodeList) xpath.evaluate(query.toString(), document, XPathConstants.NODESET); } if (currentNode < nodeList.getLength()) { diff --git a/server/src/test/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptorTest.java b/server/src/test/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptorTest.java new file mode 100644 index 0000000000..5932f080dd --- /dev/null +++ b/server/src/test/java/com/mirth/connect/plugins/datatypes/xml/XMLBatchAdaptorTest.java @@ -0,0 +1,69 @@ +// SPDX-License-Identifier: MPL-2.0 +// SPDX-FileCopyrightText: Mitch Gaffigan +package com.mirth.connect.plugins.datatypes.xml; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +import java.io.File; + +import org.apache.commons.io.FileUtils; +import org.junit.Test; +import org.xml.sax.SAXParseException; + +import com.mirth.connect.donkey.model.message.BatchRawMessage; +import com.mirth.connect.donkey.server.message.batch.BatchMessageException; +import com.mirth.connect.donkey.server.message.batch.BatchMessageReader; +import com.mirth.connect.plugins.datatypes.xml.XMLBatchProperties.SplitType; + +/** Covered end to end by ci/tests/210-xml-batch-xxe. */ +public class XMLBatchAdaptorTest { + + @Test + public void externalEntityIsNotResolved() throws Exception { + File secret = File.createTempFile("xxe", ".txt"); + secret.deleteOnExit(); + FileUtils.write(secret, "canary", UTF_8); + + XMLBatchAdaptor adaptor = adaptor(" ]>&xxe;"); + + try { + fail("Expected the external entity to be rejected, but got: " + adaptor.getMessage()); + } catch (BatchMessageException e) { + assertTrue(String.valueOf(e.getCause()), e.getCause() instanceof SAXParseException); + assertFalse(String.valueOf(e), String.valueOf(e).contains("canary")); + } + } + + @Test + public void batchSplitsOnElementName() throws Exception { + XMLBatchAdaptor adaptor = adaptor("alphabravo"); + + assertEquals("alpha", adaptor.getMessage().trim()); + assertEquals("bravo", adaptor.getMessage().trim()); + assertNull(adaptor.getMessage()); + } + + /** Splitting is namespace aware, as it was when XPath parsed the InputSource itself. */ + @Test + public void batchSplitKeepsNamespaceDeclarations() throws Exception { + XMLBatchAdaptor adaptor = adaptor("alpha"); + + assertEquals("alpha", adaptor.getMessage().trim()); + assertNull(adaptor.getMessage()); + } + + private static XMLBatchAdaptor adaptor(String xml) { + XMLBatchProperties properties = new XMLBatchProperties(); + properties.setSplitType(SplitType.Element_Name); + properties.setElementName("Message"); + + XMLBatchAdaptor adaptor = new XMLBatchAdaptor(null, null, new BatchRawMessage(new BatchMessageReader(xml))); + adaptor.setBatchProperties(properties); + return adaptor; + } +}