gnodet-bot commented on code in PR #27173: URL: https://github.com/apache/camel/pull/27173#discussion_r4224315132
########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchConsumerWorker.java: ########## @@ -0,0 +1,178 @@ +/* + * 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.sjms.consumer; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import jakarta.jms.JMSException; +import jakarta.jms.Message; +import jakarta.jms.MessageConsumer; +import jakarta.jms.Session; + +import org.apache.camel.component.sjms.SjmsEndpoint; +import org.apache.camel.component.sjms.jms.SessionAcknowledgementType; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import static org.apache.camel.component.sjms.SjmsHelper.commitIfNeeded; +import static org.apache.camel.component.sjms.SjmsHelper.rollbackIfNeeded; + +public class BatchConsumerWorker implements Runnable { + + private static final Logger LOG = LoggerFactory.getLogger(BatchConsumerWorker.class); + + private static final long JMS_CONSUMER_RECEIVE_WAKE_INTERVAL_TIMEOUT = 1000L; + private static final long JMS_CONSUMER_RECEIVE_MIN_TIMEOUT = 100L; + + private final SjmsEndpoint endpoint; + private final BatchEndpointMessageListener batchListener; + private final MessageConsumer consumer; + private final Session session; + private final AtomicBoolean running = new AtomicBoolean(true); + private volatile boolean jmsUnhealthy; + + BatchConsumerWorker(SjmsEndpoint endpoint, BatchEndpointMessageListener batchListener, + MessageConsumer consumer, Session session) { + this.endpoint = endpoint; + this.batchListener = batchListener; + this.consumer = consumer; + this.session = session; + } + + void shutdown() { + LOG.debug("Shutdown requested"); + running.set(false); + } + + boolean isShutdownRequested() { + return !running.get(); + } + + void invalidate() { + jmsUnhealthy = true; + shutdown(); + } + + private boolean isRedeliverable() { + return endpoint.isTransacted() + || endpoint.getAcknowledgementMode() == SessionAcknowledgementType.CLIENT_ACKNOWLEDGE; + } + + @Override + public void run() { + LOG.debug("run START"); + int batchSize = endpoint.getBatchSize(); + long batchInterval = endpoint.getBatchInterval(); + List<Message> buffer = new ArrayList<>(); + long batchStartTime = 0L; + + try { + while (running.get()) { + long waitMillis + = computeWaitMillis(buffer.isEmpty(), batchStartTime, batchInterval); + + Message msg = consumer.receive(waitMillis); + + if (msg != null) { + if (buffer.isEmpty()) { + batchStartTime = System.nanoTime(); + } + buffer.add(msg); + } + + boolean sizeReached = batchSize > 0 && buffer.size() >= batchSize; + boolean intervalElapsed = batchInterval > 0 && !buffer.isEmpty() + && TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - batchStartTime) >= batchInterval; + + if (sizeReached || intervalElapsed) { + dispatchOrDiscard(buffer); + buffer = new ArrayList<>(); + } + } + + if (!buffer.isEmpty()) { + dispatchOrDiscard(buffer); // graceful-stop drain + } + } catch (JMSException e) { + jmsUnhealthy = true; + if (!buffer.isEmpty()) { + dispatchOrDiscard(buffer); + } + throw new BatchConsumerWorkerException(e); + } + LOG.info("run STOP"); + } + + private void dispatchOrDiscard(List<Message> batch) { + if (jmsUnhealthy) { + if (isRedeliverable()) { + LOG.warn("JMS is unhealthy on {}: discarding {} buffered message(s); " + + "unacknowledged/uncommitted, will be redelivered", + endpoint.getEndpointUri(), batch.size()); + return; + } + LOG.warn("JMS is unhealthy on {}: attempting best-effort dispatch of {} already-acknowledged " + + "message(s) since they cannot be redelivered", + endpoint.getEndpointUri(), batch.size()); + } + onBatch(batch); + } + + private void onBatch(List<Message> batch) { + try { + doOnBatch(batch); + } catch (Exception e) { + if (e instanceof JMSException jmsException) { + if (endpoint.getExceptionListener() != null) { + endpoint.getExceptionListener().onException(jmsException); + } + } else { + LOG.warn("Execution of JMS message listener failed. This exception is ignored.", e); Review Comment: ⚠️ **Non-JMS exceptions are silently swallowed here.** The `else` branch only logs the exception and continues. The route's `ExceptionHandler` is never notified, so the application has no way to react to unexpected failures in batch processing. Consider routing non-JMS exceptions through the consumer's error handler: ```suggestion } else { getExceptionHandler().handleException("Execution of batch JMS listener failed", e); } ``` ########## catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/sjms.json: ########## @@ -94,7 +95,10 @@ "recoveryInterval": { "index": 44, "kind": "parameter", "displayName": "Recovery Interval", "group": "advanced", "label": "advanced", "required": false, "type": "duration", "javaType": "long", "deprecated": false, "autowired": false, "secret": false, "defaultValue": "5000", "description": "Specifies the interval between recovery attempts, i.e. when a connection is being refreshed, in milliseconds. The default is 5000 ms, that is, 5 seconds." }, "synchronous": { "index": 45, "kind": "parameter", "displayName": "Synchronous", "group": "advanced", "label": "advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Sets whether synchronous processing should be strictly used" }, "transferException": { "index": 46, "kind": "parameter", "displayName": "Transfer Exception", "group": "advanced", "label": "advanced", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "security": "insecure:serialization", "defaultValue": false, "description": "If enabled and you are using Request Reply messaging (InOut) and an Exchange failed on the consumer side, then the caused Exception will be send back in response as a jakarta.jms.ObjectMessage. If the client is Camel, the returned Exception is rethrown. This allows you to use Camel JMS as a bridge in your routing - for example, using persistent queues to enable robust routing. Notice that if you also have transferExchange enabled, this option takes precedence. The caught exception is required to be serializable. The original Exception on the consumer side can be wrapped in an outer exception such as org.apache.camel.RuntimeCamelException when returned t o the producer. Use this with caution as the data is using Java Object serialization and requires the received to be able to deserialize the data at Class level, which forces a strong coupling between the producers and consumer!" }, - "deserializationFilter": { "index": 47, "kind": "parameter", "displayName": "Deserialization Filter", "group": "security", "label": "advanced,security", "required": false, "type": "string", "javaType": "java.lang.String", "deprecated": false, "autowired": false, "secret": false, "description": "Sets an ObjectInputFilter pattern (jdk.serialFilter syntax) applied as a defense-in-depth check on the class of the body returned by jakarta.jms.ObjectMessage.getObject(). The pattern is evaluated after the JMS provider has deserialized the payload, so this option alone does not prevent gadget-chain execution that happens inside the provider's ObjectInputStream; to block such attacks, also configure the JMS provider's own deserialization filter and\/or the JVM-wide -Djdk.serialFilter. When this option is not set and no JVM-wide filter is configured, a conservative default filter denying java.net. and otherwise allowing java., javax. and org.apache.camel. is applied." }, - "transacted": { "index": 48, "kind": "parameter", "displayName": "Transacted", "group": "transaction", "label": "transaction", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Specifies whether to use transacted mode" } + "batching": { "index": 47, "kind": "parameter", "displayName": "Batching", "group": "batch", "label": "consumer,batch", "required": false, "type": "boolean", "javaType": "boolean", "deprecated": false, "autowired": false, "secret": false, "defaultValue": false, "description": "Enable batch consuming. The route receives one Exchange per batch, whose body is a List of the individual JMS messages, instead of one Exchange per message." }, + "batchInterval": { "index": 48, "kind": "parameter", "displayName": "Batch Interval", "group": "batch", "label": "consumer,batch", "required": false, "type": "duration", "javaType": "long", "deprecated": false, "autowired": false, "secret": false, "defaultValue": "1000", "description": "Time in millis, measured from the first message received into a new batch, after which the batch is dispatched even if batchSize has not been reached comparable to the Aggregator EIP's completionInterval. Default is 1000 ms, that is 1 second. Interval should be a postive value. Set to 0 for unlimited (not recommended)." }, Review Comment: 💬 **Typo in generated file:** `postive` → `positive`. The source `SjmsEndpoint.java` has the correct spelling, but the generated catalog JSON and DSL builder Javadoc were not regenerated — `sjms.json`, `sjms2.json`, `SjmsEndpointBuilderFactory.java`, and `Sjms2EndpointBuilderFactory.java` all still carry the typo. Please run the code generator to pick up the corrected spelling. ########## components/camel-sjms/src/main/docs/sjms-component.adoc: ########## @@ -312,3 +312,23 @@ Currently, the only correlation strategy is to use the `JMSCorrelationId`. The _InOut_ Consumer uses this strategy as well ensuring that all response messages to the included `JMSReplyTo` destination also have the `JMSCorrelationId` copied from the request as well. + +=== Batch consuming + +The consumer can be configured to receive messages in batches instead of one at a time, by setting +`batching=true`. Instead of one Exchange per message, the route receives a single Exchange whose body +is a `List<Exchange>`, one per JMS message in the batch, with a `CamelSjmsBatchSize` header giving the +batch's size. Each Exchange in the batch is comparable to the Exchange generated by a non-batching consumer. + +Batching only supports the InOnly exchange pattern. There is no reply-to/request-response support for +batched consumption — a batch of N unrelated messages has no well-defined single reply, so InOut +routes are not supported when `batching=true`. + +A batch completes when either `batchSize` messages have been received, or the `batchInterval` — +measured from the first message in the batch - has elapsed. + +Commit/acknowledgement happens once per batch rather than once per message. With +`acknowledgementMode=AUTO_ACKNOWLEDGE`, messages are acknowledged as they're received into the batch, +before routing — so a failed batch cannot be redelivered. Use `transacted=true` or +`acknowledgementMode=CLIENT_ACKNOWLEDGE` if the whole batch should be redelivered together on failure. + Review Comment: 💬 **The `concurrentConsumers > 1` behavior is worth documenting.** When `concurrentConsumers > 1`, each consumer thread runs its own independent `BatchConsumerWorker` — separate JMS session, separate batch accumulator, independent `batchSize`/`batchInterval` timers. This means the route will receive multiple independent batches simultaneously, not one merged batch. Consider adding a note like: > Setting `concurrentConsumers` greater than 1 creates parallel independent batch listeners, each accumulating and routing its own batches. ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerRollbackOnlyTest.java: ########## @@ -0,0 +1,94 @@ +/* + * 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.sjms.batch; + +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.Processor; +import org.apache.camel.RoutesBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.support.JmsTestSupport; +import org.junit.jupiter.api.Test; + +import static java.lang.String.format; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_FINISH; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_START; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.assertBatchSizesInOrder; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.createBatchRoute; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.sendMessages; + +public class BatchConsumerRollbackOnlyTest extends JmsTestSupport { + + private static final String QUEUE_NAME_TEMPLATE = "batch.consumer.%s.BatchConsumerRollbackOnlyTest"; + + private static final String ROUTE_ID_SESSION_TX = "tx"; + private static final String ROUTE_ID_CLIENT_ACK_NO_TX = "no-tx-client-ack"; + + @Test + public void testSessionTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_SESSION_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_SESSION_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + } + + @Test + public void testClientAcknowledgedNotTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, + format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + } + + @Override + protected RoutesBuilder[] createRouteBuilders() { + return new org.apache.camel.RoutesBuilder[] { Review Comment: 💬 **Redundant FQN — `RoutesBuilder` is already imported at line 22.** Use the short name: ```suggestion return new RoutesBuilder[] { ``` ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerTransactedTest.java: ########## @@ -0,0 +1,117 @@ +/* + * 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.sjms.batch; + +import java.util.Collections; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.Processor; +import org.apache.camel.RoutesBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.support.JmsTestSupport; +import org.junit.jupiter.api.Test; + +import static java.lang.String.format; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_FINISH; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_START; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.DEFAULT_MESSAGE_TEXT; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.assertBatchSizesInOrder; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.createBatchRoute; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.getBatchBodiesAsString; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.sendMessages; +import static org.junit.jupiter.api.Assertions.assertEquals; + +public class BatchConsumerTransactedTest extends JmsTestSupport { + + private static final String QUEUE_NAME_TEMPLATE = "batch.consumer.%s.BatchConsumerTransactedTest"; + + private static final String ROUTE_ID_SESSION_TX = "tx"; + private static final String ROUTE_ID_CLIENT_ACK_NO_TX = "no-tx-client-ack"; + private static final String ROUTE_ID_AUTO_ACK_NO_TX = "no-tx-auto"; + + @Test + public void testSessionTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_SESSION_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_SESSION_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + assertEquals(Collections.nCopies(5, DEFAULT_MESSAGE_TEXT), getBatchBodiesAsString(mockFinish.getExchanges().get(0))); + } + + @Test + public void testClientAcknowledgedNotTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, + format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + assertEquals(Collections.nCopies(5, DEFAULT_MESSAGE_TEXT), getBatchBodiesAsString(mockFinish.getExchanges().get(0))); + } + + @Test + public void testAutoAcknowledgedNotTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_AUTO_ACK_NO_TX)); + mockStart.expectedMessageCount(1); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_AUTO_ACK_NO_TX)); + mockFinish.expectedMessageCount(0); + + sendMessages(template, + format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_AUTO_ACK_NO_TX), 5); + + MockEndpoint.assertIsSatisfied(context); + } + + @Override + protected RoutesBuilder[] createRouteBuilders() { + return new org.apache.camel.RoutesBuilder[] { Review Comment: 💬 **Redundant FQN — `RoutesBuilder` is already imported.** Use `RoutesBuilder[]`. ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerRollbackOnlyTest.java: ########## @@ -0,0 +1,94 @@ +/* + * 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.sjms.batch; + +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.Processor; +import org.apache.camel.RoutesBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.support.JmsTestSupport; +import org.junit.jupiter.api.Test; + +import static java.lang.String.format; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_FINISH; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_START; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.assertBatchSizesInOrder; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.createBatchRoute; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.sendMessages; + +public class BatchConsumerRollbackOnlyTest extends JmsTestSupport { + + private static final String QUEUE_NAME_TEMPLATE = "batch.consumer.%s.BatchConsumerRollbackOnlyTest"; + + private static final String ROUTE_ID_SESSION_TX = "tx"; + private static final String ROUTE_ID_CLIENT_ACK_NO_TX = "no-tx-client-ack"; + + @Test + public void testSessionTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_SESSION_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_SESSION_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + } + + @Test + public void testClientAcknowledgedNotTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, + format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + } + + @Override + protected RoutesBuilder[] createRouteBuilders() { + return new org.apache.camel.RoutesBuilder[] { + createBatchRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX, 5, + 1000, + true, null, 1, new MarkRollBackProcessor()), + createBatchRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX, + 5, 1000, false, "CLIENT_ACKNOWLEDGE", 1, new MarkRollBackProcessor()) + }; + } + + private static class MarkRollBackProcessor implements Processor { + private final AtomicInteger counter = new AtomicInteger(); + + @Override + public void process(org.apache.camel.Exchange exchange) { Review Comment: 💬 **Redundant FQN — `Exchange` is already imported.** Use `Exchange exchange`. ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerTransactedTest.java: ########## @@ -0,0 +1,117 @@ +/* + * 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.sjms.batch; + +import java.util.Collections; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.camel.Processor; +import org.apache.camel.RoutesBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.support.JmsTestSupport; +import org.junit.jupiter.api.Test; + +import static java.lang.String.format; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_FINISH; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.BATCH_ROUTEBUILDER_MOCK_START; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.DEFAULT_MESSAGE_TEXT; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.assertBatchSizesInOrder; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.createBatchRoute; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.getBatchBodiesAsString; +import static org.apache.camel.component.sjms.batch.BatchTestHelper.sendMessages; +import static org.junit.jupiter.api.Assertions.assertEquals; + +public class BatchConsumerTransactedTest extends JmsTestSupport { + + private static final String QUEUE_NAME_TEMPLATE = "batch.consumer.%s.BatchConsumerTransactedTest"; + + private static final String ROUTE_ID_SESSION_TX = "tx"; + private static final String ROUTE_ID_CLIENT_ACK_NO_TX = "no-tx-client-ack"; + private static final String ROUTE_ID_AUTO_ACK_NO_TX = "no-tx-auto"; + + @Test + public void testSessionTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_SESSION_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_SESSION_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + assertEquals(Collections.nCopies(5, DEFAULT_MESSAGE_TEXT), getBatchBodiesAsString(mockFinish.getExchanges().get(0))); + } + + @Test + public void testClientAcknowledgedNotTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockStart.expectedMessageCount(2); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_CLIENT_ACK_NO_TX)); + mockFinish.expectedMessageCount(1); + + sendMessages(template, + format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX), + 5); + + MockEndpoint.assertIsSatisfied(context); + assertBatchSizesInOrder(mockFinish, 5); + assertEquals(Collections.nCopies(5, DEFAULT_MESSAGE_TEXT), getBatchBodiesAsString(mockFinish.getExchanges().get(0))); + } + + @Test + public void testAutoAcknowledgedNotTransacted() throws Exception { + MockEndpoint mockStart = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_START, ROUTE_ID_AUTO_ACK_NO_TX)); + mockStart.expectedMessageCount(1); + + MockEndpoint mockFinish = getMockEndpoint(format(BATCH_ROUTEBUILDER_MOCK_FINISH, ROUTE_ID_AUTO_ACK_NO_TX)); + mockFinish.expectedMessageCount(0); + + sendMessages(template, + format("sjms:queue:" + QUEUE_NAME_TEMPLATE, ROUTE_ID_AUTO_ACK_NO_TX), 5); + + MockEndpoint.assertIsSatisfied(context); + } + + @Override + protected RoutesBuilder[] createRouteBuilders() { + return new org.apache.camel.RoutesBuilder[] { + createBatchRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_SESSION_TX, 5, + 1000, + true, null, 1, new ThrowExceptionProcessor()), + createBatchRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_CLIENT_ACK_NO_TX, + 5, 1000, false, "CLIENT_ACKNOWLEDGE", 1, new ThrowExceptionProcessor()), + createBatchRoute(QUEUE_NAME_TEMPLATE, ROUTE_ID_AUTO_ACK_NO_TX, 5, 1000, false, + "AUTO_ACKNOWLEDGE", 1, new ThrowExceptionProcessor()) + }; + } + + private static class ThrowExceptionProcessor implements Processor { + private final AtomicInteger counter = new AtomicInteger(); + + @Override + public void process(org.apache.camel.Exchange exchange) { Review Comment: 💬 **Redundant FQN — `Exchange` is already imported.** Use `Exchange exchange`. ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerValidationTest.java: ########## @@ -0,0 +1,89 @@ +/* + * 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.sjms.batch; + +import org.apache.camel.Consumer; +import org.apache.camel.component.sjms.SjmsEndpoint; +import org.apache.camel.component.sjms.support.JmsTestSupport; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import static org.junit.jupiter.api.Assertions.*; Review Comment: 💬 **Wildcard static import breaks project convention.** All other test files in this diff use explicit static imports (e.g. `import static org.junit.jupiter.api.Assertions.assertEquals`). Replace the wildcard: ```suggestion import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; ``` (Adjust to the assertions actually used in the file.) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
