From 7835c2e69be81ce1e44d41451c0faadd5b2dda30 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 13 Aug 2026 00:22:47 +0900 Subject: [PATCH 1/4] test(cdc): reproduce duplicate connector authority overwrite --- .../cdc/spi/CdcRegistryIdentityTest.java | 100 ++++++++++++++++++ 1 file changed, 100 insertions(+) create mode 100644 cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java new file mode 100644 index 00000000..8d6ad0c1 --- /dev/null +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java @@ -0,0 +1,100 @@ +package com.xtrmetl.cdc.spi; + +import org.junit.jupiter.api.Test; + +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class CdcRegistryIdentityTest { + + @Test + void duplicateSourceConnectorIdsFailClosedInsteadOfReplacingRegistration() { + CdcSourceConnector first = source("duplicate-source"); + CdcSourceConnector second = source("duplicate-source"); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> new CdcSourceRegistry(List.of(first, second)) + ); + + assertEquals("Duplicate CDC source connector id: duplicate-source", failure.getMessage()); + } + + @Test + void duplicateTargetConnectorIdsFailClosedInsteadOfReplacingRegistration() { + CdcTargetRegistry registry = new CdcTargetRegistry(); + CdcTargetConnector duplicateKafka = target(KafkaCdcTargetConnector.ID); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> registry.register(duplicateKafka) + ); + + assertEquals("Duplicate CDC target connector id: kafka", failure.getMessage()); + } + + @Test + void nullSourceConnectorFailsBeforeRegistryMutation() { + CdcSourceRegistry registry = new CdcSourceRegistry(); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> registry.register(null) + ); + + assertEquals("CDC source connector must not be null", failure.getMessage()); + } + + @Test + void nullTargetConnectorFailsBeforeRegistryMutation() { + CdcTargetRegistry registry = new CdcTargetRegistry(); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> registry.register(null) + ); + + assertEquals("CDC target connector must not be null", failure.getMessage()); + } + + @Test + void blankSourceConnectorIdFailsBeforeRegistryMutation() { + CdcSourceConnector blank = source(" "); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> new CdcSourceRegistry(List.of(blank)) + ); + + assertEquals("CDC source connector id must not be blank", failure.getMessage()); + } + + @Test + void blankTargetConnectorIdFailsBeforeRegistryMutation() { + CdcTargetRegistry registry = new CdcTargetRegistry(); + CdcTargetConnector blank = target(""); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> registry.register(blank) + ); + + assertEquals("CDC target connector id must not be blank", failure.getMessage()); + } + + private static CdcSourceConnector source(String id) { + CdcSourceConnector connector = mock(CdcSourceConnector.class); + when(connector.id()).thenReturn(id); + return connector; + } + + private static CdcTargetConnector target(String id) { + CdcTargetConnector connector = mock(CdcTargetConnector.class); + when(connector.id()).thenReturn(id); + return connector; + } +} From 31be969bc7458783f908e502e54fee6fca305711 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 13 Aug 2026 03:13:33 +0900 Subject: [PATCH 2/4] fix(cdc): reject ambiguous connector identities --- .../xtrmetl/cdc/spi/CdcSourceRegistry.java | 48 ++++++++++++++++++- .../xtrmetl/cdc/spi/CdcTargetRegistry.java | 38 ++++++++++++++- 2 files changed, 83 insertions(+), 3 deletions(-) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java index fd00e11a..626ead5e 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java @@ -12,12 +12,21 @@ /** * Registry of CDC source connector types discovered as Spring beans, plus a safe * fallback for unit tests without a Spring context. + * + *

Connector identifiers are registration authority. Null connectors, blank + * identifiers, and duplicate identifiers are rejected before the registry is + * mutated so one implementation cannot silently replace another.

*/ @Component public class CdcSourceRegistry { private final Map byId = new LinkedHashMap<>(); + /** + * Builds a registry from Spring-discovered connector beans. + * + * @param connectors provider of discovered source connectors + */ public CdcSourceRegistry(ObjectProvider connectors) { connectors.orderedStream().forEach(this::register); if (byId.isEmpty()) { @@ -27,7 +36,10 @@ public CdcSourceRegistry(ObjectProvider connectors) { } /** - * Explicit list for tests. + * Builds a registry from an explicit connector list, primarily for tests and + * standalone embedding. + * + * @param connectors source connectors to register; a null list is treated as empty */ public CdcSourceRegistry(List connectors) { if (connectors != null) { @@ -38,18 +50,50 @@ public CdcSourceRegistry(List connectors) { } } + /** + * Builds a standalone registry with the built-in PostgreSQL/Debezium source. + */ public CdcSourceRegistry() { this(List.of()); } + /** + * Registers a CDC source connector without permitting ambiguous authority. + * + * @param connector connector to register + * @throws IllegalArgumentException when the connector is null, its identifier + * is blank, or its identifier is already registered + */ public final void register(CdcSourceConnector connector) { - byId.put(connector.id(), connector); + if (connector == null) { + throw new IllegalArgumentException("CDC source connector must not be null"); + } + + String id = connector.id(); + if (id == null || id.isBlank()) { + throw new IllegalArgumentException("CDC source connector id must not be blank"); + } + + if (byId.putIfAbsent(id, connector) != null) { + throw new IllegalArgumentException("Duplicate CDC source connector id: " + id); + } } + /** + * Finds a registered source connector by its exact identifier. + * + * @param id connector identifier + * @return the connector when registered, otherwise empty + */ public Optional find(String id) { return Optional.ofNullable(byId.get(id)); } + /** + * Returns all registered source connectors in registration order. + * + * @return registered source connectors + */ public Collection all() { return byId.values(); } diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java index 5d42b186..cbec1a45 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java @@ -9,25 +9,61 @@ /** * Registry of CDC target types (Kafka, JDBC replica, future warehouses). + * + *

Connector identifiers are registration authority. Null connectors, blank + * identifiers, and duplicate identifiers are rejected before the registry is + * mutated so one implementation cannot silently replace another.

*/ @Component public class CdcTargetRegistry { private final Map byId = new LinkedHashMap<>(); + /** + * Builds a target registry with the built-in Kafka and JDBC replica targets. + */ public CdcTargetRegistry() { register(new KafkaCdcTargetConnector()); register(new JdbcReplicaCdcTargetConnector()); } + /** + * Registers a CDC target connector without permitting ambiguous authority. + * + * @param connector connector to register + * @throws IllegalArgumentException when the connector is null, its identifier + * is blank, or its identifier is already registered + */ public final void register(CdcTargetConnector connector) { - byId.put(connector.id(), connector); + if (connector == null) { + throw new IllegalArgumentException("CDC target connector must not be null"); + } + + String id = connector.id(); + if (id == null || id.isBlank()) { + throw new IllegalArgumentException("CDC target connector id must not be blank"); + } + + if (byId.putIfAbsent(id, connector) != null) { + throw new IllegalArgumentException("Duplicate CDC target connector id: " + id); + } } + /** + * Finds a registered target connector by its exact identifier. + * + * @param id connector identifier + * @return the connector when registered, otherwise empty + */ public Optional find(String id) { return Optional.ofNullable(byId.get(id)); } + /** + * Returns all registered target connectors in registration order. + * + * @return registered target connectors + */ public Collection all() { return byId.values(); } From ad63482ec871395071d91e4bc167438cd60ceb48 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 13 Aug 2026 04:10:01 +0900 Subject: [PATCH 3/4] test(cdc): expose mutable registry collection authority --- .../cdc/spi/CdcRegistryIdentityTest.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java index 8d6ad0c1..6225287d 100644 --- a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java @@ -6,6 +6,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -86,6 +87,23 @@ void blankTargetConnectorIdFailsBeforeRegistryMutation() { assertEquals("CDC target connector id must not be blank", failure.getMessage()); } + @Test + void sourceConnectorCollectionCannotDeleteRegistrationAuthority() { + CdcSourceRegistry registry = new CdcSourceRegistry(List.of(source("immutable-source"))); + + assertThrows(UnsupportedOperationException.class, () -> registry.all().clear()); + assertTrue(registry.find("immutable-source").isPresent()); + } + + @Test + void targetConnectorCollectionCannotDeleteRegistrationAuthority() { + CdcTargetRegistry registry = new CdcTargetRegistry(); + + assertThrows(UnsupportedOperationException.class, () -> registry.all().clear()); + assertTrue(registry.find(KafkaCdcTargetConnector.ID).isPresent()); + assertTrue(registry.find(JdbcReplicaCdcTargetConnector.ID).isPresent()); + } + private static CdcSourceConnector source(String id) { CdcSourceConnector connector = mock(CdcSourceConnector.class); when(connector.id()).thenReturn(id); From 34f85a7151471cd1ae018ada0b409d2bd7bcfd6c Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 13 Aug 2026 04:14:14 +0900 Subject: [PATCH 4/4] fix(cdc): protect registry connector authority --- .../main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java | 7 ++++--- .../main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java | 7 ++++--- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java index 626ead5e..78edf83f 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceRegistry.java @@ -4,6 +4,7 @@ import org.springframework.stereotype.Component; import java.util.Collection; +import java.util.Collections; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -90,11 +91,11 @@ public Optional find(String id) { } /** - * Returns all registered source connectors in registration order. + * Returns an unmodifiable live view of all registered source connectors in registration order. * - * @return registered source connectors + * @return unmodifiable registered source connectors */ public Collection all() { - return byId.values(); + return Collections.unmodifiableCollection(byId.values()); } } diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java index cbec1a45..cbe8fff1 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcTargetRegistry.java @@ -3,6 +3,7 @@ import org.springframework.stereotype.Component; import java.util.Collection; +import java.util.Collections; import java.util.LinkedHashMap; import java.util.Map; import java.util.Optional; @@ -60,11 +61,11 @@ public Optional find(String id) { } /** - * Returns all registered target connectors in registration order. + * Returns an unmodifiable live view of all registered target connectors in registration order. * - * @return registered target connectors + * @return unmodifiable registered target connectors */ public Collection all() { - return byId.values(); + return Collections.unmodifiableCollection(byId.values()); } }