gnodet-bot commented on code in PR #27173: URL: https://github.com/apache/camel/pull/27173#discussion_r4152655023
########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsEndpoint.java: ########## @@ -24,18 +24,8 @@ import jakarta.jms.Message; import jakarta.jms.Session; -import org.apache.camel.AsyncEndpoint; -import org.apache.camel.Category; -import org.apache.camel.Component; -import org.apache.camel.Consumer; -import org.apache.camel.Exchange; -import org.apache.camel.ExchangePattern; -import org.apache.camel.MultipleConsumersSupport; -import org.apache.camel.PollingConsumer; -import org.apache.camel.Processor; -import org.apache.camel.Producer; -import org.apache.camel.component.sjms.consumer.EndpointMessageListener; -import org.apache.camel.component.sjms.consumer.SimpleMessageListenerContainer; +import org.apache.camel.*; +import org.apache.camel.component.sjms.consumer.*; Review Comment: 🔴 **Convention violation:** Star imports (`org.apache.camel.*` and `org.apache.camel.component.sjms.consumer.*`) are not used anywhere else in the camel-sjms source tree. The existing code uses explicit imports. This appears to be an IDE auto-import artifact from the AI-assisted development. ```suggestion import org.apache.camel.AggregationStrategy; import org.apache.camel.AsyncEndpoint; import org.apache.camel.Category; import org.apache.camel.Component; import org.apache.camel.Consumer; import org.apache.camel.Exchange; import org.apache.camel.ExchangePattern; import org.apache.camel.MultipleConsumersSupport; import org.apache.camel.PollingConsumer; import org.apache.camel.Processor; import org.apache.camel.Producer; import org.apache.camel.component.sjms.consumer.BatchDefaultExchangeListAggregationStrategy; import org.apache.camel.component.sjms.consumer.BatchEndpointMessageListener; import org.apache.camel.component.sjms.consumer.BatchMessageListenerContainer; import org.apache.camel.component.sjms.consumer.EndpointMessageListener; import org.apache.camel.component.sjms.consumer.SimpleMessageListenerContainer; ``` ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchEndpointMessageListener.java: ########## @@ -0,0 +1,95 @@ +/* + * 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.List; + +import jakarta.jms.Message; +import jakarta.jms.Session; + +import org.apache.camel.AggregationStrategy; +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.component.sjms.SjmsConstants; +import org.apache.camel.component.sjms.SjmsEndpoint; +import org.apache.camel.component.sjms.SjmsHelper; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class BatchEndpointMessageListener { + + public static final String SJMS_BATCH_SIZE_HEADER = "CamelSjmsBatchSize"; + + private static final Logger LOG = LoggerFactory.getLogger(BatchEndpointMessageListener.class); + + private final SjmsEndpoint endpoint; + private final Processor processor; + private final AggregationStrategy aggregationStrategy; + + public BatchEndpointMessageListener(SjmsEndpoint endpoint, Processor processor, AggregationStrategy aggregationStrategy) { + this.endpoint = endpoint; + this.processor = processor; + + this.aggregationStrategy = aggregationStrategy; + } + + private Exchange aggregate(List<Message> rawMessages, Session session) { + Exchange result = null; + for (Message m : rawMessages) { + Exchange e = endpoint.createExchange(m, session); + // Populate the headers and body of the Exchange in message + e.getIn().getHeaders(); + e.getIn().getBody(); + + result = aggregationStrategy.aggregate(result, e); + } + if (result != null) { + aggregationStrategy.onCompletion(result); + } Review Comment: ⚠️ **Suspicious side-effect-only calls need a comment:** `e.getIn().getHeaders()` and `e.getIn().getBody()` are called purely for their side effects — to force lazy-loading of JMS message properties into the Camel Exchange before the JMS Session is committed/closed after batch dispatch. This is a known Camel pattern, but without an explanatory comment, a future maintainer will remove these apparent dead-code lines. ```suggestion // 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(); ``` ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchEndpointMessageListener.java: ########## @@ -0,0 +1,95 @@ +/* + * 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.List; + +import jakarta.jms.Message; +import jakarta.jms.Session; + +import org.apache.camel.AggregationStrategy; +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.component.sjms.SjmsConstants; +import org.apache.camel.component.sjms.SjmsEndpoint; +import org.apache.camel.component.sjms.SjmsHelper; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class BatchEndpointMessageListener { + + public static final String SJMS_BATCH_SIZE_HEADER = "CamelSjmsBatchSize"; + + private static final Logger LOG = LoggerFactory.getLogger(BatchEndpointMessageListener.class); + + private final SjmsEndpoint endpoint; + private final Processor processor; + private final AggregationStrategy aggregationStrategy; + + public BatchEndpointMessageListener(SjmsEndpoint endpoint, Processor processor, AggregationStrategy aggregationStrategy) { + this.endpoint = endpoint; + this.processor = processor; + + this.aggregationStrategy = aggregationStrategy; + } + + private Exchange aggregate(List<Message> rawMessages, Session session) { + Exchange result = null; + for (Message m : rawMessages) { + Exchange e = endpoint.createExchange(m, session); + // Populate the headers and body of the Exchange in message + e.getIn().getHeaders(); + e.getIn().getBody(); + + result = aggregationStrategy.aggregate(result, e); + } + if (result != null) { + aggregationStrategy.onCompletion(result); + } + return result; + } + + void onBatch(List<Message> rawMessages, Session session) throws Exception { + Exchange batchExchange = null; + Exception failure = null; + try { + batchExchange = aggregate(rawMessages, session); + if (batchExchange != null) { + batchExchange.setProperty(SjmsConstants.JMS_SESSION, session); + batchExchange.getMessage().setHeader(SJMS_BATCH_SIZE_HEADER, rawMessages.size()); + processor.process(batchExchange); + } + } catch (Exception e) { + failure = e; + } + + boolean failed = failure != null + || (batchExchange != null && (batchExchange.isFailed() || batchExchange.isRollbackOnly())); + Message lastMessage = rawMessages.get(rawMessages.size() - 1); + + if (!failed) { + SjmsHelper.commitIfNeeded(session, lastMessage); + } else { + Exception cause = failure != null ? failure : batchExchange.getException(); Review Comment: ⚠️ **`rawMessages.get(rawMessages.size() - 1)` throws on empty list:** While `BatchConsumerWorker.dispatch()` only calls `onBatch` when the buffer is non-empty, `onBatch` is package-private and could be called from other code paths. A guard clause is cheap insurance: ```suggestion if (rawMessages.isEmpty()) { return; } Message lastMessage = rawMessages.get(rawMessages.size() - 1); ``` ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchEndpointMessageListener.java: ########## @@ -0,0 +1,95 @@ +/* + * 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.List; + +import jakarta.jms.Message; +import jakarta.jms.Session; + +import org.apache.camel.AggregationStrategy; +import org.apache.camel.Exchange; +import org.apache.camel.Processor; +import org.apache.camel.component.sjms.SjmsConstants; +import org.apache.camel.component.sjms.SjmsEndpoint; +import org.apache.camel.component.sjms.SjmsHelper; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +public class BatchEndpointMessageListener { + + public static final String SJMS_BATCH_SIZE_HEADER = "CamelSjmsBatchSize"; + + private static final Logger LOG = LoggerFactory.getLogger(BatchEndpointMessageListener.class); + + private final SjmsEndpoint endpoint; + private final Processor processor; + private final AggregationStrategy aggregationStrategy; + + public BatchEndpointMessageListener(SjmsEndpoint endpoint, Processor processor, AggregationStrategy aggregationStrategy) { + this.endpoint = endpoint; + this.processor = processor; + + this.aggregationStrategy = aggregationStrategy; + } + + private Exchange aggregate(List<Message> rawMessages, Session session) { + Exchange result = null; + for (Message m : rawMessages) { + Exchange e = endpoint.createExchange(m, session); + // Populate the headers and body of the Exchange in message + e.getIn().getHeaders(); + e.getIn().getBody(); + + result = aggregationStrategy.aggregate(result, e); + } + if (result != null) { + aggregationStrategy.onCompletion(result); + } + return result; + } + + void onBatch(List<Message> rawMessages, Session session) throws Exception { + Exchange batchExchange = null; + Exception failure = null; + try { + batchExchange = aggregate(rawMessages, session); + if (batchExchange != null) { + batchExchange.setProperty(SjmsConstants.JMS_SESSION, session); + batchExchange.getMessage().setHeader(SJMS_BATCH_SIZE_HEADER, rawMessages.size()); + processor.process(batchExchange); + } + } catch (Exception e) { + failure = e; + } + + boolean failed = failure != null + || (batchExchange != null && (batchExchange.isFailed() || batchExchange.isRollbackOnly())); + Message lastMessage = rawMessages.get(rawMessages.size() - 1); + + if (!failed) { + SjmsHelper.commitIfNeeded(session, lastMessage); + } else { + Exception cause = failure != null ? failure : batchExchange.getException(); + if (cause != null) { + LOG.warn("Batch of {} message(s) failed processing on {}: {}", rawMessages.size(), + endpoint.getEndpointUri(), cause.getMessage(), cause); + } + SjmsHelper.rollbackIfNeeded(session); + } + } +} Review Comment: ⚠️ **Incomplete failure logging when `isRollbackOnly()` is set without an exception:** When `failed=true` because `batchExchange.isRollbackOnly()` is set (without an exception), `cause` is null so the `LOG.warn` is skipped entirely. The session is rolled back correctly, but there's no log trace of *why* the batch was rolled back. ```suggestion if (!failed) { SjmsHelper.commitIfNeeded(session, lastMessage); } else { Exception cause = failure != null ? failure : batchExchange.getException(); if (cause != null) { LOG.warn("Batch of {} message(s) failed processing on {}: {}", rawMessages.size(), endpoint.getEndpointUri(), cause.getMessage(), cause); } else { LOG.warn("Batch of {} message(s) rolled back on {} (rollback-only or failed exchange)", rawMessages.size(), endpoint.getEndpointUri()); } SjmsHelper.rollbackIfNeeded(session); } ``` ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchConsumerWorkerException.java: ########## @@ -0,0 +1,25 @@ +/* + * 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; + +public class BatchConsumerWorkerException + extends RuntimeException { + + public BatchConsumerWorkerException(Throwable err) { + super(err); + } +} Review Comment: 💡 **Public exception class with no Javadoc:** This is public API in an ASF project. It should at minimum have a class-level Javadoc explaining when this exception is thrown. ```suggestion /** * Unchecked exception thrown when a {@link BatchConsumerWorker} encounters * a JMS-level failure that terminates its polling loop. */ public class BatchConsumerWorkerException extends RuntimeException { public BatchConsumerWorkerException(Throwable err) { super(err); } } ``` ########## components/camel-sjms2/src/generated/java/org/apache/camel/component/sjms2/Sjms2EndpointUriFactory.java: ########## @@ -24,13 +24,16 @@ public class Sjms2EndpointUriFactory extends org.apache.camel.support.component. private static final Set<String> ENDPOINT_IDENTITY_PROPERTY_NAMES; private static final Map<String, String> MULTI_VALUE_PREFIXES; static { - Set<String> props = new HashSet<>(52); + Set<String> props = new HashSet<>(55); props.add("acknowledgementMode"); props.add("allowNullBody"); props.add("asyncConsumer"); props.add("asyncStartListener"); props.add("asyncStopListener"); Review Comment: 🔴 **Inconsistent generated files for sjms2:** The sjms2 URI factory adds `batchSize`, `batchTimeout`, `batching` (3 properties) while sjms's URI factory adds `batching`, `batchingAggregationStrategy`, `batchingInterval`, `batchingSize`, `batchingTimeout` (5 properties). The sjms2 version is missing `batchingAggregationStrategy`, `batchingInterval`, `batchingSize`, `batchingTimeout` and has `batchSize`/`batchTimeout` instead of `batchingSize`/`batchingTimeout`. This naming mismatch means the sjms2 catalog and URI factory won't match the Java bean property names. The same divergence appears in the catalog JSON files. This likely means the code generator wasn't re-run for sjms2 after the property names were finalized. Please regenerate: `mvn generate-sources -pl components/camel-sjms2`. ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsEndpoint.java: ########## @@ -68,6 +58,8 @@ public class SjmsEndpoint extends DefaultEndpoint implements AsyncEndpoint, MultipleConsumersSupport, HeaderFilterStrategyAware { + public static final long DEFAULT_CONSUMER_BATCHING_INTERVAL_MILLIS = 1000L; + private boolean topic; Review Comment: 💡 **Public constant is only used internally:** `DEFAULT_CONSUMER_BATCHING_INTERVAL_MILLIS` is `public` but only referenced from `BatchConsumerWorker` (a package-private class). Consider making it package-private to avoid exposing an implementation detail as API. ########## components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/BatchConsumerWorker.java: ########## @@ -0,0 +1,155 @@ +/* + * 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.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; + +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); Review Comment: 💡 **`System.currentTimeMillis()` for relative deadlines:** `batchStartTime` and `lastMessageTime` use `System.currentTimeMillis()` which is subject to NTP clock adjustments. If the system clock jumps, batch timing can be thrown off. For relative deadline calculations, `System.nanoTime()` (monotonic) is more appropriate. Low severity since JMS batching typically uses second-level granularity. -- 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]
