This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 4e3f1878c315 CAMEL-25019: camel-quickfix - send InOut replies on the
session the request arrived on
4e3f1878c315 is described below
commit 4e3f1878c315b8788786dd96fb7487e32867efcc
Author: Andrea Cosentino <[email protected]>
AuthorDate: Mon Sep 28 20:12:07 2026 +0200
CAMEL-25019: camel-quickfix - send InOut replies on the session the request
arrived on
With exchangePattern=InOut the QuickFIX/J consumer read the SessionID
header when the route completed, so a route that changed that header also
changed which session received the reply. The reply is now sent on the
session the request was received on. A route that wants to send a message
to another session must use a QuickFIX/J producer endpoint with the
sessionID option; the upgrade guide describes it.
Closes #26886
Co-authored-by: Claude Opus 5.5 (1M context) <[email protected]>
---
.../camel/catalog/docs/quickfix-component.adoc | 20 +++--
.../src/main/docs/quickfix-component.adoc | 20 +++--
.../component/quickfixj/QuickfixjConsumer.java | 10 ++-
.../component/quickfixj/QuickfixjConsumerTest.java | 92 ++++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 7 ++
5 files changed, 135 insertions(+), 14 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/quickfix-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/quickfix-component.adoc
index b8e2a6d6737c..aca3b4f40c1c 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/quickfix-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/quickfix-component.adoc
@@ -91,11 +91,16 @@ include::partial$component-endpoint-headers.adoc[]
== Usage
-The DataDictionary header is useful if string messages are being
-received and need to be parsed in a route. QuickFIX/J requires a data
-dictionary to parse certain types of messages (with repeating groups,
-for example). By injecting a DataDictionary header in the route after
-receiving a message string, the FIX engine can properly parse the data.
+The `DataDictionary` exchange property is useful if string messages are
+being received and need to be parsed in a route. QuickFIX/J requires a
+data dictionary to parse certain types of messages (with repeating
+groups, for example). By setting the `DataDictionary` exchange property
+(`QuickfixjEndpoint.DATA_DICTIONARY_KEY`) in the route after receiving a
+message string, the FIX engine can properly parse the data. The property
+can hold a `quickfix.DataDictionary` instance or the name of a data
+dictionary resource, such as `FIX44.xml`. When the property is not set,
+the data dictionary of the session identified by the `SessionID` header
+is used, if that session exists.
=== QuickFIX/J Configuration Extensions
@@ -246,6 +251,11 @@ to the requestor session.
.bean(new MarketOrderStatusService());
----
+The reply is always sent on the session the request was received on.
+Changing the `SessionID` header while routing does not change where the
+reply goes. To send a message to another session, use a QuickFIX/J
+producer endpoint with the `sessionID` option.
+
==== Implementing InOut Exchanges for Producers
For producers, sending a message will block until a reply is received or
diff --git a/components/camel-quickfix/src/main/docs/quickfix-component.adoc
b/components/camel-quickfix/src/main/docs/quickfix-component.adoc
index b8e2a6d6737c..aca3b4f40c1c 100644
--- a/components/camel-quickfix/src/main/docs/quickfix-component.adoc
+++ b/components/camel-quickfix/src/main/docs/quickfix-component.adoc
@@ -91,11 +91,16 @@ include::partial$component-endpoint-headers.adoc[]
== Usage
-The DataDictionary header is useful if string messages are being
-received and need to be parsed in a route. QuickFIX/J requires a data
-dictionary to parse certain types of messages (with repeating groups,
-for example). By injecting a DataDictionary header in the route after
-receiving a message string, the FIX engine can properly parse the data.
+The `DataDictionary` exchange property is useful if string messages are
+being received and need to be parsed in a route. QuickFIX/J requires a
+data dictionary to parse certain types of messages (with repeating
+groups, for example). By setting the `DataDictionary` exchange property
+(`QuickfixjEndpoint.DATA_DICTIONARY_KEY`) in the route after receiving a
+message string, the FIX engine can properly parse the data. The property
+can hold a `quickfix.DataDictionary` instance or the name of a data
+dictionary resource, such as `FIX44.xml`. When the property is not set,
+the data dictionary of the session identified by the `SessionID` header
+is used, if that session exists.
=== QuickFIX/J Configuration Extensions
@@ -246,6 +251,11 @@ to the requestor session.
.bean(new MarketOrderStatusService());
----
+The reply is always sent on the session the request was received on.
+Changing the `SessionID` header while routing does not change where the
+reply goes. To send a message to another session, use a QuickFIX/J
+producer endpoint with the `sessionID` option.
+
==== Implementing InOut Exchanges for Producers
For producers, sending a message will block until a reply is received or
diff --git
a/components/camel-quickfix/src/main/java/org/apache/camel/component/quickfixj/QuickfixjConsumer.java
b/components/camel-quickfix/src/main/java/org/apache/camel/component/quickfixj/QuickfixjConsumer.java
index bfbb90657365..7f60ad11da4c 100644
---
a/components/camel-quickfix/src/main/java/org/apache/camel/component/quickfixj/QuickfixjConsumer.java
+++
b/components/camel-quickfix/src/main/java/org/apache/camel/component/quickfixj/QuickfixjConsumer.java
@@ -56,10 +56,14 @@ public class QuickfixjConsumer extends DefaultConsumer {
public void onExchange(Exchange exchange) {
if (isStarted()) {
try {
+ // the reply goes back on the session the request arrived on,
so capture it before routing
+ // can change the header
+ SessionID messageSessionID =
exchange.getIn().getHeader(QuickfixjEndpoint.SESSION_ID_KEY, SessionID.class);
+
getProcessor().process(exchange);
if (exchange.getPattern().isOutCapable() && exchange.hasOut())
{
- sendOutMessage(exchange);
+ sendOutMessage(exchange, messageSessionID);
}
} catch (Exception e) {
exchange.setException(e);
@@ -67,14 +71,12 @@ public class QuickfixjConsumer extends DefaultConsumer {
}
}
- private void sendOutMessage(Exchange exchange) throws QFJException {
+ private void sendOutMessage(Exchange exchange, SessionID messageSessionID)
throws QFJException {
Message camelMessage = exchange.getMessage();
quickfix.Message quickfixjMessage =
camelMessage.getBody(quickfix.Message.class);
LOG.debug("Sending FIX message reply: {}", quickfixjMessage);
- SessionID messageSessionID =
exchange.getIn().getHeader(QuickfixjEndpoint.SESSION_ID_KEY, SessionID.class);
-
Session session = getSession(messageSessionID);
if (session == null) {
throw new IllegalStateException("Unknown session: " +
messageSessionID);
diff --git
a/components/camel-quickfix/src/test/java/org/apache/camel/component/quickfixj/QuickfixjConsumerTest.java
b/components/camel-quickfix/src/test/java/org/apache/camel/component/quickfixj/QuickfixjConsumerTest.java
new file mode 100644
index 000000000000..e072018e11f4
--- /dev/null
+++
b/components/camel-quickfix/src/test/java/org/apache/camel/component/quickfixj/QuickfixjConsumerTest.java
@@ -0,0 +1,92 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.quickfixj;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePattern;
+import org.apache.camel.Processor;
+import org.apache.camel.component.quickfixj.converter.QuickfixjConverters;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+import quickfix.Message;
+import quickfix.Session;
+import quickfix.SessionID;
+
+import static org.hamcrest.CoreMatchers.nullValue;
+import static org.hamcrest.MatcherAssert.assertThat;
+
+public class QuickfixjConsumerTest extends CamelTestSupport {
+
+ private final SessionID requestSessionID = new SessionID("FIX.4.4",
"MARKET", "TRADER");
+ private final SessionID otherSessionID = new SessionID("FIX.4.4",
"MARKET", "OTHER");
+
+ @Test
+ public void replyIsSentOnTheSessionTheRequestArrivedOn() throws Exception {
+ Message reply = new Message();
+ Session requestSession = Mockito.mock(Session.class);
+ Mockito.when(requestSession.send(reply)).thenReturn(true);
+
+ QuickfixjConsumer consumer = startConsumer(replyWith(reply, null),
requestSession);
+ Exchange exchange = receive(consumer);
+
+ assertThat(exchange.getException(), nullValue());
+ Mockito.verify(requestSession).send(reply);
+ }
+
+ @Test
+ public void replyIgnoresSessionIdHeaderChangedDuringRouting() throws
Exception {
+ Message reply = new Message();
+ Session requestSession = Mockito.mock(Session.class);
+ Mockito.when(requestSession.send(reply)).thenReturn(true);
+
+ QuickfixjConsumer consumer = startConsumer(replyWith(reply,
otherSessionID), requestSession);
+ Exchange exchange = receive(consumer);
+
+ assertThat(exchange.getException(), nullValue());
+ Mockito.verify(requestSession).send(reply);
+ Mockito.verify(consumer, Mockito.never()).getSession(otherSessionID);
+ }
+
+ private QuickfixjConsumer startConsumer(Processor processor, Session
requestSession) throws Exception {
+ QuickfixjEndpoint endpoint = Mockito.mock(QuickfixjEndpoint.class);
+ Mockito.when(endpoint.getCamelContext()).thenReturn(context);
+
+ QuickfixjConsumer consumer = Mockito.spy(new
QuickfixjConsumer(endpoint, processor));
+
Mockito.doReturn(requestSession).when(consumer).getSession(requestSessionID);
+ consumer.start();
+ return consumer;
+ }
+
+ private Exchange receive(QuickfixjConsumer consumer) {
+ Exchange exchange = QuickfixjConverters.toExchange(consumer,
requestSessionID, new Message(),
+ QuickfixjEventCategory.AppMessageReceived,
ExchangePattern.InOut);
+ consumer.onExchange(exchange);
+ return exchange;
+ }
+
+ @SuppressWarnings("deprecation")
+ private static Processor replyWith(Message reply, SessionID
sessionIdHeader) {
+ return exchange -> {
+ if (sessionIdHeader != null) {
+ exchange.getIn().setHeader(QuickfixjEndpoint.SESSION_ID_KEY,
sessionIdHeader);
+ }
+ // an InOut reply is carried on the OUT message, as the bean
component does for InOut exchanges
+ exchange.getOut().setBody(reply);
+ };
+ }
+}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 966f5c9e6171..973341679cbd 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -3343,6 +3343,13 @@ coap://0.0.0.0:5683/my/resource?muteException=false
For the Rest DSL, set it with
`restConfiguration().endpointProperty("muteException", "false")`.
+=== camel-quickfix - InOut replies are sent on the session the request arrived
on
+
+With `exchangePattern=InOut`, the QuickFIX/J consumer now sends the reply on
the session the request was received on.
+Previously it read the `SessionID` header when the route completed, so a route
that changed that header also changed
+which session received the reply. A route that changed the `SessionID` header
to send the reply to another session
+must now send that message with a QuickFIX/J producer endpoint that sets the
`sessionID` option.
+
== ThrottlingExceptionRoutePolicy
`ThrottlingExceptionRoutePolicy.setKeepOpen(true)` now opens the circuit
immediately and synchronously (the consumer