From 5a087ca32b319dd2ebb8708fffe565e8e42eda0b Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 06:12:46 +0900 Subject: [PATCH 1/3] test(cdc): reject invalid Kafka retry settings --- .../xtrmetl/cdc/config/KafkaConfigTest.java | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java index 16d9f959..0efd35f8 100644 --- a/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java @@ -18,6 +18,7 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; @@ -53,6 +54,32 @@ void configuresErrorHandlerWithDeadLetterRecovererAndBackOff() { assertFalse(classifier.classify(new IllegalStateException("test"))); } + @Test + void rejectsNegativeRetryBackoffBeforeBuildingErrorHandler() { + KafkaConfig config = new KafkaConfig(); + KafkaTemplate kafkaTemplate = mock(KafkaTemplate.class); + + IllegalArgumentException exception = assertThrows( + IllegalArgumentException.class, + () -> config.kafkaListenerErrorHandler(kafkaTemplate, -1L, 30L) + ); + + assertTrue(exception.getMessage().contains("xtrmetl.replica.kafka.retry-backoff-ms")); + } + + @Test + void rejectsNegativeRetryAttemptsBeforeBuildingErrorHandler() { + KafkaConfig config = new KafkaConfig(); + KafkaTemplate kafkaTemplate = mock(KafkaTemplate.class); + + IllegalArgumentException exception = assertThrows( + IllegalArgumentException.class, + () -> config.kafkaListenerErrorHandler(kafkaTemplate, 1000L, -1L) + ); + + assertTrue(exception.getMessage().contains("xtrmetl.replica.kafka.retry-max-attempts")); + } + @Test void configuresListenerFactoryWithRecordAckModeAndCommonErrorHandler() { KafkaConfig config = new KafkaConfig(); From 8c87b7835af77a2082ab5d7ae3fcc3b9bba3dc69 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 06:15:16 +0900 Subject: [PATCH 2/3] fix(cdc): fail fast on invalid Kafka retry settings --- .../com/xtrmetl/cdc/config/KafkaConfig.java | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java b/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java index 5c77a06b..d6946a80 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/config/KafkaConfig.java @@ -23,13 +23,31 @@ public class KafkaConfig { private static final Logger log = LoggerFactory.getLogger(KafkaConfig.class); private static final int MAX_CONCURRENCY = 32; + private static final String RETRY_BACKOFF_KEY = "xtrmetl.replica.kafka.retry-backoff-ms"; + private static final String RETRY_MAX_ATTEMPTS_KEY = "xtrmetl.replica.kafka.retry-max-attempts"; + /** + * Builds the replica-listener error handler with terminal dead-letter recovery. + * + *

Retry settings are deployment-owned, so this boundary rejects negative values before + * delegating to Spring's {@link FixedBackOff}. Zero is valid: a zero interval retries + * immediately, and zero maximum attempts sends a retryable failure directly to recovery.

+ * + * @param kafkaTemplate template used to publish exhausted records to the dead-letter topic + * @param retryBackoffMs fixed delay between retry attempts in milliseconds; must be non-negative + * @param retryMaxAttempts maximum retry attempts after the original delivery; must be non-negative + * @return configured listener error handler + * @throws IllegalArgumentException when either retry setting is negative + */ @Bean public DefaultErrorHandler kafkaListenerErrorHandler( @NonNull KafkaTemplate kafkaTemplate, @Value("${xtrmetl.replica.kafka.retry-backoff-ms:1000}") long retryBackoffMs, @Value("${xtrmetl.replica.kafka.retry-max-attempts:30}") long retryMaxAttempts ) { + requireNonNegative(RETRY_BACKOFF_KEY, retryBackoffMs); + requireNonNegative(RETRY_MAX_ATTEMPTS_KEY, retryMaxAttempts); + DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer( kafkaTemplate, (record, ex) -> new TopicPartition(record.topic() + ".DLT", record.partition()) @@ -43,6 +61,14 @@ public DefaultErrorHandler kafkaListenerErrorHandler( return errorHandler; } + /** + * Builds the replica Kafka listener factory with record-level acknowledgement semantics. + * + * @param consumerFactory Kafka consumer factory for replica records + * @param kafkaListenerErrorHandler bounded retry and dead-letter error handler + * @param concurrency requested listener concurrency + * @return configured listener container factory + */ @Bean public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory( @NonNull ConsumerFactory consumerFactory, @@ -77,4 +103,10 @@ public ConcurrentKafkaListenerContainerFactory kafkaListenerCont factory.setCommonErrorHandler(kafkaListenerErrorHandler); return factory; } + + private static void requireNonNegative(String key, long value) { + if (value < 0) { + throw new IllegalArgumentException(key + " must be greater than or equal to 0"); + } + } } From 13d3bf6372abf8ae6cf214e3d5728aceaace2101 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 06:15:52 +0900 Subject: [PATCH 3/3] test(cdc): preserve zero Kafka retry semantics --- .../com/xtrmetl/cdc/config/KafkaConfigTest.java | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java index 0efd35f8..4c116b14 100644 --- a/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/config/KafkaConfigTest.java @@ -54,6 +54,21 @@ void configuresErrorHandlerWithDeadLetterRecovererAndBackOff() { assertFalse(classifier.classify(new IllegalStateException("test"))); } + @Test + void preservesZeroRetrySettings() { + KafkaConfig config = new KafkaConfig(); + KafkaTemplate kafkaTemplate = mock(KafkaTemplate.class); + + DefaultErrorHandler errorHandler = config.kafkaListenerErrorHandler(kafkaTemplate, 0L, 0L); + + Object failureTracker = ReflectionTestUtils.getField(errorHandler, "failureTracker"); + assertNotNull(failureTracker); + FixedBackOff fixedBackOff = (FixedBackOff) ReflectionTestUtils.getField(failureTracker, "backOff"); + assertNotNull(fixedBackOff); + assertEquals(0L, fixedBackOff.getInterval()); + assertEquals(0L, fixedBackOff.getMaxAttempts()); + } + @Test void rejectsNegativeRetryBackoffBeforeBuildingErrorHandler() { KafkaConfig config = new KafkaConfig();