From 2ab61de46cb5c0bf45e7b2411c3b8aab1b593abd Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Tue, 7 May 2024 16:09:33 -0500 Subject: [PATCH] Add tests for container initiated mixed txn * Adds tests for using a transactional PulsarTemplate from within a `@Transactional` `@PulsarListener` method. The test sends a message to Pulsar and inserts a row in DB and does some form of rollback (or not) to make sure things are as expected. Resolves #661 --- .../PulsarListenerWithDbTransactionTests.java | 185 ++++++++++++++++++ ...PulsarTemplateWithDbTransactionTests.java} | 29 +-- 2 files changed, 200 insertions(+), 14 deletions(-) create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarListenerWithDbTransactionTests.java rename spring-pulsar/src/test/java/org/springframework/pulsar/transaction/{PulsarMixedTransactionTests.java => PulsarTemplateWithDbTransactionTests.java} (71%) diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarListenerWithDbTransactionTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarListenerWithDbTransactionTests.java new file mode 100644 index 00000000..9bff338a --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarListenerWithDbTransactionTests.java @@ -0,0 +1,185 @@ +/* + * Copyright 2023-2024 the original author or authors. + * + * Licensed 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 + * + * https://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.springframework.pulsar.transaction; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.common.util.ObjectMapperFactory; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Configuration; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.pulsar.annotation.PulsarListener; +import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.listener.AckMode; +import org.springframework.pulsar.transaction.PulsarListenerWithDbTransactionTests.WithDbAndPulsarTransactionCommit.WithDbAndPulsarTransactionCommitConfig; +import org.springframework.pulsar.transaction.PulsarListenerWithDbTransactionTests.WithDbTransactionRollback.WithDbTransactionRollbackConfig; +import org.springframework.pulsar.transaction.PulsarListenerWithDbTransactionTests.WithPulsarTransactionRollback.WithPulsarTransactionRollbackConfig; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.transaction.annotation.EnableTransactionManagement; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.interceptor.TransactionAspectSupport; + +/** + * Tests transaction support of {@link PulsarListener} when mixed with database + * transactions. + * + * @author Chris Bono + */ +class PulsarListenerWithDbTransactionTests extends PulsarTxnWithDbTxnTestsBase { + + @Nested + @ContextConfiguration(classes = WithDbAndPulsarTransactionCommitConfig.class) + class WithDbAndPulsarTransactionCommit { + + static final CountDownLatch latch = new CountDownLatch(1); + static final String topicIn = "plwdbtxn-happy-in"; + static final String topicOut = "plwdbtxn-happy-out"; + + @Test + void whenDbTxnIsCommittedThenMessagesAreCommitted() throws Exception { + var nonTransactionalTemplate = newNonTransactionalTemplate(false, 1); + var thing = new Thing(1L, "msg1"); + var thingJson = ObjectMapperFactory.getMapper().getObjectMapper().writeValueAsString(thing); + nonTransactionalTemplate.send(topicIn, thingJson); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThatMessagesAreInTopic(topicOut, thing.name()); + assertThatMessagesAreInDb(thing); + } + + @EnableTransactionManagement + @Configuration(proxyBeanMethods = false) + static class WithDbAndPulsarTransactionCommitConfig { + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private PulsarTemplate transactionalPulsarTemplate; + + @Transactional("dataSourceTransactionManager") + @PulsarListener(topics = topicIn, ackMode = AckMode.RECORD) + void listen(String msgJson) throws Exception { + var thing = ObjectMapperFactory.getMapper().getObjectMapper().readValue(msgJson, Thing.class); + this.transactionalPulsarTemplate.send(topicOut, thing.name()); + PulsarTxnWithDbTxnTestsBase.insertThingIntoDb(jdbcTemplate, thing); + latch.countDown(); + } + + } + + } + + @Nested + @ContextConfiguration(classes = WithDbTransactionRollbackConfig.class) + class WithDbTransactionRollback { + + static final CountDownLatch latch = new CountDownLatch(1); + static final String topicIn = "plwdbtxn-dbr-in"; + static final String topicOut = "plwdbtxn-dbr-out"; + + @Test + void whenDbTxnIsSetRollbackOnlyThenMessageCommittedInPulsarButNotInDb() throws Exception { + var nonTransactionalTemplate = newNonTransactionalTemplate(false, 1); + var thing = new Thing(2L, "msg2"); + var thingJson = ObjectMapperFactory.getMapper().getObjectMapper().writeValueAsString(thing); + nonTransactionalTemplate.send(topicIn, thingJson); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThatMessagesAreNotInDb(thing); + assertThatMessagesAreInTopic(topicOut, thing.name()); + } + + @EnableTransactionManagement + @Configuration(proxyBeanMethods = false) + static class WithDbTransactionRollbackConfig { + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private PulsarTemplate transactionalPulsarTemplate; + + @Transactional("dataSourceTransactionManager") + @PulsarListener(topics = topicIn, ackMode = AckMode.RECORD) + void listen(String msgJson) throws Exception { + var thing = ObjectMapperFactory.getMapper().getObjectMapper().readValue(msgJson, Thing.class); + this.transactionalPulsarTemplate.send(topicOut, thing.name()); + PulsarTxnWithDbTxnTestsBase.insertThingIntoDb(jdbcTemplate, thing); + TransactionAspectSupport.currentTransactionStatus().setRollbackOnly(); + latch.countDown(); + } + + } + + } + + @Nested + @ContextConfiguration(classes = WithPulsarTransactionRollbackConfig.class) + class WithPulsarTransactionRollback { + + static final CountDownLatch latch = new CountDownLatch(1); + static final String topicIn = "plwdbtxn-pr-in"; + static final String topicOut = "plwdbtxn-pr-out"; + + @Test + void whenPulsarTxnIsSetRollbackOnlyThenMessageCommittedInDbButNotInPulsar() throws Exception { + var nonTransactionalTemplate = newNonTransactionalTemplate(false, 1); + var thing = new Thing(3L, "msg3"); + var thingJson = ObjectMapperFactory.getMapper().getObjectMapper().writeValueAsString(thing); + nonTransactionalTemplate.send(topicIn, thingJson); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThatMessagesAreInDb(thing); + assertThatMessagesAreNotInTopic(topicOut, thing.name()); + } + + @EnableTransactionManagement + @Configuration(proxyBeanMethods = false) + static class WithPulsarTransactionRollbackConfig { + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private PulsarTemplate transactionalPulsarTemplate; + + @Autowired + private PulsarClient pulsarClient; + + @Transactional("dataSourceTransactionManager") + @PulsarListener(topics = topicIn, ackMode = AckMode.RECORD) + void listen(String msgJson) throws Exception { + if (latch.getCount() == 0) { + return; + } + var thing = ObjectMapperFactory.getMapper().getObjectMapper().readValue(msgJson, Thing.class); + this.transactionalPulsarTemplate.send(topicOut, thing.name()); + PulsarTxnWithDbTxnTestsBase.insertThingIntoDb(jdbcTemplate, thing); + PulsarTransactionUtils.getResourceHolder(this.pulsarClient).setRollbackOnly(); + latch.countDown(); + } + + } + + } + +} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarMixedTransactionTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTemplateWithDbTransactionTests.java similarity index 71% rename from spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarMixedTransactionTests.java rename to spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTemplateWithDbTransactionTests.java index 6321c9f7..2396a37c 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarMixedTransactionTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTemplateWithDbTransactionTests.java @@ -26,8 +26,8 @@ import org.springframework.context.annotation.Configuration; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.pulsar.PulsarException; import org.springframework.pulsar.core.PulsarTemplate; -import org.springframework.pulsar.transaction.PulsarMixedTransactionTests.PulsarProducerWithDbTransaction.PulsarProducerWithDbTransactionConfig; -import org.springframework.pulsar.transaction.PulsarMixedTransactionTests.PulsarProducerWithDbTransaction.PulsarProducerWithDbTransactionConfig.ProducerOnlyService; +import org.springframework.pulsar.transaction.PulsarTemplateWithDbTransactionTests.PulsarTemplateSynchronizedWithDbTransaction.PulsarTemplateSynchronizedWithDbTransactionConfig; +import org.springframework.pulsar.transaction.PulsarTemplateWithDbTransactionTests.PulsarTemplateSynchronizedWithDbTransaction.PulsarTemplateSynchronizedWithDbTransactionConfig.TestService; import org.springframework.stereotype.Service; import org.springframework.test.context.ContextConfiguration; import org.springframework.transaction.annotation.EnableTransactionManagement; @@ -35,39 +35,40 @@ import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.interceptor.TransactionAspectSupport; /** - * Tests for Pulsar transaction support with other resource transactions. + * Tests transaction support of {@link PulsarTemplate} when mixed with database + * transactions. * * @author Chris Bono */ -class PulsarMixedTransactionTests extends PulsarTxnWithDbTxnTestsBase { +class PulsarTemplateWithDbTransactionTests extends PulsarTxnWithDbTxnTestsBase { @Nested - @ContextConfiguration(classes = PulsarProducerWithDbTransactionConfig.class) - class PulsarProducerWithDbTransaction { + @ContextConfiguration(classes = PulsarTemplateSynchronizedWithDbTransactionConfig.class) + class PulsarTemplateSynchronizedWithDbTransaction { static final String topic = "ppwdbt-topic"; @Test - void whenDbTxnIsCommittedThenMessagesAreCommitted(@Autowired ProducerOnlyService producerService) { + void whenDbTxnIsCommittedThenMessagesAreCommitted(@Autowired TestService transactionalService) { var thing1 = new Thing(1L, "msg1"); - producerService.handleRequest(thing1, false, false); + transactionalService.handleRequest(thing1, false, false); assertThatMessagesAreInTopic(topic, thing1.name()); assertThatMessagesAreInDb(thing1); } @Test - void whenDbTxnIsSetRollbackOnlyThenMessagesAreNotCommitted(@Autowired ProducerOnlyService producerService) { + void whenDbTxnIsSetRollbackOnlyThenMessagesAreNotCommitted(@Autowired TestService transactionalService) { var thing2 = new Thing(2L, "msg2"); - producerService.handleRequest(thing2, true, false); + transactionalService.handleRequest(thing2, true, false); assertThatMessagesAreNotInTopic(topic, thing2.name()); assertThatMessagesAreNotInDb(thing2); } @Test - void whenServiceThrowsExceptionThenMessagesAreNotCommitted(@Autowired ProducerOnlyService producerService) { + void whenServiceThrowsExceptionThenMessagesAreNotCommitted(@Autowired TestService transactionalService) { var thing3 = new Thing(3L, "msg3"); assertThatExceptionOfType(PulsarException.class) - .isThrownBy(() -> producerService.handleRequest(thing3, false, true)) + .isThrownBy(() -> transactionalService.handleRequest(thing3, false, true)) .withMessage("Failed to commit due to chaos"); assertThatMessagesAreNotInTopic(topic, thing3.name()); assertThatMessagesAreNotInDb(thing3); @@ -75,10 +76,10 @@ class PulsarMixedTransactionTests extends PulsarTxnWithDbTxnTestsBase { @EnableTransactionManagement @Configuration - static class PulsarProducerWithDbTransactionConfig { + static class PulsarTemplateSynchronizedWithDbTransactionConfig { @Service - class ProducerOnlyService { + class TestService { @Autowired private JdbcTemplate jdbcTemplate;