From fbff39b62b68bd76f1384f19868c383d904946b1 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Tue, 23 Apr 2024 14:22:09 -0500 Subject: [PATCH] Add tests for transactions package See #661 --- ...ltPulsarMessageListenerContainerTests.java | 4 +- .../listener/TransactionSettingsTests.java | 63 ++++++++ .../PulsarResourceHolderTests.java | 55 +++++++ .../PulsarResourceSynchronizationTests.java | 63 ++++++++ .../PulsarTransactionManagerTests.java | 119 ++++++++++++++ .../PulsarTransactionUtilsTests.java | 153 ++++++++++++++++++ 6 files changed, 456 insertions(+), 1 deletion(-) create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/listener/TransactionSettingsTests.java create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceHolderTests.java create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceSynchronizationTests.java create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionManagerTests.java create mode 100644 spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionUtilsTests.java diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index 0fab6435..9a65e660 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -54,6 +54,7 @@ import org.springframework.pulsar.core.DefaultPulsarConsumerFactory; import org.springframework.pulsar.core.DefaultPulsarProducerFactory; import org.springframework.pulsar.core.PulsarTemplate; import org.springframework.pulsar.test.support.PulsarTestContainerSupport; +import org.springframework.pulsar.transaction.PulsarAwareTransactionManager; import org.springframework.test.util.ReflectionTestUtils; /** @@ -405,6 +406,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS var containerProps = new PulsarContainerProperties(); containerProps.setSchema(Schema.STRING); containerProps.transactions().setEnabled(true); + containerProps.transactions().setTransactionManager(mock(PulsarAwareTransactionManager.class)); containerProps.setBatchListener(true); containerProps.setAckMode(AckMode.RECORD); containerProps.setMessageListener((PulsarBatchMessageListener) (consumer, msgs) -> { @@ -413,7 +415,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS var consumerFactory = new DefaultPulsarConsumerFactory(mock(PulsarClient.class), List.of()); var container = new DefaultPulsarMessageListenerContainer<>(consumerFactory, containerProps); assertThatIllegalStateException().isThrownBy(() -> container.start()) - .withMessage("Batch record listeners do not support AckMode.RECORD"); + .withMessage("Transactional batch listeners do not support AckMode.RECORD"); } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/TransactionSettingsTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/TransactionSettingsTests.java new file mode 100644 index 00000000..47d0880a --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/TransactionSettingsTests.java @@ -0,0 +1,63 @@ +/* + * 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.listener; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.time.Duration; + +import org.junit.jupiter.api.Test; + +import org.springframework.pulsar.listener.PulsarContainerProperties.TransactionSettings; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.support.DefaultTransactionDefinition; + +/** + * Unit tests for {@link TransactionSettings}. + * + * @author Chris Bono + */ +class TransactionSettingsTests { + + @Test + void whenTimeoutNotSetThenReturnsConfiguredDefinition() { + var txnSettings = new TransactionSettings(); + var txnDefinition = new DefaultTransactionDefinition(); + txnSettings.setTransactionDefinition(txnDefinition); + assertThat(txnSettings.determineTransactionDefinition()).isSameAs(txnDefinition); + } + + @Test + void whenTimeoutSetButDefinitionNotSetThenReturnsNewDefinitionWithTimeout() { + var txnSettings = new TransactionSettings(); + txnSettings.setTimeout(Duration.ofSeconds(100)); + assertThat(txnSettings.determineTransactionDefinition()).extracting(TransactionDefinition::getTimeout) + .isEqualTo(100); + } + + @Test + void whenTimeoutSetAndDefinitionSetThenReturnsCloneDefinitionUpdatedWithTimeout() { + var txnSettings = new TransactionSettings(); + txnSettings.setTimeout(Duration.ofSeconds(200)); + var txnDefinition = new DefaultTransactionDefinition(); + txnDefinition.setTimeout(100); + txnSettings.setTransactionDefinition(txnDefinition); + assertThat(txnSettings.determineTransactionDefinition()).extracting(TransactionDefinition::getTimeout) + .isEqualTo(200); + } + +} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceHolderTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceHolderTests.java new file mode 100644 index 00000000..6ee00783 --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceHolderTests.java @@ -0,0 +1,55 @@ +/* + * 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.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.util.concurrent.CompletableFuture; + +import org.apache.pulsar.client.api.transaction.Transaction; +import org.junit.jupiter.api.Test; + +/** + * Tests for {@link PulsarResourceHolder}. + * + * @author Chris Bono + */ +class PulsarResourceHolderTests { + + @Test + void rollbackAbortsTransaction() { + var txn = mock(Transaction.class); + when(txn.abort()).thenReturn(CompletableFuture.completedFuture(null)); + var holder = new PulsarResourceHolder(txn); + holder.rollback(); + verify(txn).abort(); + } + + @Test + void multipleCommitCallsCommitsTransactionOnce() { + var txn = mock(Transaction.class); + when(txn.commit()).thenReturn(CompletableFuture.completedFuture(null)); + var holder = new PulsarResourceHolder(txn); + holder.commit(); + holder.commit(); + holder.commit(); + verify(txn).commit(); + } + +} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceSynchronizationTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceSynchronizationTests.java new file mode 100644 index 00000000..40730ea4 --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarResourceSynchronizationTests.java @@ -0,0 +1,63 @@ +/* + * 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.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; + +import org.apache.pulsar.client.api.PulsarClient; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import org.springframework.transaction.support.TransactionSynchronization; + +/** + * Tests for {@link PulsarResourceSynchronization}. + * + * @author Chris Bono + */ +class PulsarResourceSynchronizationTests { + + private final PulsarClient pulsarClient = mock(PulsarClient.class); + + @Test + void processResourceAfterCommitDoesCommitOnResourceHolder() { + var holder = mock(PulsarResourceHolder.class); + var sync = new PulsarResourceSynchronization(holder, pulsarClient); + sync.processResourceAfterCommit(holder); + verify(holder).commit(); + } + + @Test + void afterCompletionDoesCommitOnHolderWhenTxnStatusIsCommitted() { + var holder = mock(PulsarResourceHolder.class); + var sync = new PulsarResourceSynchronization(holder, pulsarClient); + sync.afterCompletion(TransactionSynchronization.STATUS_COMMITTED); + verify(holder).commit(); + } + + @ParameterizedTest + @ValueSource(ints = { TransactionSynchronization.STATUS_ROLLED_BACK, TransactionSynchronization.STATUS_UNKNOWN }) + void afterCompletionDoesRollbackOnHolderWhenTxnStatusIsNotCommitted(int status) { + var holder = mock(PulsarResourceHolder.class); + var sync = new PulsarResourceSynchronization(holder, pulsarClient); + sync.afterCompletion(status); + verify(holder).rollback(); + } + +} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionManagerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionManagerTests.java new file mode 100644 index 00000000..e7bd184d --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionManagerTests.java @@ -0,0 +1,119 @@ +/* + * 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 static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import org.apache.pulsar.client.api.PulsarClient; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.support.DefaultTransactionStatus; +import org.springframework.transaction.support.TransactionSynchronizationManager; + +/** + * Tests for {@link PlatformTransactionManager}. + * + * @author Chris Bono + */ +class PulsarTransactionManagerTests { + + private PulsarClient pulsarClient = mock(PulsarClient.class); + + private PulsarTransactionManager transactionManager; + + private PulsarResourceHolder resourceHolder; + + private PulsarTransactionObject transactionObject; + + private DefaultTransactionStatus transactionStatus; + + @BeforeEach + void prepareForTest() { + transactionManager = new PulsarTransactionManager(pulsarClient); + resourceHolder = mock(PulsarResourceHolder.class); + transactionObject = new PulsarTransactionObject(); + transactionObject.setResourceHolder(resourceHolder); + transactionStatus = mock(DefaultTransactionStatus.class); + when(transactionStatus.getTransaction()).thenReturn(transactionObject); + } + + @Test + void doGetTransactionReturnsPulsarTxnObject() { + TransactionSynchronizationManager.bindResource(this.pulsarClient, resourceHolder); + assertThat(transactionManager.doGetTransaction()).isInstanceOf(PulsarTransactionObject.class) + .hasFieldOrPropertyWithValue("resourceHolder", resourceHolder); + } + + @Test + void isExistingTransactionReturnsTrueWhenTxnObjectHasResourceHolder() { + var txnObject = new PulsarTransactionObject(); + txnObject.setResourceHolder(resourceHolder); + assertThat(transactionManager.isExistingTransaction(txnObject)).isTrue(); + } + + @Test + void isExistingTransactionReturnsFalseWhenTxnObjectHasNoResourceHolder() { + var txnObject = new PulsarTransactionObject(); + assertThat(transactionManager.isExistingTransaction(txnObject)).isFalse(); + } + + @Test + void doSuspendUnbindsAndNullsOutResourceHolder() { + TransactionSynchronizationManager.bindResource(this.pulsarClient, resourceHolder); + transactionManager.doSuspend(transactionObject); + assertThat(transactionObject.getResourceHolder()).isNull(); + assertThat(TransactionSynchronizationManager.getResource(this.pulsarClient)).isNull(); + } + + @Test + void doResumeBindsResourceHolder() { + transactionManager.doResume("unused", resourceHolder); + assertThat(TransactionSynchronizationManager.getResource(this.pulsarClient)).isSameAs(resourceHolder); + } + + @Test + void doCommitDoesCommitOnResourceHolder() { + transactionManager.doCommit(transactionStatus); + verify(resourceHolder).commit(); + } + + @Test + void doRollbackDoesRollbackOnResourceHolder() { + transactionManager.doRollback(transactionStatus); + verify(resourceHolder).rollback(); + } + + @Test + void doSetRollbackOnlyDoesSetRollbackOnlyOnResourceHolder() { + transactionManager.doSetRollbackOnly(transactionStatus); + verify(resourceHolder).setRollbackOnly(); + } + + @Test + void doCleanupDoesUnbindAndClearResourceHolder() { + TransactionSynchronizationManager.bindResource(this.pulsarClient, resourceHolder); + transactionManager.doCleanupAfterCompletion(transactionObject); + assertThat(TransactionSynchronizationManager.getResource(this.pulsarClient)).isNull(); + verify(transactionObject.getResourceHolder()).clear(); + } + +} diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionUtilsTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionUtilsTests.java new file mode 100644 index 00000000..fd18ca47 --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/transaction/PulsarTransactionUtilsTests.java @@ -0,0 +1,153 @@ +/* + * 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 static org.assertj.core.api.Assertions.assertThatIllegalArgumentException; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.time.Duration; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; + +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.transaction.Transaction; +import org.apache.pulsar.client.api.transaction.TransactionBuilder; +import org.assertj.core.data.Offset; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; + +import org.springframework.transaction.support.TransactionSynchronizationManager; + +/** + * Tests for {@link PulsarTransactionUtils}. + * + * @author Chris Bono + */ +class PulsarTransactionUtilsTests { + + private PulsarClient pulsarClient = mock(PulsarClient.class); + + @Nested + class InTransaction { + + @Test + void whenNoResourceThenReturnsFalse() { + assertThat(PulsarTransactionUtils.inTransaction(pulsarClient)).isFalse(); + } + + @Test + void whenResourceThenReturnsTrue() { + TransactionSynchronizationManager.bindResource(pulsarClient, "some-fake-txn-object"); + assertThat(PulsarTransactionUtils.inTransaction(pulsarClient)).isTrue(); + } + + @Nested + class WithActualTransactionActive { + + // NOTE: Because this test sets the thread local 'actualTransactionActive' + // which interferes w/ the other InTransaction tests it is nested so that it + // executes after the other tests. + @Test + void whenNoResourceThenReturnsTrue() { + TransactionSynchronizationManager.setActualTransactionActive(true); + assertThat(PulsarTransactionUtils.inTransaction(pulsarClient)).isTrue(); + } + + } + + } + + @Nested + class GetResourceHolder { + + @Test + void whenNoResourceThenReturnsNull() { + assertThat(PulsarTransactionUtils.getResourceHolder(pulsarClient)).isNull(); + } + + @Test + void whenResourceThenReturnsResource() { + var resourceHolder = new PulsarResourceHolder(mock(Transaction.class)); + TransactionSynchronizationManager.bindResource(pulsarClient, resourceHolder); + assertThat(PulsarTransactionUtils.getResourceHolder(pulsarClient)).isEqualTo(resourceHolder); + } + + } + + @Nested + class ObtainResourceHolder { + + @Test + void whenResourceThenReturnsResource() { + var resourceHolder = new PulsarResourceHolder(mock(Transaction.class)); + TransactionSynchronizationManager.bindResource(pulsarClient, resourceHolder); + assertThat(PulsarTransactionUtils.obtainResourceHolder(pulsarClient, null)).isEqualTo(resourceHolder); + } + + @Test + void whenNoResourceThenCreatesResourceWithoutTimeout() { + var txn = mock(Transaction.class); + var txnBuilder = mock(TransactionBuilder.class); + when(txnBuilder.build()).thenReturn(CompletableFuture.completedFuture(txn)); + when(pulsarClient.newTransaction()).thenReturn(txnBuilder); + assertThat(PulsarTransactionUtils.obtainResourceHolder(pulsarClient, null)) + .extracting(PulsarResourceHolder::getTransaction) + .isEqualTo(txn); + assertThat(PulsarTransactionUtils.getResourceHolder(pulsarClient)) + .extracting(PulsarResourceHolder::getTransaction) + .isEqualTo(txn); + } + + @Test + void whenNoResourceThenCreatesResourceWithTimeout() { + var txn = mock(Transaction.class); + var txnBuilder = mock(TransactionBuilder.class); + when(txnBuilder.build()).thenReturn(CompletableFuture.completedFuture(txn)); + when(pulsarClient.newTransaction()).thenReturn(txnBuilder); + long nowEpochMillis = System.currentTimeMillis(); + var resourceHolder = PulsarTransactionUtils.obtainResourceHolder(pulsarClient, Duration.ofSeconds(60)); + assertThat(resourceHolder.getTransaction()).isEqualTo(txn); + assertThat(resourceHolder.hasTimeout()).isTrue(); + long timeoutEpochMillis = resourceHolder.getDeadline().getTime(); + assertThat(nowEpochMillis + 60_000).isCloseTo(timeoutEpochMillis, Offset.offset(500L)); + verify(txnBuilder).withTransactionTimeout(61, TimeUnit.SECONDS); + } + + } + + @Nested + class Abort { + + @Test + void whenTransactionIsNotNullThenTxnIsAborted() { + var txn = mock(Transaction.class); + when(txn.abort()).thenReturn(CompletableFuture.completedFuture(null)); + PulsarTransactionUtils.abort(txn); + verify(txn).abort(); + } + + @Test + void whenTransactionIsNullThenThrowsException() { + assertThatIllegalArgumentException().isThrownBy(() -> PulsarTransactionUtils.abort(null)); + } + + } + +}