From 70a748a967fd6ad8e10e2d19c9ca7d3ce996da6d Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Tue, 1 Sep 2026 21:47:24 +0900
Subject: [PATCH 01/36] test(cdc): require semantic source identifiers
---
.../CdcSourceFactoryIdentityContractTest.java | 26 +++++++++++++++++++
1 file changed, 26 insertions(+)
diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceFactoryIdentityContractTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceFactoryIdentityContractTest.java
index 6bd004b7..a1082e36 100644
--- a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceFactoryIdentityContractTest.java
+++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/CdcSourceFactoryIdentityContractTest.java
@@ -3,8 +3,11 @@
import org.junit.jupiter.api.Test;
import java.util.List;
+import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -13,6 +16,29 @@
*/
class CdcSourceFactoryIdentityContractTest {
+ @Test
+ void sourceSpecPublishesSemanticIdentifiers() {
+ CdcSourceFactory.SourceSpec sourceSpec =
+ new CdcSourceFactory.SourceSpec("pg-main", "postgres-debezium", true);
+
+ assertEquals("pg-main", sourceSpec.sourceId());
+ assertEquals("postgres-debezium", sourceSpec.sourceType());
+ }
+
+ @Test
+ void configuredSourceDescriptionKeepsLegacyWireKeysAtCompatibilityBoundary() {
+ CdcSourceFactory factory = factory();
+
+ Map sourceDescription = factory.describeConfigured(List.of(
+ new CdcSourceFactory.SourceSpec("pg-main", "postgres-debezium", true)
+ )).getFirst();
+
+ assertEquals("pg-main", sourceDescription.get("id"));
+ assertEquals("postgres-debezium", sourceDescription.get("type"));
+ assertFalse(sourceDescription.containsKey("sourceId"));
+ assertFalse(sourceDescription.containsKey("sourceType"));
+ }
+
@Test
void duplicateConfiguredSourceIdsFailClosed() {
CdcSourceFactory factory = factory();
From f48f800a545fff96a9b9aa7eecae53ae35585110 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Tue, 1 Sep 2026 21:49:03 +0900
Subject: [PATCH 02/36] refactor(cdc): make source spec identifiers explicit
---
.../com/xtrmetl/cdc/spi/CdcSourceFactory.java | 62 ++++++++++---------
1 file changed, 33 insertions(+), 29 deletions(-)
diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceFactory.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceFactory.java
index d9be8af1..e09918e2 100644
--- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceFactory.java
+++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/CdcSourceFactory.java
@@ -19,62 +19,66 @@
@Component
public class CdcSourceFactory {
- private final CdcSourceRegistry registry;
+ private final CdcSourceRegistry sourceRegistry;
- public CdcSourceFactory(CdcSourceRegistry registry) {
- this.registry = registry;
+ public CdcSourceFactory(CdcSourceRegistry sourceRegistry) {
+ this.sourceRegistry = sourceRegistry;
}
- public Optional resolve(String type) {
- if (type == null || type.isBlank()) {
+ public Optional resolve(String sourceType) {
+ if (sourceType == null || sourceType.isBlank()) {
return Optional.empty();
}
- String normalized = type.trim().toLowerCase(Locale.ROOT);
- return registry.find(normalized);
+ String normalizedSourceType = sourceType.trim().toLowerCase(Locale.ROOT);
+ return sourceRegistry.find(normalizedSourceType);
}
/**
* Validates and describes a multi-source configuration list without starting engines.
* Live capture remains single-source in {@code CdcService}.
*
- * @param specs configured source entries; {@code null} is treated as an empty list
+ * @param sourceSpecs configured source entries; {@code null} is treated as an empty list
* @return one descriptive row for each configured source
* @throws IllegalArgumentException when two entries declare the same source id
*/
- public List
*
* @param connectorId registered connector identifier
- * @param batch normalized change records to write
- * @throws NullPointerException when {@code connectorId} or {@code batch} is null
+ * @param changeBatch normalized change records to write
+ * @throws NullPointerException when {@code connectorId} or {@code changeBatch} is null
* @throws IllegalArgumentException when the connector identifier is unknown
* @throws IllegalStateException when the connector is disabled or the dispatcher is closed
* @throws UnsupportedOperationException when the connector is not production-supported
*/
- public void dispatch(String connectorId, List batch) {
+ public void dispatch(String connectorId, List changeBatch) {
Objects.requireNonNull(connectorId, "connectorId must not be null");
- Objects.requireNonNull(batch, "batch must not be null");
+ Objects.requireNonNull(changeBatch, "batch must not be null");
Lock dispatchLock = lifecycleGate.readLock();
dispatchLock.lock();
try {
- if (closed) {
+ if (dispatcherClosed) {
throw new IllegalStateException("Target connector dispatcher is closed");
}
- TargetConnector connector = registry.find(connectorId)
+ TargetConnector targetConnector = targetRegistry.find(connectorId)
.orElseThrow(() -> new IllegalArgumentException("Unknown connector: " + connectorId));
- if (!properties.isEnabled(connectorId)) {
+ if (!connectorProperties.isEnabled(connectorId)) {
throw new IllegalStateException(
"Connector '" + connectorId + "' is disabled. Set xtrmetl.connectors."
+ connectorId.replace("-", ".")
@@ -105,19 +108,19 @@ public void dispatch(String connectorId, List batch) {
);
}
- Map config = properties.configMap(connectorId);
- if (connector.status() != ConnectorStatus.SUPPORTED) {
- connector.validate(config);
- throw new UnsupportedOperationException(connector.writeRefusalReason());
+ Map targetConfig = connectorProperties.configMap(connectorId);
+ if (targetConnector.status() != ConnectorStatus.SUPPORTED) {
+ targetConnector.validate(targetConfig);
+ throw new UnsupportedOperationException(targetConnector.writeRefusalReason());
}
- Object connectorLock = lifecycleLocks.computeIfAbsent(
+ Object targetLock = lifecycleLocks.computeIfAbsent(
connectorId,
ignored -> new Object()
);
- synchronized (connectorLock) {
- ensureOpen(connectorId, connector, config);
- connector.write(batch);
+ synchronized (targetLock) {
+ ensureOpen(connectorId, targetConnector, targetConfig);
+ targetConnector.write(changeBatch);
}
} finally {
dispatchLock.unlock();
@@ -133,24 +136,24 @@ public List> catalog() {
Lock catalogLock = lifecycleGate.readLock();
catalogLock.lock();
try {
- List> rows = new ArrayList<>();
- for (TargetConnector connector : registry.all()) {
- Map row = new LinkedHashMap<>();
- row.put("id", connector.id());
- row.put("displayName", connector.displayName());
- row.put("status", connector.status().name());
- row.put("enabled", properties.isEnabled(connector.id()));
- row.put("writable", !closed
- && connector.status() == ConnectorStatus.SUPPORTED
- && properties.isEnabled(connector.id()));
- row.put("opened", openedConnectors.get(connector.id()) == connector);
- row.put("requiredConfigKeys", connector.requiredConfigKeys());
- row.put("optionalConfigKeys", connector.optionalConfigKeys());
- row.put("writeRefusalReason", connector.writeRefusalReason());
- row.put("integration", connector.describeIntegration());
- rows.add(row);
+ List> catalogRows = new ArrayList<>();
+ for (TargetConnector targetConnector : targetRegistry.all()) {
+ Map catalogRow = new LinkedHashMap<>();
+ catalogRow.put("id", targetConnector.targetId());
+ catalogRow.put("displayName", targetConnector.displayName());
+ catalogRow.put("status", targetConnector.status().name());
+ catalogRow.put("enabled", connectorProperties.isEnabled(targetConnector.targetId()));
+ catalogRow.put("writable", !dispatcherClosed
+ && targetConnector.status() == ConnectorStatus.SUPPORTED
+ && connectorProperties.isEnabled(targetConnector.targetId()));
+ catalogRow.put("opened", openedConnectors.get(targetConnector.targetId()) == targetConnector);
+ catalogRow.put("requiredConfigKeys", targetConnector.requiredConfigKeys());
+ catalogRow.put("optionalConfigKeys", targetConnector.optionalConfigKeys());
+ catalogRow.put("writeRefusalReason", targetConnector.writeRefusalReason());
+ catalogRow.put("integration", targetConnector.describeIntegration());
+ catalogRows.add(catalogRow);
}
- return rows;
+ return catalogRows;
} finally {
catalogLock.unlock();
}
@@ -165,25 +168,25 @@ public List> catalog() {
*/
private void ensureOpen(
String connectorId,
- TargetConnector connector,
- Map config
+ TargetConnector targetConnector,
+ Map targetConfig
) {
- TargetConnector active = openedConnectors.get(connectorId);
- if (active == connector) {
+ TargetConnector activeConnector = openedConnectors.get(connectorId);
+ if (activeConnector == targetConnector) {
return;
}
- if (active != null) {
+ if (activeConnector != null) {
throw new IllegalStateException(
"Connector registry entry changed after open: " + connectorId
);
}
- connector.validate(config);
+ targetConnector.validate(targetConfig);
try {
- connector.open(config);
+ targetConnector.open(targetConfig);
} catch (RuntimeException openFailure) {
try {
- connector.close();
+ targetConnector.close();
} catch (RuntimeException cleanupFailure) {
openFailure.addSuppressed(cleanupFailure);
log.warn("Failed to clean up target connector after open failure id={}", connectorId);
@@ -191,7 +194,7 @@ private void ensureOpen(
throw openFailure;
}
- openedConnectors.put(connectorId, connector);
+ openedConnectors.put(connectorId, targetConnector);
log.info("Opened target connector id={}", connectorId);
}
@@ -207,27 +210,27 @@ void closeOpenedConnectors() {
Lock shutdownLock = lifecycleGate.writeLock();
shutdownLock.lock();
try {
- if (closed) {
+ if (dispatcherClosed) {
return;
}
- closed = true;
+ dispatcherClosed = true;
- for (Map.Entry entry
+ for (Map.Entry connectorEntry
: new ArrayList<>(openedConnectors.entrySet())) {
- String connectorId = entry.getKey();
- TargetConnector connector = entry.getValue();
- Object connectorLock = lifecycleLocks.computeIfAbsent(
+ String connectorId = connectorEntry.getKey();
+ TargetConnector targetConnector = connectorEntry.getValue();
+ Object targetLock = lifecycleLocks.computeIfAbsent(
connectorId,
ignored -> new Object()
);
- synchronized (connectorLock) {
- if (!openedConnectors.remove(connectorId, connector)) {
+ synchronized (targetLock) {
+ if (!openedConnectors.remove(connectorId, targetConnector)) {
continue;
}
try {
- connector.close();
+ targetConnector.close();
log.info("Closed target connector id={}", connectorId);
- } catch (RuntimeException exception) {
+ } catch (RuntimeException closeFailure) {
log.error("Failed to close target connector id={}", connectorId);
}
}
From cd4bade91e8d0d32cf7c73903893b0968e904912 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Wed, 2 Sep 2026 03:10:07 +0900
Subject: [PATCH 35/36] test(etl): pin legacy target identity adapter
---
.../etl/connector/TargetConnectorSemanticIdentityTest.java | 4 +++-
1 file changed, 3 insertions(+), 1 deletion(-)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/connector/TargetConnectorSemanticIdentityTest.java b/etl-service/src/test/java/com/xtrmetl/etl/connector/TargetConnectorSemanticIdentityTest.java
index f225ece1..334e16c0 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/connector/TargetConnectorSemanticIdentityTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/connector/TargetConnectorSemanticIdentityTest.java
@@ -10,9 +10,11 @@
class TargetConnectorSemanticIdentityTest {
@Test
- void targetConnectorExposesTargetIdentityBySemanticName() {
+ @SuppressWarnings("deprecation")
+ void targetConnectorExposesSemanticIdentityAndPreservesLegacyAlias() {
TargetConnector databricksTargetConnector = new DatabricksTargetConnector();
assertEquals("databricks", databricksTargetConnector.targetId());
+ assertEquals(databricksTargetConnector.targetId(), databricksTargetConnector.id());
}
}
From 5daa6c9b5bb1c4ab23c96ce0c5af7073f6b804c7 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Mon, 7 Sep 2026 12:09:13 +0900
Subject: [PATCH 36/36] fix(test): supply CDC binding failure exception
---
.../cdc/config/XtrmetlPropertiesSecurityDefaultTest.java | 2 +-
1 file changed, 1 insertion(+), 1 deletion(-)
diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/config/XtrmetlPropertiesSecurityDefaultTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/config/XtrmetlPropertiesSecurityDefaultTest.java
index 39473019..27d6cf55 100644
--- a/cdc-service/src/test/java/com/xtrmetl/cdc/config/XtrmetlPropertiesSecurityDefaultTest.java
+++ b/cdc-service/src/test/java/com/xtrmetl/cdc/config/XtrmetlPropertiesSecurityDefaultTest.java
@@ -51,7 +51,7 @@ void legacySourceKeysBindIntoSemanticJavaIdentifiers() {
XtrmetlProperties xtrmetlProperties = new Binder(configurationSource)
.bind("xtrmetl", Bindable.of(XtrmetlProperties.class))
- .orElseThrow();
+ .orElseThrow(() -> new IllegalStateException("xtrmetl properties binding failed"));
XtrmetlProperties.Source sourceConfiguration =
xtrmetlProperties.getCdc().getSources().getFirst();