davsclaus commented on code in PR #27173: URL: https://github.com/apache/camel/pull/27173#discussion_r4227343786
########## 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: Following up on my comment in the conversation, which replaces my earlier "fine to leave it as is" on this line: please pass non-JMS failures to the consumer's `ExceptionHandler` instead of only logging them, so `bridgeErrorHandler` and a custom `exceptionHandler` also apply to batches. The rollback in `doOnBatch()` stays as it is, and redelivery limits stay with the broker. Please use the consumer's handler rather than `endpoint.getExceptionHandler()`: the endpoint's handler is `null` unless the `exceptionHandler` option is set (so the helper from 7201315 would have thrown an NPE), while the consumer's handler defaults to a logging handler. `BatchEndpointMessageListener` already holds the `SjmsConsumer`, for example: ```java // BatchEndpointMessageListener void handleException(String message, Throwable cause) { consumer.getExceptionHandler().handleException(message, cause); } // BatchConsumerWorker.onBatch(), else branch batchListener.handleException("Execution of batch JMS message listener failed", e); ``` gnodet-bot's suggestions on these lines call methods the worker does not have, so they won't compile as they are. A test would cover it: `transacted=true`, a route that always fails, a custom `exceptionHandler` that records the call, and a check that the batch is redelivered. ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsEndpoint.java: ########## @@ -285,6 +291,18 @@ public class SjmsEndpoint extends DefaultEndpoint @UriParam(defaultValue = "false", label = "advanced", description = "Sets whether synchronous processing should be strictly used") private boolean synchronous; + @UriParam(label = "consumer,batch", defaultValue = "false", + description = "Enable batch consuming. The route receives one Exchange per batch, whose body" + + " is a List<Exchange> of the individual JMS messages, instead of one Exchange per message.") + private boolean batching; + @UriParam(defaultValue = "100", label = "consumer,batch", + description = "Maximum number of messages per batch.") + private int batchSize = 100; + @UriParam(defaultValue = "1000", label = "consumer,batch", javaType = "java.time.Duration", + 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 positive value. Set to 0 for unlimited (not recommended).") Review Comment: Thanks for fixing `postive`. Two small wording points are left: the generator strips the em dash, so the catalog reads "has not been reached comparable to the Aggregator", and "unlimited" is unclear for 0. A possible wording: ```suggestion 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. Set to 0 to disable the " + "interval, so a batch is only dispatched when batchSize is reached (not recommended).") ``` Either way the generated files need a rebuild: the catalog `sjms.json`/`sjms2.json` and the Javadoc in `SjmsEndpointBuilderFactory`/`Sjms2EndpointBuilderFactory` still say `postive`, and the catalog copy of `sjms-component.adoc` still has the reply-to NOTE in its old place. Please build `components/camel-sjms` and `components/camel-sjms2`, then `catalog/camel-catalog`, then `dsl/camel-endpointdsl`, and commit what changes; otherwise CI's uncommitted-changes check fails. ########## 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 NOTE is back in place, thanks. As discussed in the conversation, the section should still cover what happens to a failed batch, `concurrentConsumers` and `JMSReplyTo`. A possible text, adjust as you like: ```suggestion `acknowledgementMode=CLIENT_ACKNOWLEDGE` if the whole batch should be redelivered together on failure. A failed batch is rolled back and redelivered as a whole, so the route should be idempotent. One bad message fails, and redelivers, the whole batch: all its messages share the delivery count, and they end up together in the broker's dead letter queue (if one is configured). How often and how fast the broker redelivers is up to its redelivery policy, for example `max-delivery-attempts` and `redelivery-delay` in ActiveMQ Artemis. To avoid this, handle failures per message inside the batch, for example with `split(body())` and `doTry`, or with `onException(...).handled(true)`. The messages of a redelivered batch have the `JMSRedelivered` header set. With `concurrentConsumers` greater than 1, each consumer collects and routes its own batches, so several batches are processed in parallel. The batch consumer never sends a reply, so a `JMSReplyTo` on the incoming messages is ignored. ``` ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchConsumerSizeTest.java: ########## @@ -0,0 +1,62 @@ +/* + * 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 org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.support.JmsTestSupport; +import org.junit.jupiter.api.Test; + +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.getBatchBodiesAsString; +import static org.junit.jupiter.api.Assertions.assertEquals; + +public class BatchConsumerSizeTest extends JmsTestSupport { + + private static final String SJMS_FROMF_URI = "%s?batching=true&batchSize=5"; Review Comment: Without an explicit `batchInterval`, the 1 s default now decides whether the first five messages land in one batch. On a slow CI agent the interval can fire first, and `assertBatchSizesInOrder(mock, 5, 2)` fails. An interval well above the send time keeps this a size test, and the second batch still arrives within the mock's default 10 s wait: ```suggestion private static final String SJMS_FROMF_URI = "%s?batching=true&batchSize=5&batchInterval=5000"; ``` ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchTestHelper.java: ########## @@ -0,0 +1,249 @@ +/* + * 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.lang.reflect.Field; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; + +import jakarta.jms.Connection; +import jakarta.jms.ConnectionFactory; +import jakarta.jms.ExceptionListener; +import jakarta.jms.JMSContext; +import jakarta.jms.JMSException; + +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.server.ServerConsumer; +import org.apache.activemq.artemis.core.server.ServerSession; +import org.apache.camel.CamelContext; +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.ProducerTemplate; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.SjmsConstants; +import org.apache.camel.component.sjms.SjmsConsumer; +import org.apache.camel.component.sjms.jms.JmsConstants; +import org.apache.camel.test.infra.artemis.services.ArtemisService; +import org.apache.camel.test.infra.artemis.services.ArtemisVMService; + +import static java.lang.String.format; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +@SuppressWarnings("unchecked") +public final class BatchTestHelper { + + static final String DEFAULT_MESSAGE_TEXT = "Hello World!"; + static final String BATCH_ROUTEBUILDER_MOCK_START = "mock:%s.start"; + static final String BATCH_ROUTEBUILDER_MOCK_FINISH = "mock:%s.complete"; + + private BatchTestHelper() { + } + + static void sendMessages(ProducerTemplate template, String endpoint, int count) { + sendMessagesWithText(template, endpoint, count, DEFAULT_MESSAGE_TEXT); + } + + static void sendMessagesWithText(ProducerTemplate template, String endpoint, int count, String text) { + for (int i = 0; i < count; i++) { + template.sendBody(endpoint, format(text, i)); + } + } + + static List<String> getBatchBodiesAsString(Exchange batchExchange) { + return BatchTestHelper.getBatchExchanges(batchExchange).stream() + .map(e -> e.getIn().getBody(String.class)) + .toList(); + } + + /** Extracts the batch body as a List<Exchange> from a batch Exchange, asserting the type. */ + static List<Exchange> getBatchExchanges(Exchange batchExchange) { + List<Exchange> body = batchExchange.getIn().getBody(List.class); + assertNotNull(body, "batch exchange body was null"); + for (Object o : body) { + assertInstanceOf(Exchange.class, o, "batch element was not an Exchange: " + o); + } + return body; + } + + /** Asserts a single batch exchange has exactly the given number of messages. */ + static void assertBatchSize(Exchange batchExchange, int expectedSize) { + List<Exchange> batch = getBatchExchanges(batchExchange); + assertEquals(expectedSize, batch.size(), "unexpected batch size"); + assertEquals(expectedSize, + batchExchange.getIn().getHeader(SjmsConstants.SJMS_BATCH_SIZE_HEADER, Integer.class), + "CamelSjmsBatchSize header did not match actual batch size"); + } + + /** + * Asserts the mock received exactly expectedSizes.length batch exchanges, in order, with each batch's size matching + * the corresponding element. Requires concurrentConsumers=1 (or an otherwise deterministic single-worker setup) so + * that arrival order is meaningful. + */ + static void assertBatchSizesInOrder(MockEndpoint mock, int... expectedSizes) { + assertEquals(expectedSizes.length, mock.getExchanges().size(), + "Number of expected sizes ddoes not match the number of exchanges"); + List<Exchange> received = mock.getExchanges(); + for (int i = 0; i < expectedSizes.length; i++) { + assertBatchSize(received.get(i), expectedSizes[i]); + } + } + + static RouteBuilder createBatchRoute( + String queueName, String id, int batchSize, int batchInterval, Boolean transacted, + String acknowledgementMode, + int concurrentConsumers, Processor processor) { + + Map<String, Object> params = new LinkedHashMap<>(); + params.put("batching", true); + params.put("batchSize", batchSize); + params.put("batchInterval", batchInterval); + params.put("transacted", transacted); + params.put("acknowledgementMode", acknowledgementMode); + params.put("concurrentConsumers", concurrentConsumers); + + String query = params.entrySet().stream() + .filter(e -> e.getValue() != null) + .map(e -> e.getKey() + "=" + e.getValue()) + .collect(Collectors.joining("&")); + + String fromJmsEndpoint = "sjms:queue:" + format(queueName, id) + (query.isEmpty() ? "" : "?" + query); + + return new RouteBuilder() { + @Override + public void configure() { + from(fromJmsEndpoint) + .id(id) + .toF(BATCH_ROUTEBUILDER_MOCK_START, id) + .process(processor) + .toF(BATCH_ROUTEBUILDER_MOCK_FINISH, id); + } + }; + } + + /** + * Asserts the JMSRedelivered flag of every message in the batch whose body equals the given body. Fails if no + * message has that body, so a typo in the body cannot make the assertion pass vacuously. + */ + static void assertBatchRedelivered(Exchange batchExchange, String body, boolean expectedRedelivered) { + List<Exchange> matching = getBatchExchanges(batchExchange).stream() + .filter(e -> body.equals(e.getIn().getBody(String.class))) + .toList(); + + assertFalse(matching.isEmpty(), "no message in the batch has body " + body); + for (Exchange e : matching) { + assertEquals(expectedRedelivered, + e.getIn().getHeader(JmsConstants.JMS_REDELIVERED, Boolean.class), + "unexpected JMSRedelivered for a message with body " + body); + } + } + + static void triggerConnectionFailure(CamelContext context, String routeId) throws Exception { + SjmsConsumer consumer = (SjmsConsumer) context + .getRoute(routeId) + .getConsumer(); + + Field field = SjmsConsumer.class.getDeclaredField("listenerContainer"); + + field.setAccessible(true); + + ExceptionListener listener = (ExceptionListener) field.get(consumer); + + listener.onException( + new JMSException("Simulated connection failure")); + } + + static class DoNothingProcessor implements org.apache.camel.Processor { Review Comment: `Processor` is imported at line 37. ```suggestion static class DoNothingProcessor implements Processor { ``` ########## 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: Same here (and `org.apache.camel.Exchange` on line 110). ```suggestion return new RoutesBuilder[] { ``` ########## components/camel-sjms/src/test/java/org/apache/camel/component/sjms/batch/BatchTestHelper.java: ########## @@ -0,0 +1,249 @@ +/* + * 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.lang.reflect.Field; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; + +import jakarta.jms.Connection; +import jakarta.jms.ConnectionFactory; +import jakarta.jms.ExceptionListener; +import jakarta.jms.JMSContext; +import jakarta.jms.JMSException; + +import org.apache.activemq.artemis.core.server.ActiveMQServer; +import org.apache.activemq.artemis.core.server.ServerConsumer; +import org.apache.activemq.artemis.core.server.ServerSession; +import org.apache.camel.CamelContext; +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.ProducerTemplate; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.sjms.SjmsConstants; +import org.apache.camel.component.sjms.SjmsConsumer; +import org.apache.camel.component.sjms.jms.JmsConstants; +import org.apache.camel.test.infra.artemis.services.ArtemisService; +import org.apache.camel.test.infra.artemis.services.ArtemisVMService; + +import static java.lang.String.format; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +@SuppressWarnings("unchecked") +public final class BatchTestHelper { + + static final String DEFAULT_MESSAGE_TEXT = "Hello World!"; + static final String BATCH_ROUTEBUILDER_MOCK_START = "mock:%s.start"; + static final String BATCH_ROUTEBUILDER_MOCK_FINISH = "mock:%s.complete"; + + private BatchTestHelper() { + } + + static void sendMessages(ProducerTemplate template, String endpoint, int count) { + sendMessagesWithText(template, endpoint, count, DEFAULT_MESSAGE_TEXT); + } + + static void sendMessagesWithText(ProducerTemplate template, String endpoint, int count, String text) { + for (int i = 0; i < count; i++) { + template.sendBody(endpoint, format(text, i)); + } + } + + static List<String> getBatchBodiesAsString(Exchange batchExchange) { + return BatchTestHelper.getBatchExchanges(batchExchange).stream() + .map(e -> e.getIn().getBody(String.class)) + .toList(); + } + + /** Extracts the batch body as a List<Exchange> from a batch Exchange, asserting the type. */ + static List<Exchange> getBatchExchanges(Exchange batchExchange) { + List<Exchange> body = batchExchange.getIn().getBody(List.class); + assertNotNull(body, "batch exchange body was null"); + for (Object o : body) { + assertInstanceOf(Exchange.class, o, "batch element was not an Exchange: " + o); + } + return body; + } + + /** Asserts a single batch exchange has exactly the given number of messages. */ + static void assertBatchSize(Exchange batchExchange, int expectedSize) { + List<Exchange> batch = getBatchExchanges(batchExchange); + assertEquals(expectedSize, batch.size(), "unexpected batch size"); + assertEquals(expectedSize, + batchExchange.getIn().getHeader(SjmsConstants.SJMS_BATCH_SIZE_HEADER, Integer.class), + "CamelSjmsBatchSize header did not match actual batch size"); + } + + /** + * Asserts the mock received exactly expectedSizes.length batch exchanges, in order, with each batch's size matching + * the corresponding element. Requires concurrentConsumers=1 (or an otherwise deterministic single-worker setup) so + * that arrival order is meaningful. + */ + static void assertBatchSizesInOrder(MockEndpoint mock, int... expectedSizes) { + assertEquals(expectedSizes.length, mock.getExchanges().size(), + "Number of expected sizes ddoes not match the number of exchanges"); Review Comment: Typo. ```suggestion "Number of expected sizes does not match the number of exchanges"); ``` ########## 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))); Review Comment: Thanks for adding the body check. All five messages have the same body, so it can't show that the redelivered batch holds the same messages as the failed one. `sendMessagesWithText(template, uri, 5, "Message %d")` gives distinct bodies; then compare `getBatchBodiesAsString(mockStart.getExchanges().get(0))` (the failed attempt) with `getBatchBodiesAsString(mockFinish.getExchanges().get(0))` (the redelivery). The so-far unused `assertBatchRedelivered(..., true)` can check the second attempt. The same applies to `testClientAcknowledgedNotTransacted`. ########## 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: `RoutesBuilder` is already imported. On line 87, `org.apache.camel.Exchange` needs `import org.apache.camel.Exchange;` (not imported in this file yet, nor in `BatchConsumerTransactedTest`, which has the same two spots). `BatchConsumerValidationTest` also still uses `Assertions.*`. ```suggestion return new RoutesBuilder[] { ``` -- 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]
