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


##########
components/camel-hibernate/src/main/java/org/apache/camel/component/hibernate/HibernateProducer.java:
##########
@@ -0,0 +1,205 @@
+/*
+ * 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.hibernate;
+
+import java.util.Map;
+import java.util.stream.Stream;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.hibernate.KeyType;
+import org.hibernate.Session;
+import org.hibernate.StatelessSession;
+import org.hibernate.Transaction;
+import org.hibernate.query.MutationQuery;
+import org.hibernate.query.SelectionQuery;
+
+public class HibernateProducer extends DefaultProducer {
+
+    private final HibernateEndpoint endpoint;
+
+    public HibernateProducer(HibernateEndpoint endpoint) {
+        super(endpoint);
+        this.endpoint = endpoint;
+    }
+
+    @Override
+    public void process(Exchange exchange) throws Exception {
+        if (endpoint.getStatelessOperation() != null) {
+            processStateless(exchange);
+            return;
+        }
+
+        Session session = 
exchange.getProperty(HibernateConstants.HIBERNATE_SESSION, Session.class);
+        boolean sessionOwned = session == null;
+
+        if (sessionOwned) {
+            session = endpoint.getTenantIdentifier() == null
+                    ? endpoint.getSessionFactory().openSession()
+                    : endpoint.getSessionFactory().withOptions()
+                            .tenantIdentifier(endpoint.getTenantIdentifier())
+                            .openSession();
+        }
+
+        final Session activeSession = session;
+        final Transaction transaction = sessionOwned
+                ? activeSession.beginTransaction()
+                : activeSession.getTransaction();
+
+        try {
+            if (endpoint.getFilters() != null) {
+                endpoint.getFilters().forEach((filterName, parameters) -> {
+                    var filter = activeSession.enableFilter(filterName);
+                    if (parameters != null) {
+                        parameters.forEach(filter::setParameter);
+                    }
+                });
+            }
+
+            @SuppressWarnings("unchecked")
+            Map<String, Object> parameters = exchange.getMessage().getHeader(
+                    HibernateConstants.HIBERNATE_PARAMETERS, Map.class);
+
+            if (endpoint.getNaturalIdParameters() != null) {
+                Object entity = activeSession.find(
+                        endpoint.getEntityType(),
+                        endpoint.getNaturalIdParameters(),
+                        KeyType.NATURAL);
+
+                exchange.getMessage().setBody(entity);
+
+                if (sessionOwned) {
+                    transaction.commit();
+                    activeSession.close();
+                }
+            } else if (endpoint.getSelectionQuery() != null) {
+                SelectionQuery<?> query = activeSession.createSelectionQuery(
+                        endpoint.getSelectionQuery(), 
endpoint.getEntityType());
+
+                activeSession.setDefaultReadOnly(endpoint.isReadOnly());
+                query.setReadOnly(endpoint.isReadOnly());
+
+                if (parameters != null) {
+                    parameters.forEach(query::setParameter);
+                }
+
+                if (endpoint.isStreaming()) {
+                    if (sessionOwned) {
+                        exchange.getExchangeExtension().addOnCompletion(new 
SynchronizationAdapter() {
+                            @Override
+                            public void onComplete(Exchange exchange) {
+                                closeStreamingSession(activeSession, 
transaction, false);
+                            }
+
+                            @Override
+                            public void onFailure(Exchange exchange) {
+                                closeStreamingSession(activeSession, 
transaction, true);
+                            }
+                        });
+                    }
+
+                    Stream<?> stream = query.getResultStream();
+
+                    if (sessionOwned) {
+                        exchange.getMessage().setBody(stream.onClose(
+                                () -> closeStreamingSession(activeSession, 
transaction, false)));

Review Comment:
   ⚠️ **Bug — Double-close on streaming session.** When `sessionOwned=true`, 
both the `SynchronizationAdapter` (lines 103–113) and `stream.onClose()` (line 
119–120) call `closeStreamingSession()`. When the exchange completes and the 
stream is closed, the session will be closed twice — the second 
`session.close()` will throw `IllegalStateException: Session/EntityManager is 
closed`.
   
   The `SynchronizationAdapter` is the correct safety net for when the stream 
is never closed. But when both fire, the second call crashes.
   
   Fix: guard with a simple `AtomicBoolean` flag, or remove the 
`stream.onClose()` callback entirely and rely solely on the 
`SynchronizationAdapter` (which already handles both success and failure paths):
   
   ```suggestion
                       Stream<?> stream = query.getResultStream();
   
                       if (sessionOwned) {
                           exchange.getMessage().setBody(stream);
                       } else {
                           exchange.getMessage().setBody(stream);
                       }
   ```
   
   (With the `SynchronizationAdapter` already registered above, the 
`stream.onClose()` is redundant and harmful.)



##########
parent/pom.xml:
##########
@@ -3892,6 +3897,26 @@
                 <artifactId>javax.inject</artifactId>
                 <version>${javax-inject-version}</version>
             </dependency>
+            <dependency>
+                <groupId>jakarta.persistence</groupId>
+                <artifactId>jakarta.persistence-api</artifactId>
+                <version>${jakarta-persistence-api-version}</version>
+            </dependency>
+            <dependency>
+                <groupId>org.hibernate.orm</groupId>
+                <artifactId>hibernate-core</artifactId>
+                <version>${hibernate-version}</version>
+            </dependency>
+            <dependency>
+                <groupId>org.springframework</groupId>
+                <artifactId>spring-orm</artifactId>
+                <version>${spring-version}</version>
+            </dependency>
+            <dependency>
+                <groupId>com.h2database</groupId>
+                <artifactId>h2</artifactId>
+                <version>${h2-version}</version>
+            </dependency>
         </dependencies>
     </dependencyManagement>
 

Review Comment:
   ⚠️ **Parent POM still polluted.** `hibernate-core` (with 
`${hibernate-version}` = 7.4.x, not the 8.x the component uses), `spring-orm`, 
`jakarta.persistence-api`, and `h2` are still added to the shared parent 
`dependencyManagement`. This was davsclaus's finding #6 from the first review.
   
   - `spring-orm` — the component doesn't use Spring at all anymore, so this 
entry is dead weight
   - `hibernate-core` — pinned to `${hibernate-version}` (7.4.x) while the 
component uses its own `${hibernate8-version}` (8.0.0.Beta3), so this managed 
entry doesn't even serve camel-hibernate — it only risks affecting `camel-jpa` 
or other modules
   - `h2` — test-scope dependency, should use the existing managed version 
(already managed elsewhere in the parent POM)
   - `jakarta.persistence-api` — already managed as 
`${jakarta-persistence-api-version}` elsewhere in the parent POM
   
   Please remove these four entries from the parent `dependencyManagement`. 
Keep the versions local to the component POM as you've already done for 
`hibernate8-version` and `jakarta-persistence4-api-version`.



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