davsclaus commented on code in PR #27314:
URL: https://github.com/apache/camel/pull/27314#discussion_r4195084739
##########
components/camel-hibernate/src/main/docs/hibernate-component.adoc:
##########
@@ -0,0 +1,731 @@
+= Hibernate Component
+:doctitle: Hibernate
+:shortname: hibernate
+:artifactid: camel-hibernate
+:description: Camel Hibernate Component
+:since: 4.23
+:supportlevel: Preview
+:tabs-sync-option:
+:component-header: Both producer and consumer are supported
+
+*Since Camel {since}*
+
+*{component-header}*
+
+The Hibernate component provides integration with relational databases using
**Hibernate ORM 8**.
+
+The component uses Hibernate's native `SessionFactory`, `Session`,
`StatelessSession`, `SelectionQuery`, and `MutationQuery` APIs.
+
+It supports HQL selection and mutation queries, natural-id lookups, Hibernate
filters, read-only queries, streaming query results, stateless insert and
upsert operations, multi-tenancy, and competing consumers using `SKIP_LOCKED`.
+
+The component can either reuse an existing Hibernate `SessionFactory` or
bootstrap one using a configured `DataSource` and explicit entity classes.
+
+The component does not use Jakarta Persistence (JPA), `camel-jpa`, or Spring
`PlatformTransactionManager`.
+
+Maven users will need to add the following dependency to their `pom.xml`
+for this component:
+
+[source,xml]
+----
+<dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-hibernate</artifactId>
+ <version>x.x.x</version>
+ <!-- use the same version as your Camel core version -->
+</dependency>
+----
+
+== URI format
+
+=== hibernate:entityClassName
+
+Where `entityClassName` is the target entity class name or entity type name.
+
+// component options: START
+include::partial$component-configure-options.adoc[]
+include::partial$component-endpoint-options.adoc[]
+include::partial$component-endpoint-headers.adoc[]
+// component options: END
+
+== Usage
+
+=== SessionFactory configuration
+
+The component can use an existing Hibernate `SessionFactory` or create one
during startup.
+
+==== Using an existing SessionFactory
+
+An existing `SessionFactory` can be configured directly on the component:
+
+[source,java]
+----
+HibernateComponent component = new HibernateComponent();
+component.setSessionFactory(sessionFactory);
+
+context.addComponent("hibernate", component);
+----
+
+The component does not own an externally supplied `SessionFactory` and
therefore does not close it when the component stops.
+
+If no `SessionFactory` is explicitly configured, the component attempts to
reuse a `SessionFactory` available in the Camel registry.
+
+==== Bootstrapping a SessionFactory
+
+The component can bootstrap a `SessionFactory` using a `DataSource` and
explicit entity classes:
+
+[source,java]
+----
+HibernateComponent component = new HibernateComponent();
+component.setDataSource(dataSource);
+component.setEntityClasses(new Class[] {
+ MyEntity.class
+});
+component.setSchemaAction("update");
+
+context.addComponent("hibernate", component);
+----
+
+The `DataSource` can be supplied directly or as a Camel registry name.
+
+Supported values for `schemaAction` are:
+
+* `none`
+* `validate`
+* `update`
+* `create`
+
+When the component creates the `SessionFactory`, it owns the resulting
`SessionFactory` and closes it when the component stops.
+
+Hibernate determines the database dialect automatically; the component does
not require an explicit dialect configuration.
+
+[NOTE]
+====
+Entity package scanning is not currently provided. Use `entityClasses` to
explicitly register the entity classes used by the `SessionFactory`.
+====
+
+==== Hibernate properties
+
+Additional Hibernate configuration properties can be supplied through
`hibernateProperties`.
+
+[source,java]
+----
+component.setHibernateProperties(Map.of(
+ "hibernate.show_sql", true,
+ "hibernate.format_sql", true
+));
+----
+
+The supplied properties are passed directly to Hibernate during
`SessionFactory` bootstrap.
+
+=== Producer operations
+
+When acting as a producer, the endpoint performs exactly one configured
operation.
+
+The supported operations are:
+
+* `selectionQuery`
+* `mutationQuery`
+* `naturalIdParameters`
+* `statelessOperation`
+
+The endpoint validates this configuration during startup.
+
+==== Selection query
+
+Use `selectionQuery` to execute an HQL selection query using Hibernate's
`SelectionQuery` API.
+
+[source,java]
+----
+from("direct:findEntities")
+ .to("hibernate:com.example.MyEntity?selectionQuery=from MyEntity where
status = :status");
+----
+
+Query parameters are supplied through the `CamelHibernateParameters` message
header:
+
+[source,java]
+----
+exchange.getMessage().setHeader(
+ HibernateConstants.HIBERNATE_PARAMETERS,
+ Map.of("status", "ACTIVE"));
+----
+
+The query result is placed in the message body.
+
+==== Mutation query
+
+Use `mutationQuery` to execute an HQL mutation query using Hibernate's
`MutationQuery` API.
+
+[source,java]
+----
+from("direct:updateEntities")
+ .to("hibernate:com.example.MyEntity?mutationQuery=update MyEntity set
status = :status where status = :oldStatus");
+----
+
+Parameters are supplied through the `CamelHibernateParameters` header:
+
+[source,java]
+----
+exchange.getMessage().setHeader(
+ HibernateConstants.HIBERNATE_PARAMETERS,
+ Map.of(
+ "status", "INACTIVE",
+ "oldStatus", "EXPIRED"));
+----
+
+The mutation update count is placed in the message body.
+
+==== Natural-id lookup
+
+Hibernate entities that define a natural identifier can be looked up using
`naturalIdParameters`.
+
+The lookup uses Hibernate's native natural-id functionality rather than
constructing an HQL query.
+
+==== Read-only queries
+
+Set `readOnly=true` for a selection query when the returned entities are not
intended to be modified.
+
+[source,java]
+----
+from("direct:readEntities")
+ .to("hibernate:com.example.MyEntity?selectionQuery=from
MyEntity&readOnly=true");
+----
+
+The component configures both the Hibernate `Session` and the `SelectionQuery`
as read-only.
+
+==== Hibernate filters
+
+Hibernate filters can be configured through the `filters` endpoint option.
+
+The configuration contains the filter name and its parameter values. The
component enables the configured filters on the Hibernate `Session` before
executing the operation.
+
+For example, an application can configure a filter named `tenantFilter` with a
tenant parameter:
+
+[source,java]
+----
+Map<String, Map<String, Object>> filters = Map.of(
+ "tenantFilter", Map.of("tenantId", "tenant-a"));
+
+endpoint.setFilters(filters);
+----
+
+When a producer reuses a `Session` supplied by a Hibernate consumer, its
filter configuration must match the filters associated with that `Session`. If
the producer's filters differ, the shared `Session` is not reused and the
producer opens a separate `Session` so that its filters can be applied without
changing the consumer's session state.
+
+==== Multi-tenancy
+
+A tenant identifier can be supplied through `tenantIdentifier`.
+
+[source,java]
+----
+from("direct:tenantEntities")
+ .to("hibernate:com.example.MyEntity?selectionQuery=from
MyEntity&tenantIdentifier=tenant-a");
+----
+
+The tenant identifier is used when creating the Hibernate `Session`.
+
+The tenant is configured by the route rather than being taken from an
untrusted message header.
+
+==== Stateless insert
+
+Use `statelessOperation=insert` to insert the entity contained in the message
body using Hibernate's native `StatelessSession`.
+
+[source,java]
+----
+from("direct:insertEntity")
+ .to("hibernate:com.example.MyEntity?statelessOperation=insert");
+----
+
+==== Stateless upsert
+
+Use `statelessOperation=upsert` to perform a stateless upsert using
Hibernate's native `StatelessSession`.
+
+[source,java]
+----
+from("direct:upsertEntity")
+ .to("hibernate:com.example.MyEntity?statelessOperation=upsert");
+----
+
+The supported stateless operations are `insert` and `upsert`.
+
+=== Consumer operations
+
+When acting as a consumer, the endpoint periodically executes its configured
`selectionQuery`.
+
+Each selected entity is delivered to the route as the message body.
+
+By default, successfully processed entities are deleted from the database
(`consumeDelete=true`).
+Set `consumeDelete=false` to leave consumed entities in the database.
+
+When a Hibernate consumer invokes a Hibernate producer in the same route
processing thread,
+the producer reuses the consumer's Hibernate `Session` and transaction when
the producer's
+session requirements are compatible with the consumer session. This allows the
+producer to operate on entities and locks held by the consumer without opening
a second
+session.
+
+Session sharing has limits. A copied exchange does not carry the live
Hibernate `Session`,
+so processing a copied exchange can open a separate session. Stateless
operations always
+use their own `StatelessSession`.
+
+If processing an entity fails, the consumer rolls back the active poll
transaction and
+stops processing the current batch. This prevents changes made earlier in the
same poll
+from being committed after a later entity fails.
+
+For example:
+
+[source,java]
+----
+from("hibernate:com.example.MyEntity?selectionQuery=from
MyEntity&maximumResults=100&initialDelay=1000&delay=5000")
+ .to("direct:processEntity");
+----
+
+==== Maximum results
+
+Use `maximumResults` to limit the number of entities retrieved during a poll.
+
+[source,java]
+----
+from("hibernate:com.example.MyEntity?selectionQuery=from
MyEntity&maximumResults=100")
+ .to("direct:processEntity");
+----
+
+==== Competing consumers with SKIP LOCKED
+
+Set `skipLocked=true` to configure the selection query with a pessimistic
write lock and Hibernate's `SKIP_LOCKED` timeout.
+
+[source,java]
+----
+from("hibernate:com.example.MyEntity?selectionQuery=from
MyEntity&skipLocked=true&maximumResults=100")
+ .to("direct:processEntity");
+----
+
+This allows multiple consumer instances to compete for database rows while
skipping rows already locked by another consumer.
+
+The database and Hibernate dialect must support the required locking behavior.
+
+When processing uses a copied exchange that cannot reuse the consumer's
`Session`, the
+processing may use a separate session and therefore does not share the
consumer's live
+transaction or session-level locks.
+
+=== Streaming query results
+
+Set `streaming=true` on a selection query to return the results as a Java
`Stream`.
+
+[source,java]
+----
+from("direct:streamEntities")
+ .to("hibernate:com.example.MyEntity?selectionQuery=from
MyEntity&streaming=true")
+ .process(exchange -> {
+ Stream<?> stream = exchange.getMessage().getBody(Stream.class);
+ try (stream) {
+ stream.forEach(entity -> {
+ // process entity
+ });
+ }
+ });
+----
+
+The Hibernate `Session` and transaction remain open while the stream is being
consumed.
+
+The returned stream must be closed after consumption so that the Hibernate
session and transaction can be closed.
+
+Streaming is not supported inside a transacted Camel exchange because the
stream can
+outlive the route transaction while the Hibernate `Session` and transaction
remain open.
+
+== Query parameters
+
+Query parameters are supplied through the `CamelHibernateParameters` message
header.
+
+For example:
+
+[source,java]
+----
+Map<String, Object> parameters = Map.of(
+ "status", "ACTIVE",
+ "country", "US");
+
+exchange.getMessage().setHeader(
+ HibernateConstants.HIBERNATE_PARAMETERS,
+ parameters);
+----
+
+The parameters are bound to the configured Hibernate query.
+
+The component does not support replacing the configured query through a
message header. Query text remains part of the route configuration.
+
+== Transactions
+
+The component uses native Hibernate transactions through `Session` and
`StatelessSession`.
+
+Regular producer and consumer operations create a Hibernate `Session` and
transaction. Each operation commits its own transaction unless the configured
`SessionFactory` is backed by JTA.
+
+Stateless operations use a Hibernate `StatelessSession` and transaction.
Stateless operations do not participate in consumer session reuse.
+
+The component does not provide its own JTA or Spring transaction-manager
integration. A transacted Camel exchange requires a JTA-backed
`SessionFactory`; using a resource-local `SessionFactory` in a transacted
exchange fails fast.
Review Comment:
This says a transacted exchange requires a JTA-backed `SessionFactory`,
while the limitations list further down says "JTA-backed transaction
coordination is not currently provided". Please make the two say the same thing
(what a user with a JTA `SessionFactory` gets today), and add a test for the
JTA path if it is supported.
##########
components/camel-hibernate/src/test/java/org/apache/camel/component/hibernate/HibernateBootstrapTest.java:
##########
@@ -0,0 +1,1121 @@
+/*
+ * 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.List;
+import java.util.Map;
+import java.util.stream.Stream;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.component.hibernate.entity.HibernateTestEntity;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.impl.engine.DefaultUnitOfWork;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.h2.jdbcx.JdbcDataSource;
+import org.hibernate.LockMode;
+import org.hibernate.Session;
+import org.hibernate.SessionBuilder;
+import org.hibernate.SessionFactory;
+import org.hibernate.Timeouts;
+import org.hibernate.Transaction;
+import org.hibernate.query.MutationQuery;
+import org.hibernate.query.SelectionQuery;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+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;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class HibernateBootstrapTest extends CamelTestSupport {
+
+ @Test
+ public void testExplicitSessionFactoryReuseAndNoClose() throws Exception {
+ SessionFactory mockSf = Mockito.mock(SessionFactory.class);
+
+ HibernateComponent comp = new HibernateComponent();
+ comp.setCamelContext(context);
+ comp.setSessionFactory(mockSf);
+ comp.start();
+
+ assertSame(mockSf, comp.getSessionFactory());
+ comp.stop();
+
+ Mockito.verify(mockSf, Mockito.never()).close();
+ }
+
+ @Test
+ public void testRegistrySessionFactoryReuseAndNoClose() throws Exception {
+ SessionFactory mockSf = Mockito.mock(SessionFactory.class);
+ context.getRegistry().bind("registrySf", mockSf);
+
+ HibernateComponent comp = new HibernateComponent();
+ comp.setCamelContext(context);
+ comp.start();
+
+ assertSame(mockSf, comp.getSessionFactory());
+ comp.stop();
+
+ Mockito.verify(mockSf, Mockito.never()).close();
+ }
+
+ @Test
+ public void testComponentCreatedSessionFactoryClosesOnStop() throws
Exception {
+ HibernateComponent comp = createComponent("testDs",
"jdbc:h2:mem:testdb;DB_CLOSE_DELAY=-1");
+
+ SessionFactory sf = comp.getSessionFactory();
+ assertNotNull(sf);
+ assertFalse(sf.isClosed());
+
+ comp.stop();
+ assertTrue(sf.isClosed());
+ }
+
+ @Test
+ public void testInvalidDataSourceRegistryNameFails() {
+ HibernateComponent comp = new HibernateComponent();
+ comp.setCamelContext(context);
+ comp.setDataSource("nonExistentDs");
+
+ Exception ex = assertThrows(Exception.class, comp::start);
+ assertTrue(ex.getMessage().contains("DataSource bean with name
'nonExistentDs' could not be found"));
+ }
+
+ @Test
+ public void testInvalidSchemaActionFails() {
+ HibernateComponent comp = new HibernateComponent();
+ comp.setCamelContext(context);
+ comp.setSchemaAction("invalid-action");
+
+ Exception ex = assertThrows(Exception.class, comp::start);
+ assertTrue(ex.getMessage().contains("Invalid schemaAction"));
+ }
+
+ @Test
+ public void testSelectionQueryProducer() throws Exception {
+ HibernateComponent comp = createComponent(
+ "selectionDs",
+ "jdbc:h2:mem:selectiondb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+
+ endpoint.start();
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+
exchange.getMessage().setHeader(HibernateConstants.HIBERNATE_PARAMETERS,
Map.of());
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ assertNotNull(exchange.getMessage().getBody());
+ assertTrue(exchange.getMessage().getBody() instanceof java.util.List);
+ assertTrue(((java.util.List<?>)
exchange.getMessage().getBody()).isEmpty());
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ public void testMutationQueryProducer() throws Exception {
+ HibernateComponent comp = createComponent(
+ "mutationDs",
+ "jdbc:h2:mem:mutationdb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+ endpoint.setMutationQuery("delete from HibernateTestEntity");
+
+ endpoint.start();
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+
exchange.getMessage().setHeader(HibernateConstants.HIBERNATE_PARAMETERS,
Map.of());
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ assertEquals(0, exchange.getMessage().getBody());
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ public void testSchemaActionsAcceptedValues() {
+ String[] actions = { "none", "validate", "update", "create" };
+ for (String action : actions) {
+ HibernateComponent comp = new HibernateComponent();
+ comp.setSchemaAction(action);
+ assertEquals(action, comp.getSchemaAction());
+ }
+ }
+
+ @Test
+ public void testExplicitEntityClassesBootstrap() throws Exception {
+ HibernateComponent comp = createComponent(
+ "entityDs",
+ "jdbc:h2:mem:entitydb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ SessionFactory sf = comp.getSessionFactory();
+ assertNotNull(sf);
+ assertNotNull(sf.getMetamodel().entity(HibernateTestEntity.class));
+
+ comp.stop();
+ }
+
+ @Test
+ public void testSelectionQueryProducerReadOnly() throws Exception {
+ HibernateComponent comp = createComponent(
+ "readOnlyDs",
+ "jdbc:h2:mem:readonlydb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setReadOnly(true);
+
+ endpoint.start();
+
+ assertTrue(endpoint.isReadOnly());
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+
exchange.getMessage().setHeader(HibernateConstants.HIBERNATE_PARAMETERS,
Map.of());
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ assertNotNull(exchange.getMessage().getBody());
+ assertTrue(exchange.getMessage().getBody() instanceof java.util.List);
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ public void testNaturalIdLookupProducer() throws Exception {
+ HibernateComponent comp = createComponent(
+ "naturalIdDs",
+ "jdbc:h2:mem:naturaliddb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ Session session = comp.getSessionFactory().openSession();
+ Transaction transaction = session.beginTransaction();
+
+ HibernateTestEntity entity = new HibernateTestEntity();
+ entity.setId(1L);
+ entity.setName("test");
+ session.persist(entity);
+
+ transaction.commit();
+ session.close();
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+ endpoint.setNaturalIdParameters(Map.of("name", "test"));
+
+ endpoint.start();
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ HibernateTestEntity result =
exchange.getMessage().getBody(HibernateTestEntity.class);
+
+ assertNotNull(result);
+ assertEquals(1L, result.getId());
+ assertEquals("test", result.getName());
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ public void testSelectionQueryWithFilter() throws Exception {
+ HibernateComponent comp = createComponent(
+ "filterDs",
+ "jdbc:h2:mem:filterdb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ Session session = comp.getSessionFactory().openSession();
+ Transaction transaction = session.beginTransaction();
+
+ HibernateTestEntity first = new HibernateTestEntity();
+ first.setId(1L);
+ first.setName("test");
+
+ HibernateTestEntity second = new HibernateTestEntity();
+ second.setId(2L);
+ second.setName("other");
+
+ session.persist(first);
+ session.persist(second);
+
+ transaction.commit();
+ session.close();
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setFilters(Map.of("nameFilter", Map.of("name", "test")));
+
+ endpoint.start();
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ java.util.List<?> results =
exchange.getMessage().getBody(java.util.List.class);
+
+ assertEquals(1, results.size());
+ HibernateTestEntity result = (HibernateTestEntity) results.get(0);
+ assertEquals("test", result.getName());
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testTenantIdentifierUsedForSession() throws Exception {
+ SessionFactory sessionFactory = Mockito.mock(SessionFactory.class);
+ SessionBuilder sessionBuilder = Mockito.mock(SessionBuilder.class);
+ Session session = Mockito.mock(Session.class);
+ Transaction transaction = Mockito.mock(Transaction.class);
+ SelectionQuery<HibernateTestEntity> query =
Mockito.mock(SelectionQuery.class);
+
+ Mockito.when(sessionFactory.withOptions()).thenReturn(sessionBuilder);
+
Mockito.when(sessionBuilder.tenantIdentifier("tenant1")).thenReturn(sessionBuilder);
+ Mockito.when(sessionBuilder.openSession()).thenReturn(session);
+ Mockito.when(session.beginTransaction()).thenReturn(transaction);
+ Mockito.when(session.createSelectionQuery(
+ "from HibernateTestEntity",
HibernateTestEntity.class)).thenReturn(query);
+ Mockito.when(query.getResultList()).thenReturn(java.util.List.of());
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setTenantIdentifier("tenant1");
+
+ endpoint.start();
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ Mockito.verify(sessionBuilder).tenantIdentifier("tenant1");
+ Mockito.verify(sessionBuilder).openSession();
+
+ endpoint.stop();
+ }
+
+ @Test
+ public void testStatelessInsertProducer() throws Exception {
+ HibernateComponent comp = createComponent(
+ "statelessInsertDs",
+ "jdbc:h2:mem:statelessinsertdb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+
endpoint.setStatelessOperation(HibernateEndpoint.StatelessOperation.INSERT);
+
+ endpoint.start();
+
+ HibernateTestEntity entity = new HibernateTestEntity();
+ entity.setId(1L);
+ entity.setName("stateless");
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+ exchange.getMessage().setBody(entity);
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ Session session = comp.getSessionFactory().openSession();
+ HibernateTestEntity result = session.find(HibernateTestEntity.class,
1L);
+ session.close();
+
+ assertNotNull(result);
+ assertEquals(1L, result.getId());
+ assertEquals("stateless", result.getName());
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ public void testStatelessUpsertProducer() throws Exception {
+ HibernateComponent comp = createComponent(
+ "statelessUpsertDs",
+ "jdbc:h2:mem:statelessupsertdb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(comp.getSessionFactory());
+ endpoint.setEntityClassName(HibernateTestEntity.class.getName());
+
endpoint.setStatelessOperation(HibernateEndpoint.StatelessOperation.UPSERT);
+
+ endpoint.start();
+
+ HibernateTestEntity entity = new HibernateTestEntity();
+ entity.setId(1L);
+ entity.setName("upsert");
+
+ Exchange exchange =
context.getEndpoint("direct:test").createExchange();
+ exchange.getMessage().setBody(entity);
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint)) {
+ producer.process(exchange);
+ }
+
+ Session session = comp.getSessionFactory().openSession();
+ HibernateTestEntity result = session.find(HibernateTestEntity.class,
1L);
+ session.close();
+
+ assertNotNull(result);
+ assertEquals(1L, result.getId());
+ assertEquals("upsert", result.getName());
+
+ endpoint.stop();
+ comp.stop();
+ }
+
+ @Test
+ void shouldConfigureSkipLockedQuery() throws Exception {
+ SessionFactory sessionFactory = Mockito.mock(SessionFactory.class);
+ Session session = Mockito.mock(Session.class);
+ Transaction transaction = Mockito.mock(Transaction.class);
+
+ @SuppressWarnings("unchecked")
+ SelectionQuery<HibernateTestEntity> query =
Mockito.mock(SelectionQuery.class);
+
+ Mockito.when(sessionFactory.openSession()).thenReturn(session);
+ Mockito.when(session.beginTransaction()).thenReturn(transaction);
+ Mockito.when(session.createSelectionQuery(
+ "from HibernateTestEntity",
HibernateTestEntity.class)).thenReturn(query);
+ Mockito.when(query.getResultList()).thenReturn(List.of());
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setSkipLocked(true);
+
+ HibernateConsumer consumer = new HibernateConsumer(endpoint,
exchange -> {
+ });
+
+ consumer.poll();
+
+
Mockito.verify(query).setHibernateLockMode(LockMode.PESSIMISTIC_WRITE);
+ Mockito.verify(query).setLockTimeout(Timeouts.SKIP_LOCKED);
+ Mockito.verify(transaction).commit();
+ }
+ }
+
+ @Test
+ void shouldSkipLockedRowsWithH2() throws Exception {
+ HibernateComponent component = createComponent(
+ "skipLockedDs",
+ "jdbc:h2:mem:skiplockeddb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ SessionFactory sessionFactory = component.getSessionFactory();
+
+ Session lockingSession = sessionFactory.openSession();
+ Transaction lockingTransaction = lockingSession.beginTransaction();
+
+ HibernateTestEntity entity = new HibernateTestEntity();
+ entity.setId(1L);
+ lockingSession.persist(entity);
+ lockingTransaction.commit();
+
+ lockingTransaction = lockingSession.beginTransaction();
+ lockingSession.find(
+ HibernateTestEntity.class,
+ entity.getId(),
+ LockMode.PESSIMISTIC_WRITE);
+
+ try {
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setSkipLocked(true);
+ endpoint.setConsumeDelete(false);
+
+ HibernateConsumer consumer = new HibernateConsumer(endpoint,
exchange -> {
+ });
+
+ assertEquals(0, consumer.poll());
+ } finally {
+ lockingTransaction.rollback();
+ lockingSession.close();
+ component.stop();
+ }
+ }
+
+ @Test
+ void shouldStreamSelectionQueryResults() throws Exception {
+ SessionFactory sessionFactory = Mockito.mock(SessionFactory.class);
+ Session session = Mockito.mock(Session.class);
+ Transaction transaction = Mockito.mock(Transaction.class);
+
+ @SuppressWarnings("unchecked")
+ SelectionQuery<HibernateTestEntity> query =
Mockito.mock(SelectionQuery.class);
+
+ HibernateTestEntity entity1 = new HibernateTestEntity();
+ HibernateTestEntity entity2 = new HibernateTestEntity();
+
+ Stream<HibernateTestEntity> stream = Stream.of(entity1, entity2);
+
+ Mockito.when(sessionFactory.openSession()).thenReturn(session);
+ Mockito.when(session.beginTransaction()).thenReturn(transaction);
+ Mockito.when(transaction.isActive()).thenReturn(true);
+ Mockito.when(session.createSelectionQuery(
+ "from HibernateTestEntity",
HibernateTestEntity.class)).thenReturn(query);
+ Mockito.when(query.getResultStream()).thenReturn(stream);
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setStreaming(true);
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint))
{
+ Exchange exchange = endpoint.createExchange();
+
+ producer.process(exchange);
+
+ Object body = exchange.getMessage().getBody();
+
+ assertInstanceOf(Stream.class, body);
+
+ @SuppressWarnings("unchecked")
+ Stream<HibernateTestEntity> result =
(Stream<HibernateTestEntity>) body;
+
+ assertEquals(List.of(entity1, entity2), result.toList());
+
+ result.close();
+ }
+ }
+
+ Mockito.verify(query).getResultStream();
+ Mockito.verify(transaction).commit();
+ Mockito.verify(session).close();
+ }
+
+ @Test
+ void shouldNotFailWhenStreamingExchangeCompletesBeforeStreamIsClosed()
throws Exception {
+ SessionFactory sessionFactory = Mockito.mock(SessionFactory.class);
+ Session session = Mockito.mock(Session.class);
+ Transaction transaction = Mockito.mock(Transaction.class);
+
+ @SuppressWarnings("unchecked")
+ SelectionQuery<HibernateTestEntity> query =
Mockito.mock(SelectionQuery.class);
+
+ Stream<HibernateTestEntity> stream = Stream.of(new
HibernateTestEntity());
+
+ Mockito.when(sessionFactory.openSession()).thenReturn(session);
+ Mockito.when(session.beginTransaction()).thenReturn(transaction);
+ Mockito.when(transaction.isActive()).thenReturn(true);
+ Mockito.when(session.createSelectionQuery(
+ "from HibernateTestEntity",
HibernateTestEntity.class)).thenReturn(query);
+ Mockito.when(query.getResultStream()).thenReturn(stream);
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setStreaming(true);
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint))
{
+ Exchange exchange = endpoint.createExchange();
+
+ producer.process(exchange);
+
+ @SuppressWarnings("unchecked")
+ Stream<HibernateTestEntity> result =
(Stream<HibernateTestEntity>) exchange.getMessage().getBody();
+
+ DefaultUnitOfWork unitOfWork = new DefaultUnitOfWork(exchange);
+ exchange.getExchangeExtension().setUnitOfWork(unitOfWork);
+ exchange.getUnitOfWork().done(exchange);
+
+ result.close();
+ }
+ }
+
+ Mockito.verify(transaction, Mockito.times(1)).commit();
+ Mockito.verify(session, Mockito.times(1)).close();
+ }
+
+ @Test
+ void shouldCloseSessionWhenStreamingTransactionIsInactive() throws
Exception {
+ SessionFactory sessionFactory = Mockito.mock(SessionFactory.class);
+ Session session = Mockito.mock(Session.class);
+ Transaction transaction = Mockito.mock(Transaction.class);
+
+ @SuppressWarnings("unchecked")
+ SelectionQuery<HibernateTestEntity> query =
Mockito.mock(SelectionQuery.class);
+
+ HibernateTestEntity entity = new HibernateTestEntity();
+ Stream<HibernateTestEntity> stream = Stream.of(entity);
+
+ Mockito.when(sessionFactory.openSession()).thenReturn(session);
+ Mockito.when(session.beginTransaction()).thenReturn(transaction);
+ Mockito.when(transaction.isActive()).thenReturn(false);
+ Mockito.when(session.createSelectionQuery(
+ "from HibernateTestEntity",
HibernateTestEntity.class)).thenReturn(query);
+ Mockito.when(query.getResultStream()).thenReturn(stream);
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setStreaming(true);
+
+ try (HibernateProducer producer = new HibernateProducer(endpoint))
{
+ Exchange exchange = endpoint.createExchange();
+
+ producer.process(exchange);
+
+ @SuppressWarnings("unchecked")
+ Stream<HibernateTestEntity> result =
(Stream<HibernateTestEntity>) exchange.getMessage().getBody();
+
+ result.close();
+ }
+ }
+
+ Mockito.verify(transaction, Mockito.never()).commit();
+ Mockito.verify(session).close();
+ }
+
+ @Test
+ void shouldRollbackBatchAndAbortOnProcessorFailure() throws Exception {
+ SessionFactory sessionFactory = Mockito.mock(SessionFactory.class);
+ Session session = Mockito.mock(Session.class);
+ Transaction transaction = Mockito.mock(Transaction.class);
+
+ @SuppressWarnings("unchecked")
+ SelectionQuery<HibernateTestEntity> query =
Mockito.mock(SelectionQuery.class);
+
+ HibernateTestEntity entity1 = new HibernateTestEntity();
+ HibernateTestEntity entity2 = new HibernateTestEntity();
+
+ Mockito.when(sessionFactory.openSession()).thenReturn(session);
+ Mockito.when(session.beginTransaction()).thenReturn(transaction);
+ Mockito.when(transaction.isActive()).thenReturn(true);
+ Mockito.when(session.createSelectionQuery(
+ "from HibernateTestEntity",
HibernateTestEntity.class)).thenReturn(query);
+ Mockito.when(query.getResultList()).thenReturn(List.of(entity1,
entity2));
+
+ List<HibernateTestEntity> processed = new java.util.ArrayList<>();
+
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ HibernateEndpoint endpoint = new HibernateEndpoint();
+ endpoint.setCamelContext(context);
+ endpoint.setSessionFactory(sessionFactory);
+ endpoint.setEntityType(HibernateTestEntity.class);
+ endpoint.setSelectionQuery("from HibernateTestEntity");
+ endpoint.setConsumeDelete(true);
+
+ HibernateConsumer consumer = new HibernateConsumer(endpoint,
exchange -> {
+ HibernateTestEntity entity =
exchange.getMessage().getBody(HibernateTestEntity.class);
+ processed.add(entity);
+
+ if (entity == entity2) {
+ throw new IllegalStateException("poison row");
+ }
+ });
+
+ assertThrows(IllegalStateException.class, consumer::poll);
+ }
+
+ assertEquals(List.of(entity1, entity2), processed);
+ Mockito.verify(session).remove(entity1);
+ Mockito.verify(session, Mockito.never()).remove(entity2);
+ Mockito.verify(transaction).rollback();
+ Mockito.verify(transaction, Mockito.never()).commit();
+ Mockito.verify(session).close();
+ }
+
+ @Test
+ void shouldReuseConsumerSessionInRouteWithSkipLocked() throws Exception {
+ HibernateComponent component = createComponent(
+ "routeReuseDs",
+ "jdbc:h2:mem:routeReuseDb;DB_CLOSE_DELAY=-1",
+ HibernateTestEntity.class);
+
+ SessionFactory sessionFactory = component.getSessionFactory();
+
+ Session setupSession = sessionFactory.openSession();
+ Transaction setupTransaction = setupSession.beginTransaction();
+
+ HibernateTestEntity entity = new HibernateTestEntity();
+ entity.setId(1L);
+ entity.setName("route-reuse");
+ setupSession.persist(entity);
+
+ setupTransaction.commit();
+ setupSession.close();
+
+ context.addComponent("hibernate", component);
+
+ context.addRoutes(new org.apache.camel.builder.RouteBuilder() {
Review Comment:
Nit: please import `RouteBuilder`, `Route`, `Consumer` (below) and
`org.hibernate.Filter` (further down) instead of the fully qualified names;
CLAUDE.md asks for imports, and the build's OpenRewrite step rewrites them,
which leaves uncommitted changes in CI.
--
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]