Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,10 @@ protected ConcurrentPulsarMessageListenerContainer<T> createContainerInstance(Pu
containerProps.setTopicResolver(factoryProps.getTopicResolver());
containerProps.setSubscriptionType(factoryProps.getSubscriptionType());
containerProps.setSubscriptionName(factoryProps.getSubscriptionName());
// Set the retry template before the policy as setting a non-null template
// implicitly switches the policy to RETRY
containerProps.setStartupFailureRetryTemplate(factoryProps.getStartupFailureRetryTemplate());
containerProps.setStartupFailurePolicy(factoryProps.getStartupFailurePolicy());
var factoryTxnProps = factoryProps.transactions();
var containerTxnProps = containerProps.transactions();
containerTxnProps.setEnabled(factoryTxnProps.isEnabled());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;

import org.springframework.core.retry.RetryPolicy;
import org.springframework.core.retry.RetryTemplate;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.pulsar.listener.PulsarContainerProperties;
Expand Down Expand Up @@ -151,4 +153,58 @@ void factoryPropsUsedWhenSpecified() {

}

@SuppressWarnings("unchecked")
@Nested
class StartupFailurePolicyFrom {

@Test
void factoryPropsUsedWhenSpecified() {
var factoryProps = new PulsarContainerProperties();
factoryProps.setStartupFailurePolicy(StartupFailurePolicy.CONTINUE);
var containerFactory = new ConcurrentPulsarListenerContainerFactory<String>(
mock(PulsarConsumerFactory.class), factoryProps);
var endpoint = mock(PulsarListenerEndpoint.class);
when(endpoint.getConcurrency()).thenReturn(1);
var createdContainer = containerFactory.createRegisteredContainer(endpoint);
assertThat(createdContainer.getContainerProperties().getStartupFailurePolicy())
.isEqualTo(StartupFailurePolicy.CONTINUE);
}

@Test
void defaultUsedWhenNotSetOnFactoryProps() {
var factoryProps = new PulsarContainerProperties();
var containerFactory = new ConcurrentPulsarListenerContainerFactory<String>(
mock(PulsarConsumerFactory.class), factoryProps);
var endpoint = mock(PulsarListenerEndpoint.class);
when(endpoint.getConcurrency()).thenReturn(1);
var createdContainer = containerFactory.createRegisteredContainer(endpoint);
assertThat(createdContainer.getContainerProperties().getStartupFailurePolicy())
.isEqualTo(StartupFailurePolicy.STOP);
}

}

@SuppressWarnings("unchecked")
@Nested
class StartupFailureRetryTemplateFrom {

@Test
void factoryPropsUsedWhenSpecified() {
var factoryProps = new PulsarContainerProperties();
var retryTemplate = new RetryTemplate(RetryPolicy.builder().maxRetries(3).build());
factoryProps.setStartupFailureRetryTemplate(retryTemplate);
var containerFactory = new ConcurrentPulsarListenerContainerFactory<String>(
mock(PulsarConsumerFactory.class), factoryProps);
var endpoint = mock(PulsarListenerEndpoint.class);
when(endpoint.getConcurrency()).thenReturn(1);
var createdContainer = containerFactory.createRegisteredContainer(endpoint);
assertThat(createdContainer.getContainerProperties().getStartupFailureRetryTemplate())
.isSameAs(retryTemplate);
// Setting the retry template implicitly switches the policy to RETRY
assertThat(createdContainer.getContainerProperties().getStartupFailurePolicy())
.isEqualTo(StartupFailurePolicy.RETRY);
}

}

}