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/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 + 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/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/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. * 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; + } +} 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. */