ramu11 commented on code in PR #27314: URL: https://github.com/apache/camel/pull/27314#discussion_r4197047139
########## 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: addressed ########## 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: done -- 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]
