From dba14059d270176c30093c450388b5d45cf1f395 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Mon, 10 Aug 2026 19:07:33 +0900 Subject: [PATCH 1/3] test(connectors): expose nested ChangeRecord aliasing --- .../etl/connector/ChangeRecordTest.java | 62 +++++++++++++++++++ 1 file changed, 62 insertions(+) diff --git a/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java b/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java index 83e814db..e2afcd45 100644 --- a/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java +++ b/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java @@ -2,7 +2,9 @@ import org.junit.jupiter.api.Test; +import java.util.ArrayList; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -41,6 +43,66 @@ void snapshotsMutableInputMaps() { assertEquals(7L, record.getPk().get("id")); } + @Test + void recursivelySnapshotsNestedJsonContainers() { + List lineItems = new ArrayList<>(); + lineItems.add("first"); + lineItems.add(null); + Map details = new LinkedHashMap<>(); + details.put("line_items", lineItems); + Map after = new LinkedHashMap<>(); + after.put("details", details); + + ChangeRecord record = new ChangeRecord( + "source", + "u", + "public", + "orders", + 123L, + Map.of(), + after, + Map.of("id", 7L) + ); + + lineItems.set(0, "mutated"); + lineItems.add("late"); + details.put("late_field", "mutated"); + + Map snapshottedDetails = (Map) record.getAfter().get("details"); + List snapshottedLineItems = (List) snapshottedDetails.get("line_items"); + assertEquals(2, snapshottedLineItems.size()); + assertEquals("first", snapshottedLineItems.get(0)); + assertNull(snapshottedLineItems.get(1)); + assertTrue(!snapshottedDetails.containsKey("late_field")); + } + + @Test + @SuppressWarnings({"rawtypes", "unchecked"}) + void exposesNestedJsonContainersAsUnmodifiable() { + List lineItems = new ArrayList<>(); + lineItems.add("first"); + Map details = new LinkedHashMap<>(); + details.put("line_items", lineItems); + + ChangeRecord record = new ChangeRecord( + "source", + "c", + "public", + "orders", + 123L, + Map.of(), + Map.of("details", details), + Map.of("id", 7L) + ); + + Map nestedMap = (Map) record.getAfter().get("details"); + List nestedList = (List) nestedMap.get("line_items"); + assertThrows(UnsupportedOperationException.class, + () -> nestedMap.put("late_field", "mutated")); + assertThrows(UnsupportedOperationException.class, + () -> nestedList.add("mutated")); + } + @Test void exposesUnmodifiableMaps() { ChangeRecord record = new ChangeRecord( From 401bc6833887b2244a3f46438467e7d3dc0b67a3 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Mon, 10 Aug 2026 19:14:49 +0900 Subject: [PATCH 2/3] fix(connectors): snapshot nested ChangeRecord containers --- .../xtrmetl/etl/connector/ChangeRecord.java | 29 +++++++++++++++++-- 1 file changed, 26 insertions(+), 3 deletions(-) diff --git a/etl-service/src/main/java/com/xtrmetl/etl/connector/ChangeRecord.java b/etl-service/src/main/java/com/xtrmetl/etl/connector/ChangeRecord.java index d62bc6c0..c6a51934 100644 --- a/etl-service/src/main/java/com/xtrmetl/etl/connector/ChangeRecord.java +++ b/etl-service/src/main/java/com/xtrmetl/etl/connector/ChangeRecord.java @@ -1,15 +1,20 @@ package com.xtrmetl.etl.connector; +import java.util.ArrayList; import java.util.Collections; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import java.util.Objects; /** * Normalized change event for target connectors (canonical CDC/ETL record). * - *

Row maps are shallow-snapshotted at construction time so callers cannot mutate a - * record after it has entered a connector pipeline. Null database values remain supported.

+ *

JSON-shaped row maps and nested map/list containers are recursively snapshotted at + * construction time so later caller mutations cannot change a record that has entered a + * connector pipeline. Null database values remain supported. Scalar and other non-container + * values are retained as supplied; this type does not claim to clone arbitrary mutable Java + * objects.

*/ public final class ChangeRecord { @@ -101,6 +106,24 @@ private static Map snapshot(Map source) { if (source == null || source.isEmpty()) { return Map.of(); } - return Collections.unmodifiableMap(new LinkedHashMap<>(source)); + Map copy = new LinkedHashMap<>(); + source.forEach((key, value) -> copy.put(key, snapshotValue(value))); + return Collections.unmodifiableMap(copy); + } + + private static Object snapshotValue(Object value) { + if (value instanceof Map mapValue) { + Map copy = new LinkedHashMap<>(); + mapValue.forEach((key, nestedValue) -> copy.put(key, snapshotValue(nestedValue))); + return Collections.unmodifiableMap(copy); + } + if (value instanceof List listValue) { + List copy = new ArrayList<>(listValue.size()); + for (Object element : listValue) { + copy.add(snapshotValue(element)); + } + return Collections.unmodifiableList(copy); + } + return value; } } From c4dd764ded629aed111131c518acef8a2499d67d Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Mon, 10 Aug 2026 19:23:20 +0900 Subject: [PATCH 3/3] test(connectors): reject cyclic ChangeRecord containers --- .../etl/connector/ChangeRecordTest.java | 22 +++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java b/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java index e2afcd45..f8750fbd 100644 --- a/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java +++ b/etl-service/src/test/java/com/xtrmetl/etl/connector/ChangeRecordTest.java @@ -103,6 +103,28 @@ void exposesNestedJsonContainersAsUnmodifiable() { () -> nestedList.add("mutated")); } + @Test + void rejectsCyclicNestedContainersDeterministically() { + Map cyclic = new LinkedHashMap<>(); + cyclic.put("self", cyclic); + + IllegalArgumentException failure = assertThrows( + IllegalArgumentException.class, + () -> new ChangeRecord( + "source", + "c", + "public", + "orders", + 123L, + Map.of(), + Map.of("cycle", cyclic), + Map.of("id", 7L) + ) + ); + + assertEquals("ChangeRecord JSON containers must not contain cycles", failure.getMessage()); + } + @Test void exposesUnmodifiableMaps() { ChangeRecord record = new ChangeRecord(