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.
+ }
+ }
+}