Skip to content

Commit bfbd9e9

Browse files
Fix Infinite loop during application startup when a queue does not exist (awspring#1512) (awspring#1653)
1 parent 2f4e08c commit bfbd9e9

3 files changed

Lines changed: 219 additions & 1 deletion

File tree

Original file line numberDiff line numberDiff line change
@@ -0,0 +1,138 @@
1+
/*
2+
* Copyright 2013-2026 the original author or authors.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* https://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
package io.awspring.cloud.autoconfigure.sqs;
17+
18+
import static org.assertj.core.api.Assertions.assertThat;
19+
20+
import com.amazon.sqs.javamessaging.AmazonSQSExtendedAsyncClient;
21+
import io.awspring.cloud.autoconfigure.core.AwsAutoConfiguration;
22+
import io.awspring.cloud.autoconfigure.core.CredentialsProviderAutoConfiguration;
23+
import io.awspring.cloud.autoconfigure.core.RegionProviderAutoConfiguration;
24+
import io.awspring.cloud.sqs.annotation.SqsListener;
25+
import io.awspring.cloud.sqs.config.SqsBeanNames;
26+
import io.awspring.cloud.sqs.listener.DefaultListenerContainerRegistry;
27+
import java.util.concurrent.atomic.AtomicInteger;
28+
import org.junit.jupiter.api.Test;
29+
import org.springframework.boot.autoconfigure.AutoConfigurations;
30+
import org.springframework.boot.test.context.FilteredClassLoader;
31+
import org.springframework.boot.test.context.SpringBootTest;
32+
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
33+
import org.springframework.context.ApplicationContextException;
34+
import org.springframework.context.annotation.Bean;
35+
import org.springframework.context.annotation.Configuration;
36+
import org.testcontainers.junit.jupiter.Container;
37+
import org.testcontainers.junit.jupiter.Testcontainers;
38+
import org.testcontainers.localstack.LocalStackContainer;
39+
import org.testcontainers.shaded.org.bouncycastle.util.Arrays;
40+
import org.testcontainers.utility.DockerImageName;
41+
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials;
42+
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider;
43+
import software.amazon.awssdk.regions.Region;
44+
import software.amazon.awssdk.services.sqs.SqsAsyncClient;
45+
import software.amazon.awssdk.services.sqs.model.QueueDoesNotExistException;
46+
47+
/**
48+
* Integration tests for SQS listener startup.
49+
*
50+
* @author Bruno Garcia
51+
*/
52+
@Testcontainers
53+
@SpringBootTest
54+
class SqsListenerContainerStartupIntegrationTest {
55+
56+
private static final String EXISTING_QUEUE_NAME = "messaging-greetings-notifications";
57+
58+
private static final String MISSING_QUEUE_NAME = "not-existing-queue";
59+
60+
@Container
61+
static LocalStackContainer localstack = new LocalStackContainer(
62+
DockerImageName.parse("localstack/localstack:4.4.0"));
63+
64+
static {
65+
localstack.start();
66+
}
67+
68+
private static final String[] BASE_PARAMS = { "spring.cloud.aws.sqs.region=eu-west-1",
69+
"spring.cloud.aws.sqs.endpoint=" + localstack.getEndpoint(), "spring.cloud.aws.credentials.access-key=noop",
70+
"spring.cloud.aws.credentials.secret-key=noop", "spring.cloud.aws.region.static=eu-west-1" };
71+
72+
private static final AutoConfigurations BASE_CONFIGURATIONS = AutoConfigurations.of(
73+
RegionProviderAutoConfiguration.class, CredentialsProviderAutoConfiguration.class,
74+
SqsAutoConfiguration.class, AwsAutoConfiguration.class, ExistingAndMissingQueueListenerConfiguration.class);
75+
76+
private final ApplicationContextRunner applicationContextRunnerWithFailStrategy = new ApplicationContextRunner()
77+
.withClassLoader(new FilteredClassLoader(AmazonSQSExtendedAsyncClient.class))
78+
.withPropertyValues(Arrays.append(BASE_PARAMS, "spring.cloud.aws.sqs.queue-not-found-strategy=fail"))
79+
.withConfiguration(BASE_CONFIGURATIONS);
80+
81+
@Test
82+
void stopsRegistryWhenOneSqsListenerFailsToResolveQueueOnStartup() {
83+
createQueue(EXISTING_QUEUE_NAME);
84+
TrackingListenerContainerRegistry registry = new TrackingListenerContainerRegistry();
85+
86+
applicationContextRunnerWithFailStrategy.withBean(SqsBeanNames.ENDPOINT_REGISTRY_BEAN_NAME,
87+
TrackingListenerContainerRegistry.class, () -> registry).run(context -> {
88+
assertThat(context.getStartupFailure()).isInstanceOf(ApplicationContextException.class)
89+
.hasRootCauseInstanceOf(QueueDoesNotExistException.class);
90+
91+
assertThat(registry.stopInvocations).hasValue(1);
92+
assertThat(registry.getListenerContainers()).hasSize(2)
93+
.allSatisfy(container -> assertThat(container.isRunning())
94+
.as("Container %s should be stopped", container.getId()).isFalse());
95+
assertThat(registry.isRunning()).isFalse();
96+
});
97+
}
98+
99+
private static void createQueue(String queueName) {
100+
try (SqsAsyncClient client = SqsAsyncClient.builder()
101+
.credentialsProvider(StaticCredentialsProvider.create(AwsBasicCredentials.create("noop", "noop")))
102+
.region(Region.EU_WEST_1).endpointOverride(localstack.getEndpoint()).build()) {
103+
client.createQueue(request -> request.queueName(queueName)).join();
104+
}
105+
}
106+
107+
static class ExistingAndMissingQueueListener {
108+
109+
@SqsListener(EXISTING_QUEUE_NAME)
110+
void listenToExistingQueue(String message) {
111+
}
112+
113+
@SqsListener(MISSING_QUEUE_NAME)
114+
void listenToMissingQueue(String message) {
115+
}
116+
}
117+
118+
@Configuration
119+
static class ExistingAndMissingQueueListenerConfiguration {
120+
121+
@Bean
122+
ExistingAndMissingQueueListener existingAndMissingQueueListener() {
123+
return new ExistingAndMissingQueueListener();
124+
}
125+
}
126+
127+
static class TrackingListenerContainerRegistry extends DefaultListenerContainerRegistry {
128+
129+
private final AtomicInteger stopInvocations = new AtomicInteger();
130+
131+
@Override
132+
public void stop() {
133+
this.stopInvocations.incrementAndGet();
134+
super.stop();
135+
}
136+
}
137+
138+
}

spring-cloud-aws-sqs/src/main/java/io/awspring/cloud/sqs/listener/DefaultListenerContainerRegistry.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,10 +80,10 @@ public MessageListenerContainer<?> getContainerById(String id) {
8080
public void start() {
8181
synchronized (this.lifecycleMonitor) {
8282
logger.debug("Starting {}", getClass().getSimpleName());
83+
this.running = true;
8384
List<MessageListenerContainer<?>> containersToStart = this.listenerContainers.values().stream()
8485
.filter(SmartLifecycle::isAutoStartup).collect(Collectors.toList());
8586
LifecycleHandler.get().start(containersToStart);
86-
this.running = true;
8787
logger.debug("{} started", getClass().getSimpleName());
8888
}
8989
}

spring-cloud-aws-sqs/src/test/java/io/awspring/cloud/sqs/listener/DefaultListenerContainerRegistryTests.java

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,16 @@
1919
import static org.assertj.core.api.Assertions.assertThatThrownBy;
2020
import static org.mockito.BDDMockito.given;
2121
import static org.mockito.BDDMockito.then;
22+
import static org.mockito.BDDMockito.willAnswer;
2223
import static org.mockito.Mockito.mock;
2324
import static org.mockito.Mockito.times;
2425

26+
import java.util.concurrent.CompletionException;
27+
import java.util.concurrent.CountDownLatch;
28+
import java.util.concurrent.TimeUnit;
2529
import org.junit.jupiter.api.Test;
30+
import org.springframework.context.ApplicationContextException;
31+
import org.springframework.context.support.GenericApplicationContext;
2632

2733
/**
2834
* Tests for {@link DefaultListenerContainerRegistry}.
@@ -129,4 +135,78 @@ void shouldThrowIfIdAlreadyPresent() {
129135
.isInstanceOf(IllegalArgumentException.class);
130136
}
131137

138+
@Test
139+
void shouldRemainRunningAfterContainerFailsToStart() {
140+
CountDownLatch healthyContainerUpSignal = new CountDownLatch(1);
141+
MessageListenerContainer<Object> healthyContainer = mock(MessageListenerContainer.class);
142+
MessageListenerContainer<Object> faultyContainer = mock(MessageListenerContainer.class);
143+
String healthyId = "healthy-test-container-id";
144+
String faultyId = "faulty-test-container-id";
145+
given(healthyContainer.getId()).willReturn(healthyId);
146+
given(healthyContainer.isAutoStartup()).willReturn(true);
147+
given(faultyContainer.getId()).willReturn(faultyId);
148+
given(faultyContainer.isAutoStartup()).willReturn(true);
149+
willAnswer(invocation -> {
150+
healthyContainerUpSignal.countDown();
151+
return null;
152+
}).given(healthyContainer).start();
153+
String containerFailedToStartMessage = "Container failed to start";
154+
willAnswer(invocation -> {
155+
// Wait until the healthy container has started before failing the faulty container.
156+
if (!healthyContainerUpSignal.await(5, TimeUnit.SECONDS)) {
157+
throw new IllegalStateException("Test setup failed: healthy container start was not called");
158+
}
159+
throw new IllegalStateException(containerFailedToStartMessage);
160+
}).given(faultyContainer).start();
161+
DefaultListenerContainerRegistry registry = new DefaultListenerContainerRegistry();
162+
registry.registerListenerContainer(healthyContainer);
163+
registry.registerListenerContainer(faultyContainer);
164+
assertThatThrownBy(registry::start).isInstanceOf(CompletionException.class).cause()
165+
.isInstanceOf(RuntimeException.class).hasMessage(containerFailedToStartMessage);
166+
assertThat(registry.isRunning()).isTrue();
167+
registry.stop();
168+
then(healthyContainer).should(times(1)).stop();
169+
then(faultyContainer).should(times(1)).stop();
170+
assertThat(registry.isRunning()).isFalse();
171+
}
172+
173+
@Test
174+
void shouldBeStoppedBySpringLifecycleProcessorWhenOneContainerFailsToStart() {
175+
CountDownLatch healthyContainerUpSignal = new CountDownLatch(1);
176+
MessageListenerContainer<Object> healthyContainer = mock(MessageListenerContainer.class);
177+
MessageListenerContainer<Object> faultyContainer = mock(MessageListenerContainer.class);
178+
String healthyId = "healthy-test-container-id";
179+
String faultyId = "faulty-test-container-id";
180+
given(healthyContainer.getId()).willReturn(healthyId);
181+
given(healthyContainer.isAutoStartup()).willReturn(true);
182+
given(faultyContainer.getId()).willReturn(faultyId);
183+
given(faultyContainer.isAutoStartup()).willReturn(true);
184+
willAnswer(invocation -> {
185+
healthyContainerUpSignal.countDown();
186+
return null;
187+
}).given(healthyContainer).start();
188+
String containerFailedToStartMessage = "Container failed to start";
189+
willAnswer(invocation -> {
190+
if (!healthyContainerUpSignal.await(5, TimeUnit.SECONDS)) {
191+
throw new IllegalStateException("Test setup failed: healthy container start was not called");
192+
}
193+
throw new IllegalStateException(containerFailedToStartMessage);
194+
}).given(faultyContainer).start();
195+
DefaultListenerContainerRegistry registry = new DefaultListenerContainerRegistry();
196+
registry.registerListenerContainer(healthyContainer);
197+
registry.registerListenerContainer(faultyContainer);
198+
GenericApplicationContext context = new GenericApplicationContext();
199+
200+
try (context) {
201+
context.registerBean(DefaultListenerContainerRegistry.class, () -> registry);
202+
assertThatThrownBy(context::refresh).isInstanceOf(ApplicationContextException.class)
203+
.hasMessageContaining(DefaultListenerContainerRegistry.class.getName()).cause()
204+
.isInstanceOf(CompletionException.class).cause().isInstanceOf(IllegalStateException.class)
205+
.hasMessage(containerFailedToStartMessage);
206+
207+
then(healthyContainer).should(times(1)).stop();
208+
then(faultyContainer).should(times(1)).stop();
209+
assertThat(registry.isRunning()).isFalse();
210+
}
211+
}
132212
}

0 commit comments

Comments
 (0)