Add tests for transactions package

See #661
This commit is contained in:
Chris Bono
2024-04-23 14:22:09 -05:00
parent d92deb3da6
commit fbff39b62b
6 changed files with 456 additions and 1 deletions

View File

@@ -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<String>(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");
}
}

View File

@@ -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);
}
}

View File

@@ -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();
}
}

View File

@@ -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();
}
}

View File

@@ -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();
}
}

View File

@@ -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));
}
}
}