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..1a4e2ceb 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 @@ -1,23 +1,36 @@ package com.xtrmetl.cdc.spi; import org.springframework.beans.factory.ObjectProvider; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import java.util.Collection; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Optional; /** * Registry of CDC source connector types discovered as Spring beans, plus a safe * fallback for unit tests without a Spring context. + * + *

Connector identifiers are configuration authority. Registration therefore fails closed for + * null connectors, blank identifiers, and duplicate identifiers instead of allowing bean order to + * replace the selected implementation.

*/ @Component public class CdcSourceRegistry { private final Map byId = new LinkedHashMap<>(); + /** + * Creates a registry from Spring-discovered source connectors. + * + * @param connectors ordered provider of source connector beans + * @throws IllegalArgumentException when a discovered connector has an invalid or duplicate id + */ + @Autowired public CdcSourceRegistry(ObjectProvider connectors) { connectors.orderedStream().forEach(this::register); if (byId.isEmpty()) { @@ -27,7 +40,10 @@ public CdcSourceRegistry(ObjectProvider connectors) { } /** - * Explicit list for tests. + * Creates a registry from an explicit connector list, primarily for standalone use and tests. + * + * @param connectors source connectors to register; a null list means no explicit connectors + * @throws IllegalArgumentException when a connector has an invalid or duplicate id */ public CdcSourceRegistry(List connectors) { if (connectors != null) { @@ -38,19 +54,50 @@ public CdcSourceRegistry(List connectors) { } } + /** + * Creates a registry containing the built-in PostgreSQL Debezium source connector. + */ public CdcSourceRegistry() { this(List.of()); } + /** + * Registers one source connector without allowing existing configuration identity to be replaced. + * + * @param connector source connector to register + * @throws IllegalArgumentException when the connector is null, its id is blank, or its id 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 = Objects.requireNonNullElse(connector.id(), ""); + if (id.isBlank()) { + throw new IllegalArgumentException("CDC source connector id must not be blank"); + } + CdcSourceConnector previous = byId.putIfAbsent(id, connector); + if (previous != null) { + throw new IllegalArgumentException("Duplicate CDC source connector id: " + id); + } } + /** + * Finds a source connector by its exact configuration identifier. + * + * @param id exact connector identifier + * @return the registered connector, or empty when the identifier is unknown + */ public Optional find(String id) { return Optional.ofNullable(byId.get(id)); } + /** + * Returns an immutable insertion-ordered snapshot of registered source connectors. + * + * @return immutable connector collection detached from registry mutation authority + */ public Collection all() { - return byId.values(); + return List.copyOf(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..746c00db 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 @@ -4,31 +4,64 @@ import java.util.Collection; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Optional; /** - * Registry of CDC target types (Kafka, JDBC replica, future warehouses). + * Registry of CDC target connector types such as Kafka and JDBC replica targets. + * + *

Connector identifiers are configuration authority. Invalid registration is rejected so + * registration order cannot silently replace a selected connector implementation.

*/ @Component public class CdcTargetRegistry { private final Map byId = new LinkedHashMap<>(); + /** Creates a registry containing the built-in Kafka and JDBC replica target connectors. */ public CdcTargetRegistry() { register(new KafkaCdcTargetConnector()); register(new JdbcReplicaCdcTargetConnector()); } + /** + * Registers one target connector without replacing an existing connector with the same id. + * + * @param connector target connector to register + * @throws IllegalArgumentException when the connector is null, its id is blank, or its id 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 = Objects.requireNonNullElse(connector.id(), ""); + if (id.isBlank()) { + throw new IllegalArgumentException("CDC target connector id must not be blank"); + } + CdcTargetConnector previous = byId.putIfAbsent(id, connector); + if (previous != null) { + throw new IllegalArgumentException("Duplicate CDC target connector id: " + id); + } } + /** + * Finds a target connector by its exact configuration identifier. + * + * @param id exact connector identifier + * @return the registered connector, or empty when the identifier is unknown + */ public Optional find(String id) { return Optional.ofNullable(byId.get(id)); } + /** + * Returns an immutable insertion-ordered snapshot of registered target connectors. + * + * @return immutable connector collection detached from registry mutation authority + */ public Collection all() { - return byId.values(); + return List.copyOf(byId.values()); } } 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..70c89562 --- /dev/null +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcRegistryIdentityTest.java @@ -0,0 +1,128 @@ +package com.xtrmetl.cdc.spi; + +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.beans.factory.annotation.Autowired; + +import java.lang.reflect.Constructor; +import java.util.Arrays; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +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; +import static org.mockito.Mockito.when; + +/** + * Fail-first contract for CDC connector registration authority. + * + *

Connector identifiers select production implementations. Invalid registration must fail before + * registry mutation so bean order or plugin code cannot silently replace or remove that authority.

+ */ +class CdcRegistryIdentityTest { + + @Test + void duplicateSourceConnectorIdsFailClosedInsteadOfReplacingRegistration() { + CdcSourceConnector first = source("duplicate-source"); + CdcSourceConnector second = source("duplicate-source"); + CdcSourceRegistry registry = new CdcSourceRegistry(List.of(first)); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> registry.register(second) + ); + + assertEquals("Duplicate CDC source connector id: duplicate-source", failure.getMessage()); + assertSame(first, registry.find("duplicate-source").orElseThrow()); + } + + @Test + void duplicateTargetConnectorIdsFailClosedInsteadOfReplacingRegistration() { + CdcTargetRegistry registry = new CdcTargetRegistry(); + CdcTargetConnector originalKafka = registry.find(KafkaCdcTargetConnector.ID).orElseThrow(); + CdcTargetConnector duplicateKafka = target(KafkaCdcTargetConnector.ID); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> registry.register(duplicateKafka) + ); + + assertEquals("Duplicate CDC target connector id: kafka", failure.getMessage()); + assertSame(originalKafka, registry.find(KafkaCdcTargetConnector.ID).orElseThrow()); + } + + @Test + void springDiscoveryConstructorIsExplicitlyAutowired() { + Constructor discoveryConstructor = Arrays.stream(CdcSourceRegistry.class.getConstructors()) + .filter(constructor -> Arrays.equals( + constructor.getParameterTypes(), + new Class[]{ObjectProvider.class} + )) + .findFirst() + .orElseThrow(); + + assertTrue(discoveryConstructor.isAnnotationPresent(Autowired.class), + "Spring discovery constructor must be explicitly selected when other public constructors exist"); + } + + @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()); + } + + @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); + return connector; + } + + private static CdcTargetConnector target(String id) { + CdcTargetConnector connector = mock(CdcTargetConnector.class); + when(connector.id()).thenReturn(id); + return connector; + } +} diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceRegistrySpringWiringTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceRegistrySpringWiringTest.java new file mode 100644 index 00000000..cb0e4f7d --- /dev/null +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceRegistrySpringWiringTest.java @@ -0,0 +1,71 @@ +package com.xtrmetl.cdc.spi; + +import org.junit.jupiter.api.Test; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; + +import java.util.Map; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.assertSame; + +/** + * Verifies that the Spring-managed source registry receives discovered connector beans. + * + *

This test reaches the actual Spring constructor-selection boundary instead of directly + * instantiating {@link CdcSourceRegistry}. It prevents a public no-argument constructor from + * silently bypassing the {@code ObjectProvider} integration path.

+ */ +class CdcSourceRegistrySpringWiringTest { + + @Test + void springContextRegistersDiscoveredSourceConnectorBean() { + TestSourceConnector connector = new TestSourceConnector(); + + try (AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext()) { + context.registerBean(CdcSourceConnector.class, () -> connector); + context.register(CdcSourceRegistry.class); + context.refresh(); + + CdcSourceRegistry registry = context.getBean(CdcSourceRegistry.class); + + assertSame( + connector, + registry.find(connector.id()).orElseThrow(), + "Spring must construct the registry through its connector-provider constructor" + ); + } + } + + private static final class TestSourceConnector implements CdcSourceConnector { + + @Override + public String id() { + return "test_source"; + } + + @Override + public String displayName() { + return "Test source"; + } + + @Override + public SourceCapabilities capabilities() { + return new SourceCapabilities("test", Set.of("test_database"), false); + } + + @Override + public void validate(Map config) { + // No configuration is required for this constructor-selection regression fixture. + } + + @Override + public void start(Map config) { + // No runtime capture is required for this constructor-selection regression fixture. + } + + @Override + public void stop() { + // No runtime capture is started by this constructor-selection regression fixture. + } + } +}