gnodet-bot commented on code in PR #27173:
URL: https://github.com/apache/camel/pull/27173#discussion_r4208459063


##########
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);
+            }
+        }
+    }
+
+    private void doOnBatch(List<Message> batch) throws Exception {
+        try {
+            batchListener.onBatch(batch, session);
+        } catch (Exception e) {
+            // unexpected error so rollback
+            rollbackIfNeeded(session);
+            throw e;
+        }
+        // success then commit if we need to
+        Message lastMessage = batch.get(batch.size() - 1);
+        commitIfNeeded(session, lastMessage);

Review Comment:
   ⚠️ **Bug (HIGH):** `doOnBatch()` calls `commitIfNeeded()` without checking 
`exchange.isFailed()` or `exchange.isRollbackOnly()`. If the route marks the 
exchange as rollback-only, this will silently commit the JMS transaction anyway.
   
   Compare with `EndpointMessageListener.processMessage()` which checks both 
flags before committing. You need equivalent logic here:
   
   ```suggestion
       private void doOnBatch(List<Message> batch) throws Exception {
           try {
               batchListener.onBatch(batch, session);
           } catch (Exception e) {
               // unexpected error so rollback
               rollbackIfNeeded(session);
               throw e;
           }
           // Check exchange failure state before committing
           Exchange batchExchange = batchListener.getLastBatchExchange();
           if (batchExchange != null && (batchExchange.isFailed() || 
batchExchange.isRollbackOnly())) {
               rollbackIfNeeded(session);
               return;
           }
           // success then commit if we need to
           Message lastMessage = batch.get(batch.size() - 1);
           commitIfNeeded(session, lastMessage);
       }
   ```
   
   (Exact fix depends on how you expose the exchange from 
`BatchEndpointMessageListener` — the key point is that 
`isFailed()`/`isRollbackOnly()` must be checked before `commitIfNeeded()`.)



##########
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:
   ⚠️ **Issue (MEDIUM):** Non-JMS exceptions are silently swallowed with a 
`WARN` log. `EndpointMessageListener` routes non-JMS failures through 
`consumer.getExceptionHandler().handleException(...)`, giving users a chance to 
react. The batch path skips this entirely — users have no way to hook in a 
dead-letter channel or custom exception handler for batch errors.
   
   Consider calling `consumer.getExceptionHandler().handleException("Batch 
listener failed", e)` here instead of (or in addition to) the WARN log.



##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchEndpointMessageListener.java:
##########
@@ -0,0 +1,87 @@
+/*
+ * 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 jakarta.jms.Message;
+import jakarta.jms.Session;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.component.sjms.SjmsConstants;
+import org.apache.camel.component.sjms.SjmsConsumer;
+import org.apache.camel.component.sjms.SjmsEndpoint;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class BatchEndpointMessageListener {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(BatchEndpointMessageListener.class);
+
+    private final SjmsConsumer consumer;
+    private final SjmsEndpoint endpoint;
+    private final Processor processor;
+
+    public BatchEndpointMessageListener(SjmsConsumer consumer, SjmsEndpoint 
endpoint, Processor processor) {
+        this.consumer = consumer;
+        this.endpoint = endpoint;
+        this.processor = processor;
+    }
+
+    private Exchange aggregate(List<Message> rawMessages, Session session) {
+        List<Exchange> exchanges = new ArrayList<>(rawMessages.size());
+        for (Message m : rawMessages) {
+            Exchange e = endpoint.createExchange(m, session);
+            // Force eager materialization of JMS message headers and body 
into the
+            // Camel Exchange, before the session is committed/closed after 
dispatch.
+            e.getIn().getHeaders();
+            e.getIn().getBody();
+
+            exchanges.add(e);
+        }
+
+        Exchange batchExchange = consumer.createExchange(false);
+        batchExchange.getIn().setBody(exchanges);
+        batchExchange.setProperty(SjmsConstants.JMS_SESSION, session);
+        
batchExchange.getMessage().setHeader(SjmsConstants.SJMS_BATCH_SIZE_HEADER, 
rawMessages.size());
+
+        return batchExchange;
+    }
+
+    void onBatch(List<Message> rawMessages, Session session) throws Exception {
+        LOG.trace("onBatch START");
+
+        LOG.debug("{} consumer received batch message: {}", endpoint, 
rawMessages);
+        Exchange batchExchange = null;
+        try {
+            batchExchange = aggregate(rawMessages, session);
+            processor.process(batchExchange);
+            consumer.releaseExchange(batchExchange, false);
+        } catch (Exception e) {
+            batchExchange.setException(e);

Review Comment:
   ⚠️ **Bug (HIGH):** NPE if `aggregate()` throws. `batchExchange` is `null` 
when entering the catch block in that case, so `batchExchange.setException(e)` 
on line 77 throws a `NullPointerException` that masks the original exception.
   
   ```suggestion
           Exchange batchExchange = null;
           try {
               batchExchange = aggregate(rawMessages, session);
               processor.process(batchExchange);
               consumer.releaseExchange(batchExchange, false);
           } catch (Exception e) {
               if (batchExchange != null) {
                   batchExchange.setException(e);
               } else {
                   throw e;
               }
           }
   ```



##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsConstants.java:
##########
@@ -38,5 +38,8 @@ public interface SjmsConstants {
               description = "Provides an explicit ReplyTo destination 
(overrides any incoming value of Message.getJMSReplyTo() in consumer)",
               javaType = "String")
     String JMS_REPLY_TO = JmsConstants.JMS_REPLY_TO;
-
+    @Metadata(label = "consumer, batch",
+              description = "The size of the batch when using the batching 
consumer option.",
+              javaType = "String")
+    String SJMS_BATCH_SIZE_HEADER = "CamelSjmsBatchSize";

Review Comment:
   ⚠️ **Bug (MEDIUM):** `javaType = "String"` is incorrect. The header value is 
set to `rawMessages.size()` which is an `int`. The `javaType` should be 
`"int"`. This mistake is propagated into all three generated catalog JSON files.
   
   ```suggestion
       @Metadata(label = "consumer,batch",
                 description = "The size of the batch when using the batching 
consumer option.",
                 javaType = "int")
       String SJMS_BATCH_SIZE_HEADER = "CamelSjmsBatchSize";
   ```
   
   (Also fixed the label spacing — see below.)



##########
components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchEndpointMessageListener.java:
##########
@@ -0,0 +1,87 @@
+/*
+ * 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 jakarta.jms.Message;
+import jakarta.jms.Session;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.component.sjms.SjmsConstants;
+import org.apache.camel.component.sjms.SjmsConsumer;
+import org.apache.camel.component.sjms.SjmsEndpoint;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class BatchEndpointMessageListener {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(BatchEndpointMessageListener.class);
+
+    private final SjmsConsumer consumer;
+    private final SjmsEndpoint endpoint;
+    private final Processor processor;
+
+    public BatchEndpointMessageListener(SjmsConsumer consumer, SjmsEndpoint 
endpoint, Processor processor) {
+        this.consumer = consumer;
+        this.endpoint = endpoint;
+        this.processor = processor;
+    }
+
+    private Exchange aggregate(List<Message> rawMessages, Session session) {
+        List<Exchange> exchanges = new ArrayList<>(rawMessages.size());
+        for (Message m : rawMessages) {
+            Exchange e = endpoint.createExchange(m, session);
+            // Force eager materialization of JMS message headers and body 
into the
+            // Camel Exchange, before the session is committed/closed after 
dispatch.
+            e.getIn().getHeaders();
+            e.getIn().getBody();
+
+            exchanges.add(e);
+        }
+
+        Exchange batchExchange = consumer.createExchange(false);
+        batchExchange.getIn().setBody(exchanges);
+        batchExchange.setProperty(SjmsConstants.JMS_SESSION, session);
+        
batchExchange.getMessage().setHeader(SjmsConstants.SJMS_BATCH_SIZE_HEADER, 
rawMessages.size());

Review Comment:
   Nit: inconsistent API usage on adjacent lines — `getIn()` on line 60, 
`getMessage()` on line 62. They return the same object on a new exchange, so 
this is harmless, but pick one for consistency.
   
   ```suggestion
           Exchange batchExchange = consumer.createExchange(false);
           batchExchange.getIn().setBody(exchanges);
           batchExchange.setProperty(SjmsConstants.JMS_SESSION, session);
           
batchExchange.getIn().setHeader(SjmsConstants.SJMS_BATCH_SIZE_HEADER, 
rawMessages.size());
   ```



-- 
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]

Reply via email to