From ddb45d7d1889ae3ea844b750ca49011381eb081c Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:10:37 +0900
Subject: [PATCH 01/32] test(etl): capture replay result contract on repaired
stack
---
.../etl/job/EtlJobReplayBoundaryTest.java | 39 +++++++++++++++++++
1 file changed, 39 insertions(+)
create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
new file mode 100644
index 00000000..0a13e47d
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -0,0 +1,39 @@
+package com.xtrmetl.etl.job;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.UUID;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * Captures the first immutable replay-result contract on the repaired durable-job stack.
+ */
+class EtlJobReplayBoundaryTest {
+
+ private static final UUID SOURCE_ID = UUID.fromString(
+ "cf4f083f-8c90-4f34-a8b6-b53761de44ef"
+ );
+
+ @Test
+ void replayResultRequiresIdentityAndStatusButPreservesCurrentState() {
+ EtlJobReplay pending = new EtlJobReplay(SOURCE_ID, EtlJobStatus.PENDING, false);
+ EtlJobReplay terminalReplay = new EtlJobReplay(
+ SOURCE_ID,
+ EtlJobStatus.SUCCEEDED,
+ true
+ );
+
+ assertEquals(EtlJobStatus.PENDING, pending.jobStatus());
+ assertEquals(EtlJobStatus.SUCCEEDED, terminalReplay.jobStatus());
+ assertThrows(
+ NullPointerException.class,
+ () -> new EtlJobReplay(null, EtlJobStatus.PENDING, false)
+ );
+ assertThrows(
+ NullPointerException.class,
+ () -> new EtlJobReplay(SOURCE_ID, null, false)
+ );
+ }
+}
From 167a7778f240e15777c7925cbeaca43015572221 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:13:13 +0900
Subject: [PATCH 02/32] feat(etl): add immutable replay result
---
.../com/xtrmetl/etl/job/EtlJobReplay.java | 29 +++++++++++++++++++
1 file changed, 29 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplay.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplay.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplay.java
new file mode 100644
index 00000000..65caf105
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplay.java
@@ -0,0 +1,29 @@
+package com.xtrmetl.etl.job;
+
+import java.util.Objects;
+import java.util.UUID;
+
+/**
+ * Reports one newly accepted or replayed immutable-lineage durable ETL job.
+ *
+ * The source terminal resource and lineage remain internal persistence evidence. This result
+ * exposes only the new opaque job identifier, its current stable lifecycle state, and whether the
+ * same principal-scoped replay request had already created it. A first creation is pending; a later
+ * idempotent retry may correctly report that the same created job has since progressed.
+ *
+ * @param jobRecordId replay-created durable job identifier
+ * @param jobStatus current stable lifecycle state of that created job
+ * @param replayed {@code true} when this response reuses an already-created replay job
+ */
+public record EtlJobReplay(
+ UUID jobRecordId,
+ EtlJobStatus jobStatus,
+ boolean replayed
+) {
+
+ /** Validates the immutable replay result. */
+ public EtlJobReplay {
+ Objects.requireNonNull(jobRecordId, "jobRecordId must not be null");
+ Objects.requireNonNull(jobStatus, "jobStatus must not be null");
+ }
+}
From 1bb07e4e0c959da42dce5dca2031a173f7b04f67 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:15:26 +0900
Subject: [PATCH 03/32] test(etl): capture replay service construction contract
---
.../etl/job/EtlJobReplayBoundaryTest.java | 31 ++++++++++++++++++-
1 file changed, 30 insertions(+), 1 deletion(-)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index 0a13e47d..60b67ed7 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -1,14 +1,18 @@
package com.xtrmetl.etl.job;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.xtrmetl.etl.service.EtlBatchProperties;
import org.junit.jupiter.api.Test;
+import org.springframework.jdbc.core.JdbcTemplate;
import java.util.UUID;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
/**
- * Captures the first immutable replay-result contract on the repaired durable-job stack.
+ * Captures immutable replay-result and collaborator-construction contracts on the repaired stack.
*/
class EtlJobReplayBoundaryTest {
@@ -36,4 +40,29 @@ void replayResultRequiresIdentityAndStatusButPreservesCurrentState() {
() -> new EtlJobReplay(SOURCE_ID, null, false)
);
}
+
+ @Test
+ void constructorsRejectMissingCollaborators() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ ObjectMapper mapper = new ObjectMapper();
+ EtlBatchProperties properties = new EtlBatchProperties();
+
+ assertThrows(
+ NullPointerException.class,
+ () -> new EtlJobReplayService(null, mapper, properties, hash -> true)
+ );
+ assertThrows(
+ NullPointerException.class,
+ () -> new EtlJobReplayService(jdbcTemplate, null, properties, hash -> true)
+ );
+ assertThrows(
+ NullPointerException.class,
+ () -> new EtlJobReplayService(jdbcTemplate, mapper, null, hash -> true)
+ );
+ assertThrows(
+ NullPointerException.class,
+ () -> new EtlJobReplayService(jdbcTemplate, mapper, properties, null)
+ );
+ new EtlJobReplayService(jdbcTemplate, mapper, properties);
+ }
}
From a0e43c603f03e84247f8c72c4607783f6a2a5a46 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:16:48 +0900
Subject: [PATCH 04/32] feat(etl): add replay service construction boundary
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 80 +++++++++++++++++++
1 file changed, 80 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
new file mode 100644
index 00000000..d0a79c40
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -0,0 +1,80 @@
+package com.xtrmetl.etl.job;
+
+import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.xtrmetl.etl.service.EtlBatchProperties;
+import com.xtrmetl.etl.service.EtlRequestLock;
+import com.xtrmetl.etl.service.PostgresEtlRequestLock;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.stereotype.Service;
+
+import java.util.Objects;
+
+/**
+ * Admits owner-scoped durable-job replay requests on the repaired stack.
+ *
+ * This initial implementation establishes only the validated collaborator boundary required by
+ * the first replay-service contract. Replay admission behavior is added in subsequent test-first
+ * increments so an incomplete branch cannot silently claim the full replay contract.
+ */
+@Service
+public class EtlJobReplayService {
+
+ private final JdbcTemplate jdbcTemplate;
+ private final ObjectMapper objectMapper;
+ private final EtlBatchProperties batchProperties;
+ private final EtlRequestLock requestLock;
+
+ /**
+ * Creates replay admission with the PostgreSQL transaction-lock implementation.
+ *
+ * @param jdbcTemplate parameterized durable job persistence
+ * @param objectMapper JSON parser configuration to copy
+ * @param batchProperties bounded request limits
+ */
+ public EtlJobReplayService(
+ JdbcTemplate jdbcTemplate,
+ ObjectMapper objectMapper,
+ EtlBatchProperties batchProperties
+ ) {
+ this(
+ jdbcTemplate,
+ objectMapper,
+ batchProperties,
+ new PostgresEtlRequestLock(Objects.requireNonNull(
+ jdbcTemplate,
+ "jdbcTemplate must not be null"
+ ))
+ );
+ }
+
+ /**
+ * Creates replay admission with an explicit transaction-lifetime request lock.
+ *
+ * @param jdbcTemplate parameterized durable job persistence
+ * @param objectMapper JSON parser configuration to copy
+ * @param batchProperties bounded request limits
+ * @param requestLock transaction-lifetime replay-key lock
+ */
+ @Autowired
+ public EtlJobReplayService(
+ JdbcTemplate jdbcTemplate,
+ ObjectMapper objectMapper,
+ EtlBatchProperties batchProperties,
+ EtlRequestLock requestLock
+ ) {
+ this.jdbcTemplate = Objects.requireNonNull(jdbcTemplate, "jdbcTemplate must not be null");
+ ObjectMapper sourceMapper = Objects.requireNonNull(
+ objectMapper,
+ "objectMapper must not be null"
+ );
+ this.objectMapper = sourceMapper.copy();
+ this.objectMapper.enable(JsonParser.Feature.STRICT_DUPLICATE_DETECTION);
+ this.batchProperties = Objects.requireNonNull(
+ batchProperties,
+ "batchProperties must not be null"
+ );
+ this.requestLock = Objects.requireNonNull(requestLock, "requestLock must not be null");
+ }
+}
From 3b28af9097aebf013901534d6f5ccbaa5fd6c1de Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:18:51 +0900
Subject: [PATCH 05/32] test(etl): capture replay validation before persistence
---
.../etl/job/EtlJobReplayBoundaryTest.java | 68 ++++++++++++++++++-
1 file changed, 67 insertions(+), 1 deletion(-)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index 60b67ed7..1efb1f84 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -2,6 +2,7 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import com.xtrmetl.etl.service.EtlBatchProperties;
+import com.xtrmetl.etl.service.EtlRequestException;
import org.junit.jupiter.api.Test;
import org.springframework.jdbc.core.JdbcTemplate;
@@ -10,15 +11,18 @@
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.verifyNoInteractions;
/**
- * Captures immutable replay-result and collaborator-construction contracts on the repaired stack.
+ * Captures immutable replay-result, construction, and pre-persistence validation contracts.
*/
class EtlJobReplayBoundaryTest {
private static final UUID SOURCE_ID = UUID.fromString(
"cf4f083f-8c90-4f34-a8b6-b53761de44ef"
);
+ private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
+ private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
@Test
void replayResultRequiresIdentityAndStatusButPreservesCurrentState() {
@@ -65,4 +69,66 @@ void constructorsRejectMissingCollaborators() {
);
new EtlJobReplayService(jdbcTemplate, mapper, properties);
}
+
+ @Test
+ void validatesIdentityKeyPrincipalAndPayloadBeforeDatabaseWork() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ new ObjectMapper(),
+ new EtlBatchProperties(),
+ hash -> true
+ );
+
+ assertThrows(
+ NullPointerException.class,
+ () -> service.replayOwned(null, PAYLOAD, REPLAY_KEY, "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_job_replay_key_required",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, null, "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_job_replay_key_required",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, "unsafe key", "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_idempotency_principal_required",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, null)
+ );
+ assertErrorCode(
+ "etl_idempotency_principal_required",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, " ".repeat(513))
+ );
+ assertErrorCode(
+ "etl_idempotency_principal_required",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, "a".repeat(513))
+ );
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, null, REPLAY_KEY, "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, "", REPLAY_KEY, "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, "null", REPLAY_KEY, "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, "not-json", REPLAY_KEY, "tenant_alpha")
+ );
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, "{}", REPLAY_KEY, "tenant_alpha")
+ );
+ verifyNoInteractions(jdbcTemplate);
+ }
+
+ private static void assertErrorCode(String expectedCode, Runnable invocation) {
+ EtlRequestException exception = assertThrows(EtlRequestException.class, invocation::run);
+ assertEquals(expectedCode, exception.getMessage());
+ }
}
From a8591a7aaa9155b0ec80f9460da91d77d5b12d44 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:21:41 +0900
Subject: [PATCH 06/32] feat(etl): classify invalid replay keys
---
.../java/com/xtrmetl/etl/service/EtlRequestError.java | 9 +++++++++
1 file changed, 9 insertions(+)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
index d28aa33d..e62e89f5 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
@@ -131,6 +131,15 @@ public enum EtlRequestError {
"Cancellation requires a supported principal-scoped Idempotency-Key."
),
+ /** The replay key is absent or outside the bounded safe idempotency profile. */
+ JOB_REPLAY_KEY_REQUIRED(
+ HttpStatus.BAD_REQUEST,
+ "etl_job_replay_key_required",
+ "urn:mightyetl:problem:etl-job-replay-key-required",
+ "ETL job replay key required",
+ "Replay requires a supported principal-scoped Idempotency-Key."
+ ),
+
/** The durable job was already cancelled with a different cancellation key. */
JOB_CANCELLATION_KEY_REUSED(
HttpStatus.UNPROCESSABLE_ENTITY,
From cad4b3c19c2d2c82e498e5f4f3b01b865acfe9bd Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:22:43 +0900
Subject: [PATCH 07/32] feat(etl): validate replay requests before persistence
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 90 ++++++++++++++++++-
1 file changed, 87 insertions(+), 3 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index d0a79c40..5e8f40ac 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -1,26 +1,42 @@
package com.xtrmetl.etl.job;
import com.fasterxml.jackson.core.JsonParser;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.xtrmetl.etl.service.EtlBatchProperties;
+import com.xtrmetl.etl.service.EtlRequestError;
+import com.xtrmetl.etl.service.EtlRequestException;
import com.xtrmetl.etl.service.EtlRequestLock;
import com.xtrmetl.etl.service.PostgresEtlRequestLock;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.lang.Nullable;
import org.springframework.stereotype.Service;
import java.util.Objects;
+import java.util.UUID;
+import java.util.regex.Pattern;
/**
* Admits owner-scoped durable-job replay requests on the repaired stack.
*
- * This initial implementation establishes only the validated collaborator boundary required by
- * the first replay-service contract. Replay admission behavior is added in subsequent test-first
- * increments so an incomplete branch cannot silently claim the full replay contract.
+ * The service validates replay identity, authenticated principal scope, and JSON shape before
+ * persistence work. Database admission and replay transitions are intentionally added in later
+ * test-first increments so an incomplete branch cannot silently claim the full replay contract.
*/
@Service
public class EtlJobReplayService {
+ private static final int MAX_PRINCIPAL_SCOPE_CODE_POINTS = 512;
+ private static final String IDEMPOTENCY_KEY_VALUE_EXPRESSION = "[A-Za-z0-9._:-]{16,128}";
+ private static final Pattern IDEMPOTENCY_KEY_VALUE_PROFILE = Pattern.compile(
+ IDEMPOTENCY_KEY_VALUE_EXPRESSION
+ );
+ private static final Pattern IDEMPOTENCY_KEY_STRUCTURED_FIELD_PROFILE = Pattern.compile(
+ "\"(" + IDEMPOTENCY_KEY_VALUE_EXPRESSION + ")\""
+ );
+
private final JdbcTemplate jdbcTemplate;
private final ObjectMapper objectMapper;
private final EtlBatchProperties batchProperties;
@@ -77,4 +93,72 @@ public EtlJobReplayService(
);
this.requestLock = Objects.requireNonNull(requestLock, "requestLock must not be null");
}
+
+ /**
+ * Validates one owner-scoped replay request before any persistence work.
+ *
+ * The database admission and replay transition are introduced in later TDD increments.
+ * Until then, an otherwise valid request fails closed rather than acquiring a lock or touching
+ * durable state.
+ *
+ * @param sourceJobRecordId terminal durable job whose intent is being replayed
+ * @param requestBody bounded JSON-array payload resupplied by the authenticated owner
+ * @param rawReplayKey principal-scoped idempotency key for this replay request
+ * @param principalName authenticated principal namespace
+ * @return replay-created durable job result after persistence support is introduced
+ * @throws NullPointerException when the source job identifier is absent
+ * @throws EtlRequestException when key, principal, or JSON validation fails
+ * @throws IllegalStateException while the persistence increment is intentionally unavailable
+ */
+ public EtlJobReplay replayOwned(
+ UUID sourceJobRecordId,
+ @Nullable String requestBody,
+ @Nullable String rawReplayKey,
+ @Nullable String principalName
+ ) {
+ Objects.requireNonNull(sourceJobRecordId, "sourceJobRecordId must not be null");
+ validateReplayKey(rawReplayKey);
+ validatePrincipalScope(principalName);
+ validateReplayPayload(requestBody);
+ throw new IllegalStateException("Durable ETL job replay persistence is not yet available");
+ }
+
+ private static String validateReplayKey(@Nullable String rawReplayKey) {
+ if (rawReplayKey == null) {
+ throw new EtlRequestException(EtlRequestError.JOB_REPLAY_KEY_REQUIRED);
+ }
+ var structuredFieldMatcher = IDEMPOTENCY_KEY_STRUCTURED_FIELD_PROFILE.matcher(rawReplayKey);
+ if (structuredFieldMatcher.matches()) {
+ return structuredFieldMatcher.group(1);
+ }
+ if (IDEMPOTENCY_KEY_VALUE_PROFILE.matcher(rawReplayKey).matches()) {
+ return rawReplayKey;
+ }
+ throw new EtlRequestException(EtlRequestError.JOB_REPLAY_KEY_REQUIRED);
+ }
+
+ private static String validatePrincipalScope(@Nullable String principalName) {
+ if (principalName == null
+ || principalName.isBlank()
+ || principalName.codePointCount(0, principalName.length())
+ > MAX_PRINCIPAL_SCOPE_CODE_POINTS) {
+ throw new EtlRequestException(EtlRequestError.IDEMPOTENCY_PRINCIPAL_REQUIRED);
+ }
+ return principalName;
+ }
+
+ private JsonNode validateReplayPayload(@Nullable String requestBody) {
+ if (requestBody == null || requestBody.isEmpty()) {
+ throw new EtlRequestException(EtlRequestError.INVALID_JSON);
+ }
+ try {
+ JsonNode root = objectMapper.readTree(requestBody);
+ if (root == null || !root.isArray()) {
+ throw new EtlRequestException(EtlRequestError.INVALID_JSON);
+ }
+ return root;
+ } catch (JsonProcessingException exception) {
+ throw new EtlRequestException(EtlRequestError.INVALID_JSON, exception);
+ }
+ }
}
From 09a66542bdb3c679b2d4e0e033caf419572e34b3 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:24:36 +0900
Subject: [PATCH 08/32] test(etl): close replay validation coverage gaps
---
.../etl/job/EtlJobReplayBoundaryTest.java | 35 +++++++++++++++++++
1 file changed, 35 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index 1efb1f84..cc5d4e3b 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -23,6 +23,8 @@ class EtlJobReplayBoundaryTest {
);
private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
+ private static final String PERSISTENCE_NOT_AVAILABLE =
+ "Durable ETL job replay persistence is not yet available";
@Test
void replayResultRequiresIdentityAndStatusButPreservesCurrentState() {
@@ -112,6 +114,10 @@ void validatesIdentityKeyPrincipalAndPayloadBeforeDatabaseWork() {
"etl_invalid_json",
() -> service.replayOwned(SOURCE_ID, "", REPLAY_KEY, "tenant_alpha")
);
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, " ", REPLAY_KEY, "tenant_alpha")
+ );
assertErrorCode(
"etl_invalid_json",
() -> service.replayOwned(SOURCE_ID, "null", REPLAY_KEY, "tenant_alpha")
@@ -127,6 +133,35 @@ void validatesIdentityKeyPrincipalAndPayloadBeforeDatabaseWork() {
verifyNoInteractions(jdbcTemplate);
}
+ @Test
+ void acceptsRawAndStructuredReplayKeysThenFailsClosedBeforePersistence() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ new ObjectMapper(),
+ new EtlBatchProperties(),
+ hash -> true
+ );
+
+ IllegalStateException rawKeyFailure = assertThrows(
+ IllegalStateException.class,
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, "tenant_alpha")
+ );
+ assertEquals(PERSISTENCE_NOT_AVAILABLE, rawKeyFailure.getMessage());
+
+ IllegalStateException structuredKeyFailure = assertThrows(
+ IllegalStateException.class,
+ () -> service.replayOwned(
+ SOURCE_ID,
+ PAYLOAD,
+ "\"" + REPLAY_KEY + "\"",
+ "tenant_alpha"
+ )
+ );
+ assertEquals(PERSISTENCE_NOT_AVAILABLE, structuredKeyFailure.getMessage());
+ verifyNoInteractions(jdbcTemplate);
+ }
+
private static void assertErrorCode(String expectedCode, Runnable invocation) {
EtlRequestException exception = assertThrows(EtlRequestException.class, invocation::run);
assertEquals(expectedCode, exception.getMessage());
From bd821f2bba181959406a6fbd338e00ddc92555ab Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:27:03 +0900
Subject: [PATCH 09/32] test(etl): cover absent replay parse root
---
.../etl/job/EtlJobReplayBoundaryTest.java | 29 +++++++++++++++++++
1 file changed, 29 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index cc5d4e3b..058d2326 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -1,5 +1,6 @@
package com.xtrmetl.etl.job;
+import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.xtrmetl.etl.service.EtlBatchProperties;
import com.xtrmetl.etl.service.EtlRequestException;
@@ -133,6 +134,34 @@ void validatesIdentityKeyPrincipalAndPayloadBeforeDatabaseWork() {
verifyNoInteractions(jdbcTemplate);
}
+ @Test
+ void rejectsAnAbsentParsedRootBeforeDatabaseWork() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ ObjectMapper absentRootMapper = new ObjectMapper() {
+ @Override
+ public ObjectMapper copy() {
+ return this;
+ }
+
+ @Override
+ public JsonNode readTree(String content) {
+ return null;
+ }
+ };
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ absentRootMapper,
+ new EtlBatchProperties(),
+ hash -> true
+ );
+
+ assertErrorCode(
+ "etl_invalid_json",
+ () -> service.replayOwned(SOURCE_ID, "[]", REPLAY_KEY, "tenant_alpha")
+ );
+ verifyNoInteractions(jdbcTemplate);
+ }
+
@Test
void acceptsRawAndStructuredReplayKeysThenFailsClosedBeforePersistence() {
JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
From cb24c6d3526be7579c09a424990cfc128d5f022f Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:38:08 +0900
Subject: [PATCH 10/32] test(etl): prove first replay persistence boundary
---
...tlJobReplayPersistenceIntegrationTest.java | 210 ++++++++++++++++++
1 file changed, 210 insertions(+)
create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
new file mode 100644
index 00000000..fe37ed89
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -0,0 +1,210 @@
+package com.xtrmetl.etl.job;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.xtrmetl.etl.service.EtlBatchProperties;
+import com.xtrmetl.etl.service.EtlRequestLock;
+import com.xtrmetl.etl.service.Sha256Digest;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.jdbc.datasource.DataSourceTransactionManager;
+import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseBuilder;
+import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType;
+import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.annotation.EnableTransactionManagement;
+
+import javax.sql.DataSource;
+import java.time.Instant;
+import java.util.UUID;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+/**
+ * Proves the first database-owned durable-job replay transition on the repaired stack.
+ */
+@SpringJUnitConfig(EtlJobReplayPersistenceIntegrationTest.TestConfiguration.class)
+class EtlJobReplayPersistenceIntegrationTest {
+
+ private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
+ private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
+ private static final String PRINCIPAL = "tenant_alpha";
+
+ private final EtlJobReplayService replayService;
+ private final JdbcTemplate jdbcTemplate;
+
+ @Autowired
+ EtlJobReplayPersistenceIntegrationTest(
+ EtlJobReplayService replayService,
+ JdbcTemplate jdbcTemplate
+ ) {
+ this.replayService = replayService;
+ this.jdbcTemplate = jdbcTemplate;
+ }
+
+ @BeforeEach
+ void createDurableJobTable() {
+ jdbcTemplate.execute("DROP TABLE IF EXISTS etl_job_records");
+ jdbcTemplate.execute("""
+ CREATE TABLE etl_job_records (
+ job_record_id UUID PRIMARY KEY,
+ principal_scope_hash CHAR(64) NOT NULL,
+ submission_key_hash CHAR(64) NOT NULL,
+ request_digest CHAR(64) NOT NULL,
+ request_payload CLOB,
+ job_status VARCHAR(32) NOT NULL,
+ attempt_count INTEGER NOT NULL DEFAULT 0,
+ failure_code VARCHAR(128),
+ lease_claim_id UUID,
+ lease_owner_id VARCHAR(128),
+ lease_expires_at TIMESTAMP WITH TIME ZONE,
+ cancellation_key_hash CHAR(64),
+ cancellation_code VARCHAR(128),
+ job_cancelled_at TIMESTAMP WITH TIME ZONE,
+ replay_source_job_record_id UUID,
+ replay_root_job_record_id UUID,
+ replay_generation_count INTEGER,
+ created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ updated_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ CONSTRAINT etl_job_submission_scope_unique
+ UNIQUE (principal_scope_hash, submission_key_hash)
+ )
+ """);
+ }
+
+ @Test
+ void createsFirstGenerationPendingReplayWithoutMutatingTerminalSource() {
+ UUID sourceJobRecordId = UUID.randomUUID();
+ jdbcTemplate.update(
+ """
+ INSERT INTO etl_job_records (
+ job_record_id, principal_scope_hash, submission_key_hash,
+ request_digest, request_payload, job_status, attempt_count,
+ failure_code
+ ) VALUES (?, ?, ?, ?, NULL, 'FAILED', 1, ?)
+ """,
+ sourceJobRecordId,
+ Sha256Digest.digest(PRINCIPAL),
+ Sha256Digest.digest("original-submission-key"),
+ Sha256Digest.digest(PAYLOAD),
+ "etl_target_failure"
+ );
+ Instant sourceUpdatedAt = instantColumn(sourceJobRecordId, "updated_at");
+
+ EtlJobReplay replay = replayService.replayOwned(
+ sourceJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ );
+
+ assertFalse(replay.replayed());
+ assertEquals(EtlJobStatus.PENDING, replay.jobStatus());
+ assertNotEquals(sourceJobRecordId, replay.jobRecordId());
+ assertEquals(
+ sourceJobRecordId,
+ uuidColumn(replay.jobRecordId(), "replay_source_job_record_id")
+ );
+ assertEquals(
+ sourceJobRecordId,
+ uuidColumn(replay.jobRecordId(), "replay_root_job_record_id")
+ );
+ assertEquals(1, integerColumn(replay.jobRecordId(), "replay_generation_count"));
+ assertEquals(PAYLOAD, textColumn(replay.jobRecordId(), "request_payload"));
+ assertEquals("FAILED", textColumn(sourceJobRecordId, "job_status"));
+ assertNull(textColumn(sourceJobRecordId, "request_payload"));
+ assertEquals(sourceUpdatedAt, instantColumn(sourceJobRecordId, "updated_at"));
+ }
+
+ private String textColumn(UUID jobRecordId, String columnName) {
+ return jdbcTemplate.queryForObject(
+ "SELECT " + columnName + " FROM etl_job_records WHERE job_record_id=?",
+ String.class,
+ jobRecordId
+ );
+ }
+
+ private UUID uuidColumn(UUID jobRecordId, String columnName) {
+ return jdbcTemplate.queryForObject(
+ "SELECT " + columnName + " FROM etl_job_records WHERE job_record_id=?",
+ UUID.class,
+ jobRecordId
+ );
+ }
+
+ private Integer integerColumn(UUID jobRecordId, String columnName) {
+ return jdbcTemplate.queryForObject(
+ "SELECT " + columnName + " FROM etl_job_records WHERE job_record_id=?",
+ Integer.class,
+ jobRecordId
+ );
+ }
+
+ private Instant instantColumn(UUID jobRecordId, String columnName) {
+ return jdbcTemplate.queryForObject(
+ "SELECT " + columnName + " FROM etl_job_records WHERE job_record_id=?",
+ (resultSet, rowNumber) -> resultSet.getTimestamp(columnName).toInstant(),
+ jobRecordId
+ );
+ }
+
+ /** Minimal transaction-enabled Spring context for replay persistence verification. */
+ @Configuration
+ @EnableTransactionManagement
+ static class TestConfiguration {
+
+ @Bean
+ DataSource dataSource() {
+ return new EmbeddedDatabaseBuilder()
+ .generateUniqueName(true)
+ .setType(EmbeddedDatabaseType.H2)
+ .build();
+ }
+
+ @Bean
+ JdbcTemplate jdbcTemplate(DataSource dataSource) {
+ return new JdbcTemplate(dataSource);
+ }
+
+ @Bean
+ PlatformTransactionManager transactionManager(DataSource dataSource) {
+ return new DataSourceTransactionManager(dataSource);
+ }
+
+ @Bean
+ ObjectMapper objectMapper() {
+ return new ObjectMapper();
+ }
+
+ @Bean
+ EtlBatchProperties etlBatchProperties() {
+ return new EtlBatchProperties();
+ }
+
+ @Bean
+ EtlRequestLock etlRequestLock() {
+ return lockHash -> true;
+ }
+
+ @Bean
+ EtlJobReplayService replayService(
+ JdbcTemplate jdbcTemplate,
+ ObjectMapper objectMapper,
+ EtlBatchProperties batchProperties,
+ EtlRequestLock requestLock
+ ) {
+ return new EtlJobReplayService(
+ jdbcTemplate,
+ objectMapper,
+ batchProperties,
+ requestLock
+ );
+ }
+ }
+}
From 97ba5563b9bc7d2af17d55a1140af3cd72189162 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:41:31 +0900
Subject: [PATCH 11/32] test(etl): transition replay persistence contract
---
.../etl/job/EtlJobReplayBoundaryTest.java | 33 +----------
...tlJobReplayPersistenceIntegrationTest.java | 55 ++++++++++++++-----
2 files changed, 41 insertions(+), 47 deletions(-)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index 058d2326..c0bac809 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -15,7 +15,7 @@
import static org.mockito.Mockito.verifyNoInteractions;
/**
- * Captures immutable replay-result, construction, and pre-persistence validation contracts.
+ * Captures immutable replay-result, construction, and pre-database validation contracts.
*/
class EtlJobReplayBoundaryTest {
@@ -24,8 +24,6 @@ class EtlJobReplayBoundaryTest {
);
private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
- private static final String PERSISTENCE_NOT_AVAILABLE =
- "Durable ETL job replay persistence is not yet available";
@Test
void replayResultRequiresIdentityAndStatusButPreservesCurrentState() {
@@ -162,35 +160,6 @@ public JsonNode readTree(String content) {
verifyNoInteractions(jdbcTemplate);
}
- @Test
- void acceptsRawAndStructuredReplayKeysThenFailsClosedBeforePersistence() {
- JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
- EtlJobReplayService service = new EtlJobReplayService(
- jdbcTemplate,
- new ObjectMapper(),
- new EtlBatchProperties(),
- hash -> true
- );
-
- IllegalStateException rawKeyFailure = assertThrows(
- IllegalStateException.class,
- () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, "tenant_alpha")
- );
- assertEquals(PERSISTENCE_NOT_AVAILABLE, rawKeyFailure.getMessage());
-
- IllegalStateException structuredKeyFailure = assertThrows(
- IllegalStateException.class,
- () -> service.replayOwned(
- SOURCE_ID,
- PAYLOAD,
- "\"" + REPLAY_KEY + "\"",
- "tenant_alpha"
- )
- );
- assertEquals(PERSISTENCE_NOT_AVAILABLE, structuredKeyFailure.getMessage());
- verifyNoInteractions(jdbcTemplate);
- }
-
private static void assertErrorCode(String expectedCode, Runnable invocation) {
EtlRequestException exception = assertThrows(EtlRequestException.class, invocation::run);
assertEquals(expectedCode, exception.getMessage());
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
index fe37ed89..1d00bda4 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -34,6 +34,7 @@ class EtlJobReplayPersistenceIntegrationTest {
private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
+ private static final String OTHER_REPLAY_KEY = "519bc126-1398-4b4e-a4e3-1fb18a00f19b";
private static final String PRINCIPAL = "tenant_alpha";
private final EtlJobReplayService replayService;
@@ -80,21 +81,7 @@ cancellation_code VARCHAR(128),
@Test
void createsFirstGenerationPendingReplayWithoutMutatingTerminalSource() {
- UUID sourceJobRecordId = UUID.randomUUID();
- jdbcTemplate.update(
- """
- INSERT INTO etl_job_records (
- job_record_id, principal_scope_hash, submission_key_hash,
- request_digest, request_payload, job_status, attempt_count,
- failure_code
- ) VALUES (?, ?, ?, ?, NULL, 'FAILED', 1, ?)
- """,
- sourceJobRecordId,
- Sha256Digest.digest(PRINCIPAL),
- Sha256Digest.digest("original-submission-key"),
- Sha256Digest.digest(PAYLOAD),
- "etl_target_failure"
- );
+ UUID sourceJobRecordId = insertFailedSource("original-submission-key");
Instant sourceUpdatedAt = instantColumn(sourceJobRecordId, "updated_at");
EtlJobReplay replay = replayService.replayOwned(
@@ -122,6 +109,44 @@ INSERT INTO etl_job_records (
assertEquals(sourceUpdatedAt, instantColumn(sourceJobRecordId, "updated_at"));
}
+ @Test
+ void acceptsStructuredReplayKeyForIndependentFailedSource() {
+ UUID sourceJobRecordId = insertFailedSource("second-original-submission-key");
+
+ EtlJobReplay replay = replayService.replayOwned(
+ sourceJobRecordId,
+ PAYLOAD,
+ "\"" + OTHER_REPLAY_KEY + "\"",
+ PRINCIPAL
+ );
+
+ assertFalse(replay.replayed());
+ assertEquals(EtlJobStatus.PENDING, replay.jobStatus());
+ assertEquals(
+ sourceJobRecordId,
+ uuidColumn(replay.jobRecordId(), "replay_source_job_record_id")
+ );
+ }
+
+ private UUID insertFailedSource(String submissionKey) {
+ UUID sourceJobRecordId = UUID.randomUUID();
+ jdbcTemplate.update(
+ """
+ INSERT INTO etl_job_records (
+ job_record_id, principal_scope_hash, submission_key_hash,
+ request_digest, request_payload, job_status, attempt_count,
+ failure_code
+ ) VALUES (?, ?, ?, ?, NULL, 'FAILED', 1, ?)
+ """,
+ sourceJobRecordId,
+ Sha256Digest.digest(PRINCIPAL),
+ Sha256Digest.digest(submissionKey),
+ Sha256Digest.digest(PAYLOAD),
+ "etl_target_failure"
+ );
+ return sourceJobRecordId;
+ }
+
private String textColumn(UUID jobRecordId, String columnName) {
return jdbcTemplate.queryForObject(
"SELECT " + columnName + " FROM etl_job_records WHERE job_record_id=?",
From 8e1978e50760c83d42f25fee20eeb8b7674af2ec Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 18:45:14 +0900
Subject: [PATCH 12/32] feat(etl): persist first-generation failed-job replay
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 85 ++++++++++++++++---
1 file changed, 72 insertions(+), 13 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index 5e8f40ac..60b47caa 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -9,10 +9,12 @@
import com.xtrmetl.etl.service.EtlRequestException;
import com.xtrmetl.etl.service.EtlRequestLock;
import com.xtrmetl.etl.service.PostgresEtlRequestLock;
+import com.xtrmetl.etl.service.Sha256Digest;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.lang.Nullable;
import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
import java.util.Objects;
import java.util.UUID;
@@ -22,8 +24,10 @@
* Admits owner-scoped durable-job replay requests on the repaired stack.
*
* The service validates replay identity, authenticated principal scope, and JSON shape before
- * persistence work. Database admission and replay transitions are intentionally added in later
- * test-first increments so an incomplete branch cannot silently claim the full replay contract.
+ * database work. The first persistence increment admits only an owner's first-generation replay of
+ * a {@code FAILED} source whose resupplied payload has the exact source request digest. Later
+ * test-first increments add the remaining terminal states, idempotent replay lookup, concurrency,
+ * and generation-depth policies before the HTTP replay resource is exposed.
*/
@Service
public class EtlJobReplayService {
@@ -36,6 +40,33 @@ public class EtlJobReplayService {
private static final Pattern IDEMPOTENCY_KEY_STRUCTURED_FIELD_PROFILE = Pattern.compile(
"\"(" + IDEMPOTENCY_KEY_VALUE_EXPRESSION + ")\""
);
+ private static final String REPLAY_KEY_HASH_DOMAIN = "mightyetl:durable-job-replay-key:v1:";
+ private static final String SELECT_FIRST_GENERATION_FAILED_SOURCE_SQL = """
+ SELECT job_record_id
+ FROM etl_job_records
+ WHERE job_record_id = ?
+ AND principal_scope_hash = ?
+ AND request_digest = ?
+ AND job_status = 'FAILED'
+ AND replay_source_job_record_id IS NULL
+ AND replay_root_job_record_id IS NULL
+ AND replay_generation_count IS NULL
+ FOR UPDATE
+ """;
+ private static final String INSERT_FIRST_GENERATION_REPLAY_SQL = """
+ INSERT INTO etl_job_records (
+ job_record_id,
+ principal_scope_hash,
+ submission_key_hash,
+ request_digest,
+ request_payload,
+ job_status,
+ attempt_count,
+ replay_source_job_record_id,
+ replay_root_job_record_id,
+ replay_generation_count
+ ) VALUES (?, ?, ?, ?, ?, 'PENDING', 0, ?, ?, 1)
+ """;
private final JdbcTemplate jdbcTemplate;
private final ObjectMapper objectMapper;
@@ -71,7 +102,7 @@ public EtlJobReplayService(
* @param jdbcTemplate parameterized durable job persistence
* @param objectMapper JSON parser configuration to copy
* @param batchProperties bounded request limits
- * @param requestLock transaction-lifetime replay-key lock
+ * @param requestLock transaction-lifetime replay-key lock reserved for the idempotency increment
*/
@Autowired
public EtlJobReplayService(
@@ -95,32 +126,60 @@ public EtlJobReplayService(
}
/**
- * Validates one owner-scoped replay request before any persistence work.
+ * Creates one first-generation replay of an owner-scoped failed durable job.
*
- * The database admission and replay transition are introduced in later TDD increments.
- * Until then, an otherwise valid request fails closed rather than acquiring a lock or touching
- * durable state.
+ * Authentication scope and byte-exact payload digest are part of the source-row predicate,
+ * and the source row is locked before the child is inserted. The source row is never mutated.
+ * Replay-key normalization is domain-separated before reuse of the existing durable submission
+ * key column so replay and ordinary submission key material cannot collide accidentally.
*
* @param sourceJobRecordId terminal durable job whose intent is being replayed
* @param requestBody bounded JSON-array payload resupplied by the authenticated owner
* @param rawReplayKey principal-scoped idempotency key for this replay request
* @param principalName authenticated principal namespace
- * @return replay-created durable job result after persistence support is introduced
+ * @return newly created pending replay job
* @throws NullPointerException when the source job identifier is absent
* @throws EtlRequestException when key, principal, or JSON validation fails
- * @throws IllegalStateException while the persistence increment is intentionally unavailable
*/
+ @Transactional
public EtlJobReplay replayOwned(
UUID sourceJobRecordId,
@Nullable String requestBody,
@Nullable String rawReplayKey,
@Nullable String principalName
) {
- Objects.requireNonNull(sourceJobRecordId, "sourceJobRecordId must not be null");
- validateReplayKey(rawReplayKey);
- validatePrincipalScope(principalName);
+ UUID validatedSourceJobRecordId = Objects.requireNonNull(
+ sourceJobRecordId,
+ "sourceJobRecordId must not be null"
+ );
+ String validatedReplayKey = validateReplayKey(rawReplayKey);
+ String validatedPrincipalName = validatePrincipalScope(principalName);
validateReplayPayload(requestBody);
- throw new IllegalStateException("Durable ETL job replay persistence is not yet available");
+
+ String principalScopeHash = Sha256Digest.digest(validatedPrincipalName);
+ String requestDigest = Sha256Digest.digest(requestBody);
+ UUID lockedSourceJobRecordId = jdbcTemplate.queryForObject(
+ SELECT_FIRST_GENERATION_FAILED_SOURCE_SQL,
+ UUID.class,
+ validatedSourceJobRecordId,
+ principalScopeHash,
+ requestDigest
+ );
+ UUID replayJobRecordId = UUID.randomUUID();
+ String replayKeyHash = Sha256Digest.digest(
+ REPLAY_KEY_HASH_DOMAIN + principalScopeHash + ":" + validatedReplayKey
+ );
+ jdbcTemplate.update(
+ INSERT_FIRST_GENERATION_REPLAY_SQL,
+ replayJobRecordId,
+ principalScopeHash,
+ replayKeyHash,
+ requestDigest,
+ requestBody,
+ lockedSourceJobRecordId,
+ lockedSourceJobRecordId
+ );
+ return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false);
}
private static String validateReplayKey(@Nullable String rawReplayKey) {
From 135d6e4b446e092e1e3ce48e30cb15f3448116e6 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 19:13:43 +0900
Subject: [PATCH 13/32] test(etl): require cancelled-source replay
---
...tlJobReplayPersistenceIntegrationTest.java | 55 +++++++++++++++++++
1 file changed, 55 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
index 1d00bda4..04308af2 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -21,6 +21,7 @@
import java.time.Instant;
import java.util.UUID;
+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.assertNotEquals;
@@ -128,6 +129,40 @@ void acceptsStructuredReplayKeyForIndependentFailedSource() {
);
}
+ @Test
+ void createsFirstGenerationPendingReplayFromCancelledSourceWithoutMutatingCancellationEvidence() {
+ UUID sourceJobRecordId = insertCancelledSource();
+ Instant sourceUpdatedAt = instantColumn(sourceJobRecordId, "updated_at");
+
+ EtlJobReplay replay = assertDoesNotThrow(() -> replayService.replayOwned(
+ sourceJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ ));
+
+ assertFalse(replay.replayed());
+ assertEquals(EtlJobStatus.PENDING, replay.jobStatus());
+ assertNotEquals(sourceJobRecordId, replay.jobRecordId());
+ assertEquals(
+ sourceJobRecordId,
+ uuidColumn(replay.jobRecordId(), "replay_source_job_record_id")
+ );
+ assertEquals(
+ sourceJobRecordId,
+ uuidColumn(replay.jobRecordId(), "replay_root_job_record_id")
+ );
+ assertEquals(1, integerColumn(replay.jobRecordId(), "replay_generation_count"));
+ assertEquals(PAYLOAD, textColumn(replay.jobRecordId(), "request_payload"));
+ assertEquals("CANCELLED", textColumn(sourceJobRecordId, "job_status"));
+ assertEquals(
+ EtlJobService.CANCELLED_BY_OWNER_CODE,
+ textColumn(sourceJobRecordId, "cancellation_code")
+ );
+ assertNull(textColumn(sourceJobRecordId, "request_payload"));
+ assertEquals(sourceUpdatedAt, instantColumn(sourceJobRecordId, "updated_at"));
+ }
+
private UUID insertFailedSource(String submissionKey) {
UUID sourceJobRecordId = UUID.randomUUID();
jdbcTemplate.update(
@@ -147,6 +182,26 @@ INSERT INTO etl_job_records (
return sourceJobRecordId;
}
+ private UUID insertCancelledSource() {
+ UUID sourceJobRecordId = UUID.randomUUID();
+ jdbcTemplate.update(
+ """
+ INSERT INTO etl_job_records (
+ job_record_id, principal_scope_hash, submission_key_hash,
+ request_digest, request_payload, job_status, attempt_count,
+ cancellation_key_hash, cancellation_code, job_cancelled_at
+ ) VALUES (?, ?, ?, ?, NULL, 'CANCELLED', 1, ?, ?, CURRENT_TIMESTAMP)
+ """,
+ sourceJobRecordId,
+ Sha256Digest.digest(PRINCIPAL),
+ Sha256Digest.digest("cancelled-original-submission-key"),
+ Sha256Digest.digest(PAYLOAD),
+ "d".repeat(64),
+ EtlJobService.CANCELLED_BY_OWNER_CODE
+ );
+ return sourceJobRecordId;
+ }
+
private String textColumn(UUID jobRecordId, String columnName) {
return jdbcTemplate.queryForObject(
"SELECT " + columnName + " FROM etl_job_records WHERE job_record_id=?",
From 13a8c1a4d914bd52c6d46073c2032c26656e5ff6 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 19:16:54 +0900
Subject: [PATCH 14/32] feat(etl): replay cancelled terminal jobs
---
.../com/xtrmetl/etl/job/EtlJobReplayService.java | 14 +++++++-------
1 file changed, 7 insertions(+), 7 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index 60b47caa..e5e859fc 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -25,9 +25,9 @@
*
* The service validates replay identity, authenticated principal scope, and JSON shape before
* database work. The first persistence increment admits only an owner's first-generation replay of
- * a {@code FAILED} source whose resupplied payload has the exact source request digest. Later
- * test-first increments add the remaining terminal states, idempotent replay lookup, concurrency,
- * and generation-depth policies before the HTTP replay resource is exposed.
+ * a {@code FAILED} or {@code CANCELLED} source whose resupplied payload has the exact source request
+ * digest. Later test-first increments add idempotent replay lookup, concurrency, and
+ * generation-depth policies before the HTTP replay resource is exposed.
*/
@Service
public class EtlJobReplayService {
@@ -41,13 +41,13 @@ public class EtlJobReplayService {
"\"(" + IDEMPOTENCY_KEY_VALUE_EXPRESSION + ")\""
);
private static final String REPLAY_KEY_HASH_DOMAIN = "mightyetl:durable-job-replay-key:v1:";
- private static final String SELECT_FIRST_GENERATION_FAILED_SOURCE_SQL = """
+ private static final String SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL = """
SELECT job_record_id
FROM etl_job_records
WHERE job_record_id = ?
AND principal_scope_hash = ?
AND request_digest = ?
- AND job_status = 'FAILED'
+ AND job_status IN ('FAILED', 'CANCELLED')
AND replay_source_job_record_id IS NULL
AND replay_root_job_record_id IS NULL
AND replay_generation_count IS NULL
@@ -126,7 +126,7 @@ public EtlJobReplayService(
}
/**
- * Creates one first-generation replay of an owner-scoped failed durable job.
+ * Creates one first-generation replay of an owner-scoped failed or cancelled durable job.
*
* Authentication scope and byte-exact payload digest are part of the source-row predicate,
* and the source row is locked before the child is inserted. The source row is never mutated.
@@ -159,7 +159,7 @@ public EtlJobReplay replayOwned(
String principalScopeHash = Sha256Digest.digest(validatedPrincipalName);
String requestDigest = Sha256Digest.digest(requestBody);
UUID lockedSourceJobRecordId = jdbcTemplate.queryForObject(
- SELECT_FIRST_GENERATION_FAILED_SOURCE_SQL,
+ SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL,
UUID.class,
validatedSourceJobRecordId,
principalScopeHash,
From f0b9ffb4845eac885b0112b2d1dcceb6b679b52f Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:10:10 +0900
Subject: [PATCH 15/32] test(etl): require replay idempotent lookup
---
...tlJobReplayPersistenceIntegrationTest.java | 32 +++++++++++++++++++
1 file changed, 32 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
index 04308af2..f8852f9c 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -26,6 +26,7 @@
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Proves the first database-owned durable-job replay transition on the repaired stack.
@@ -110,6 +111,37 @@ void createsFirstGenerationPendingReplayWithoutMutatingTerminalSource() {
assertEquals(sourceUpdatedAt, instantColumn(sourceJobRecordId, "updated_at"));
}
+ @Test
+ void replaysSameOwnerSourceKeyAndPayloadWithoutSecondInsert() {
+ UUID sourceJobRecordId = insertFailedSource("idempotent-source-submission-key");
+
+ EtlJobReplay first = replayService.replayOwned(
+ sourceJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ );
+ EtlJobReplay replay = replayService.replayOwned(
+ sourceJobRecordId,
+ PAYLOAD,
+ "\"" + REPLAY_KEY + "\"",
+ PRINCIPAL
+ );
+
+ assertFalse(first.replayed());
+ assertTrue(replay.replayed());
+ assertEquals(first.jobRecordId(), replay.jobRecordId());
+ assertEquals(EtlJobStatus.PENDING, replay.jobStatus());
+ assertEquals(
+ 1,
+ jdbcTemplate.queryForObject(
+ "SELECT COUNT(*) FROM etl_job_records "
+ + "WHERE replay_generation_count IS NOT NULL",
+ Integer.class
+ )
+ );
+ }
+
@Test
void acceptsStructuredReplayKeyForIndependentFailedSource() {
UUID sourceJobRecordId = insertFailedSource("second-original-submission-key");
From 3d6bfb606fad6266ca56fc61a6c4609163da43d2 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:11:51 +0900
Subject: [PATCH 16/32] feat(etl): replay an existing identical job
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 55 ++++++++++++++-----
1 file changed, 41 insertions(+), 14 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index e5e859fc..d8d46da7 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -16,6 +16,7 @@
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
+import java.util.List;
import java.util.Objects;
import java.util.UUID;
import java.util.regex.Pattern;
@@ -24,10 +25,11 @@
* Admits owner-scoped durable-job replay requests on the repaired stack.
*
* The service validates replay identity, authenticated principal scope, and JSON shape before
- * database work. The first persistence increment admits only an owner's first-generation replay of
- * a {@code FAILED} or {@code CANCELLED} source whose resupplied payload has the exact source request
- * digest. Later test-first increments add idempotent replay lookup, concurrency, and
- * generation-depth policies before the HTTP replay resource is exposed.
+ * database work. The first persistence increments admit an owner's first-generation replay of a
+ * {@code FAILED} or {@code CANCELLED} source and return the already-created job for an identical
+ * owner/source/key/payload replay. Later test-first increments add conflicting-key classification,
+ * concurrent replay serialization, and generation-depth policies before the HTTP replay resource
+ * is exposed.
*/
@Service
public class EtlJobReplayService {
@@ -41,6 +43,14 @@ public class EtlJobReplayService {
"\"(" + IDEMPOTENCY_KEY_VALUE_EXPRESSION + ")\""
);
private static final String REPLAY_KEY_HASH_DOMAIN = "mightyetl:durable-job-replay-key:v1:";
+ private static final String SELECT_IDENTICAL_EXISTING_REPLAY_SQL = """
+ SELECT job_record_id, job_status
+ FROM etl_job_records
+ WHERE principal_scope_hash = ?
+ AND submission_key_hash = ?
+ AND request_digest = ?
+ AND replay_source_job_record_id = ?
+ """;
private static final String SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL = """
SELECT job_record_id
FROM etl_job_records
@@ -102,7 +112,7 @@ public EtlJobReplayService(
* @param jdbcTemplate parameterized durable job persistence
* @param objectMapper JSON parser configuration to copy
* @param batchProperties bounded request limits
- * @param requestLock transaction-lifetime replay-key lock reserved for the idempotency increment
+ * @param requestLock transaction-lifetime replay-key lock reserved for the concurrency increment
*/
@Autowired
public EtlJobReplayService(
@@ -126,18 +136,19 @@ public EtlJobReplayService(
}
/**
- * Creates one first-generation replay of an owner-scoped failed or cancelled durable job.
+ * Creates or returns one first-generation replay of an owner-scoped terminal durable job.
*
- * Authentication scope and byte-exact payload digest are part of the source-row predicate,
- * and the source row is locked before the child is inserted. The source row is never mutated.
- * Replay-key normalization is domain-separated before reuse of the existing durable submission
- * key column so replay and ordinary submission key material cannot collide accidentally.
+ * An identical committed replay is looked up by the full owner/key/source/payload identity
+ * and returned without a second insert. A new replay still locks the immutable terminal source
+ * before inserting a fresh {@code PENDING} child. Conflicting replay-key reuse intentionally
+ * remains fail-closed at the database uniqueness boundary until its dedicated error contract is
+ * introduced test-first.
*
* @param sourceJobRecordId terminal durable job whose intent is being replayed
* @param requestBody bounded JSON-array payload resupplied by the authenticated owner
* @param rawReplayKey principal-scoped idempotency key for this replay request
* @param principalName authenticated principal namespace
- * @return newly created pending replay job
+ * @return newly created or previously committed replay job
* @throws NullPointerException when the source job identifier is absent
* @throws EtlRequestException when key, principal, or JSON validation fails
*/
@@ -158,6 +169,25 @@ public EtlJobReplay replayOwned(
String principalScopeHash = Sha256Digest.digest(validatedPrincipalName);
String requestDigest = Sha256Digest.digest(requestBody);
+ String replayKeyHash = Sha256Digest.digest(
+ REPLAY_KEY_HASH_DOMAIN + principalScopeHash + ":" + validatedReplayKey
+ );
+ List existingReplays = jdbcTemplate.query(
+ SELECT_IDENTICAL_EXISTING_REPLAY_SQL,
+ (resultSet, rowNumber) -> new EtlJobReplay(
+ resultSet.getObject("job_record_id", UUID.class),
+ EtlJobStatus.valueOf(resultSet.getString("job_status")),
+ true
+ ),
+ principalScopeHash,
+ replayKeyHash,
+ requestDigest,
+ validatedSourceJobRecordId
+ );
+ if (!existingReplays.isEmpty()) {
+ return existingReplays.getFirst();
+ }
+
UUID lockedSourceJobRecordId = jdbcTemplate.queryForObject(
SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL,
UUID.class,
@@ -166,9 +196,6 @@ public EtlJobReplay replayOwned(
requestDigest
);
UUID replayJobRecordId = UUID.randomUUID();
- String replayKeyHash = Sha256Digest.digest(
- REPLAY_KEY_HASH_DOMAIN + principalScopeHash + ":" + validatedReplayKey
- );
jdbcTemplate.update(
INSERT_FIRST_GENERATION_REPLAY_SQL,
replayJobRecordId,
From 884caa9917ab04d578dcc7cd1153552e943bf79d Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:15:06 +0900
Subject: [PATCH 17/32] test(etl): reject replay key reuse across sources
---
...tlJobReplayPersistenceIntegrationTest.java | 34 +++++++++++++++++++
1 file changed, 34 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
index f8852f9c..033473e8 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -2,6 +2,7 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import com.xtrmetl.etl.service.EtlBatchProperties;
+import com.xtrmetl.etl.service.EtlRequestException;
import com.xtrmetl.etl.service.EtlRequestLock;
import com.xtrmetl.etl.service.Sha256Digest;
import org.junit.jupiter.api.BeforeEach;
@@ -26,6 +27,7 @@
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
@@ -142,6 +144,38 @@ void replaysSameOwnerSourceKeyAndPayloadWithoutSecondInsert() {
);
}
+ @Test
+ void rejectsReplayKeyReuseAcrossDifferentSources() {
+ UUID firstSourceJobRecordId = insertFailedSource("reused-key-first-source");
+ replayService.replayOwned(
+ firstSourceJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ );
+ UUID secondSourceJobRecordId = insertFailedSource("reused-key-second-source");
+
+ EtlRequestException exception = assertThrows(
+ EtlRequestException.class,
+ () -> replayService.replayOwned(
+ secondSourceJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ )
+ );
+
+ assertEquals("etl_job_replay_key_reused", exception.getMessage());
+ assertEquals(
+ 1,
+ jdbcTemplate.queryForObject(
+ "SELECT COUNT(*) FROM etl_job_records "
+ + "WHERE replay_generation_count IS NOT NULL",
+ Integer.class
+ )
+ );
+ }
+
@Test
void acceptsStructuredReplayKeyForIndependentFailedSource() {
UUID sourceJobRecordId = insertFailedSource("second-original-submission-key");
From ad52079de7036c1180f841afee4b8199051c4d1b Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:17:37 +0900
Subject: [PATCH 18/32] test(etl): reject replay key reuse across payloads
---
...tlJobReplayPersistenceIntegrationTest.java | 32 +++++++++++++++++++
1 file changed, 32 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
index 033473e8..d147fa87 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -37,6 +37,7 @@
class EtlJobReplayPersistenceIntegrationTest {
private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
+ private static final String OTHER_PAYLOAD = "[{\"id\":\"record_beta\"}]";
private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
private static final String OTHER_REPLAY_KEY = "519bc126-1398-4b4e-a4e3-1fb18a00f19b";
private static final String PRINCIPAL = "tenant_alpha";
@@ -176,6 +177,37 @@ void rejectsReplayKeyReuseAcrossDifferentSources() {
);
}
+ @Test
+ void rejectsReplayKeyReuseForDifferentPayloadOnSameSource() {
+ UUID sourceJobRecordId = insertFailedSource("reused-key-payload-source");
+ replayService.replayOwned(
+ sourceJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ );
+
+ EtlRequestException exception = assertThrows(
+ EtlRequestException.class,
+ () -> replayService.replayOwned(
+ sourceJobRecordId,
+ OTHER_PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ )
+ );
+
+ assertEquals("etl_job_replay_key_reused", exception.getMessage());
+ assertEquals(
+ 1,
+ jdbcTemplate.queryForObject(
+ "SELECT COUNT(*) FROM etl_job_records "
+ + "WHERE replay_generation_count IS NOT NULL",
+ Integer.class
+ )
+ );
+ }
+
@Test
void acceptsStructuredReplayKeyForIndependentFailedSource() {
UUID sourceJobRecordId = insertFailedSource("second-original-submission-key");
From f95dd0155ea42b57c400a9b016bce3435ea134d4 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:20:20 +0900
Subject: [PATCH 19/32] feat(etl): classify replay key reuse
---
.../java/com/xtrmetl/etl/service/EtlRequestError.java | 9 +++++++++
1 file changed, 9 insertions(+)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
index e62e89f5..73d523c6 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
@@ -140,6 +140,15 @@ public enum EtlRequestError {
"Replay requires a supported principal-scoped Idempotency-Key."
),
+ /** A replay idempotency key already identifies a different source or payload. */
+ JOB_REPLAY_KEY_REUSED(
+ HttpStatus.UNPROCESSABLE_ENTITY,
+ "etl_job_replay_key_reused",
+ "urn:mightyetl:problem:etl-job-replay-key-reused",
+ "ETL job replay key reused",
+ "The Idempotency-Key already identifies a different durable job replay."
+ ),
+
/** The durable job was already cancelled with a different cancellation key. */
JOB_CANCELLATION_KEY_REUSED(
HttpStatus.UNPROCESSABLE_ENTITY,
From 5be293e4a913d017aee4bec927faab49e1db1dc5 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:20:59 +0900
Subject: [PATCH 20/32] feat(etl): reject conflicting replay key reuse
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 50 +++++++++++--------
1 file changed, 28 insertions(+), 22 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index d8d46da7..d0b6eb3b 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -26,10 +26,10 @@
*
* The service validates replay identity, authenticated principal scope, and JSON shape before
* database work. The first persistence increments admit an owner's first-generation replay of a
- * {@code FAILED} or {@code CANCELLED} source and return the already-created job for an identical
- * owner/source/key/payload replay. Later test-first increments add conflicting-key classification,
- * concurrent replay serialization, and generation-depth policies before the HTTP replay resource
- * is exposed.
+ * {@code FAILED} or {@code CANCELLED} source, return the already-created job for an identical
+ * owner/source/key/payload replay, and reject a committed replay key that identifies a different
+ * source or payload. Later test-first increments add concurrent replay serialization and
+ * generation-depth policies before the HTTP replay resource is exposed.
*/
@Service
public class EtlJobReplayService {
@@ -43,13 +43,11 @@ public class EtlJobReplayService {
"\"(" + IDEMPOTENCY_KEY_VALUE_EXPRESSION + ")\""
);
private static final String REPLAY_KEY_HASH_DOMAIN = "mightyetl:durable-job-replay-key:v1:";
- private static final String SELECT_IDENTICAL_EXISTING_REPLAY_SQL = """
- SELECT job_record_id, job_status
+ private static final String SELECT_EXISTING_REPLAY_SQL = """
+ SELECT job_record_id, request_digest, job_status, replay_source_job_record_id
FROM etl_job_records
WHERE principal_scope_hash = ?
AND submission_key_hash = ?
- AND request_digest = ?
- AND replay_source_job_record_id = ?
""";
private static final String SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL = """
SELECT job_record_id
@@ -139,10 +137,9 @@ public EtlJobReplayService(
* Creates or returns one first-generation replay of an owner-scoped terminal durable job.
*
* An identical committed replay is looked up by the full owner/key/source/payload identity
- * and returned without a second insert. A new replay still locks the immutable terminal source
- * before inserting a fresh {@code PENDING} child. Conflicting replay-key reuse intentionally
- * remains fail-closed at the database uniqueness boundary until its dedicated error contract is
- * introduced test-first.
+ * and returned without a second insert. A committed key bound to another source or payload is
+ * rejected with {@link EtlRequestError#JOB_REPLAY_KEY_REUSED}. A new replay still locks the
+ * immutable terminal source before inserting a fresh {@code PENDING} child.
*
* @param sourceJobRecordId terminal durable job whose intent is being replayed
* @param requestBody bounded JSON-array payload resupplied by the authenticated owner
@@ -150,7 +147,7 @@ public EtlJobReplayService(
* @param principalName authenticated principal namespace
* @return newly created or previously committed replay job
* @throws NullPointerException when the source job identifier is absent
- * @throws EtlRequestException when key, principal, or JSON validation fails
+ * @throws EtlRequestException when validation fails or a replay key is reused for other intent
*/
@Transactional
public EtlJobReplay replayOwned(
@@ -173,16 +170,25 @@ public EtlJobReplay replayOwned(
REPLAY_KEY_HASH_DOMAIN + principalScopeHash + ":" + validatedReplayKey
);
List existingReplays = jdbcTemplate.query(
- SELECT_IDENTICAL_EXISTING_REPLAY_SQL,
- (resultSet, rowNumber) -> new EtlJobReplay(
- resultSet.getObject("job_record_id", UUID.class),
- EtlJobStatus.valueOf(resultSet.getString("job_status")),
- true
- ),
+ SELECT_EXISTING_REPLAY_SQL,
+ (resultSet, rowNumber) -> {
+ UUID existingSourceJobRecordId = resultSet.getObject(
+ "replay_source_job_record_id",
+ UUID.class
+ );
+ String existingRequestDigest = resultSet.getString("request_digest");
+ if (!validatedSourceJobRecordId.equals(existingSourceJobRecordId)
+ || !requestDigest.equals(existingRequestDigest)) {
+ throw new EtlRequestException(EtlRequestError.JOB_REPLAY_KEY_REUSED);
+ }
+ return new EtlJobReplay(
+ resultSet.getObject("job_record_id", UUID.class),
+ EtlJobStatus.valueOf(resultSet.getString("job_status")),
+ true
+ );
+ },
principalScopeHash,
- replayKeyHash,
- requestDigest,
- validatedSourceJobRecordId
+ replayKeyHash
);
if (!existingReplays.isEmpty()) {
return existingReplays.getFirst();
From e81c0b570ba5f21983e3a0c5a1d88253270726b3 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:22:55 +0900
Subject: [PATCH 21/32] test(etl): fail replay while key is in progress
---
.../etl/job/EtlJobReplayBoundaryTest.java | 17 +++++++++++++++++
1 file changed, 17 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index c0bac809..3fd5c706 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -132,6 +132,23 @@ void validatesIdentityKeyPrincipalAndPayloadBeforeDatabaseWork() {
verifyNoInteractions(jdbcTemplate);
}
+ @Test
+ void rejectsCompetingReplayKeyBeforeDatabaseWork() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ new ObjectMapper(),
+ new EtlBatchProperties(),
+ hash -> false
+ );
+
+ assertErrorCode(
+ "etl_job_replay_in_progress",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, "tenant_alpha")
+ );
+ verifyNoInteractions(jdbcTemplate);
+ }
+
@Test
void rejectsAnAbsentParsedRootBeforeDatabaseWork() {
JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
From dd0b2dda20934d859ed01391b59122176dbbb81e Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:25:02 +0900
Subject: [PATCH 22/32] feat(etl): classify replay creation contention
---
.../java/com/xtrmetl/etl/service/EtlRequestError.java | 9 +++++++++
1 file changed, 9 insertions(+)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
index 73d523c6..0e554c85 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
@@ -149,6 +149,15 @@ public enum EtlRequestError {
"The Idempotency-Key already identifies a different durable job replay."
),
+ /** Another transaction owns the same principal-scoped replay identity. */
+ JOB_REPLAY_IN_PROGRESS(
+ HttpStatus.CONFLICT,
+ "etl_job_replay_in_progress",
+ "urn:mightyetl:problem:etl-job-replay-in-progress",
+ "ETL job replay in progress",
+ "A durable job replay with the same principal-scoped Idempotency-Key is being created."
+ ),
+
/** The durable job was already cancelled with a different cancellation key. */
JOB_CANCELLATION_KEY_REUSED(
HttpStatus.UNPROCESSABLE_ENTITY,
From a8029f2b5eeda2c862bfdf7887c525442b0eabd3 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:25:54 +0900
Subject: [PATCH 23/32] feat(etl): serialize replay creation by scoped key
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 30 ++++++++++++-------
1 file changed, 20 insertions(+), 10 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index d0b6eb3b..9dba4a45 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -26,10 +26,11 @@
*
* The service validates replay identity, authenticated principal scope, and JSON shape before
* database work. The first persistence increments admit an owner's first-generation replay of a
- * {@code FAILED} or {@code CANCELLED} source, return the already-created job for an identical
- * owner/source/key/payload replay, and reject a committed replay key that identifies a different
- * source or payload. Later test-first increments add concurrent replay serialization and
- * generation-depth policies before the HTTP replay resource is exposed.
+ * {@code FAILED} or {@code CANCELLED} source, serialize creation by principal-scoped replay key,
+ * return the already-created job for an identical owner/source/key/payload replay, and reject a
+ * committed replay key that identifies a different source or payload. Later test-first increments
+ * add source-state classification and generation-depth policies before the HTTP replay resource is
+ * exposed.
*/
@Service
public class EtlJobReplayService {
@@ -43,6 +44,7 @@ public class EtlJobReplayService {
"\"(" + IDEMPOTENCY_KEY_VALUE_EXPRESSION + ")\""
);
private static final String REPLAY_KEY_HASH_DOMAIN = "mightyetl:durable-job-replay-key:v1:";
+ private static final String REPLAY_LOCK_HASH_DOMAIN = "mightyetl:durable-job-replay-lock:v1:";
private static final String SELECT_EXISTING_REPLAY_SQL = """
SELECT job_record_id, request_digest, job_status, replay_source_job_record_id
FROM etl_job_records
@@ -110,7 +112,7 @@ public EtlJobReplayService(
* @param jdbcTemplate parameterized durable job persistence
* @param objectMapper JSON parser configuration to copy
* @param batchProperties bounded request limits
- * @param requestLock transaction-lifetime replay-key lock reserved for the concurrency increment
+ * @param requestLock transaction-lifetime principal-scoped replay-key lock
*/
@Autowired
public EtlJobReplayService(
@@ -136,10 +138,11 @@ public EtlJobReplayService(
/**
* Creates or returns one first-generation replay of an owner-scoped terminal durable job.
*
- * An identical committed replay is looked up by the full owner/key/source/payload identity
- * and returned without a second insert. A committed key bound to another source or payload is
- * rejected with {@link EtlRequestError#JOB_REPLAY_KEY_REUSED}. A new replay still locks the
- * immutable terminal source before inserting a fresh {@code PENDING} child.
+ * The principal-scoped replay identity is serialized before table access. An identical
+ * committed replay is then looked up by owner/key/source/payload identity and returned without
+ * a second insert. A committed key bound to another source or payload is rejected with
+ * {@link EtlRequestError#JOB_REPLAY_KEY_REUSED}. A new replay locks the immutable terminal
+ * source before inserting a fresh {@code PENDING} child.
*
* @param sourceJobRecordId terminal durable job whose intent is being replayed
* @param requestBody bounded JSON-array payload resupplied by the authenticated owner
@@ -147,7 +150,7 @@ public EtlJobReplayService(
* @param principalName authenticated principal namespace
* @return newly created or previously committed replay job
* @throws NullPointerException when the source job identifier is absent
- * @throws EtlRequestException when validation fails or a replay key is reused for other intent
+ * @throws EtlRequestException when validation fails, the key is busy, or intent conflicts
*/
@Transactional
public EtlJobReplay replayOwned(
@@ -169,6 +172,13 @@ public EtlJobReplay replayOwned(
String replayKeyHash = Sha256Digest.digest(
REPLAY_KEY_HASH_DOMAIN + principalScopeHash + ":" + validatedReplayKey
);
+ String replayLockHash = Sha256Digest.digest(
+ REPLAY_LOCK_HASH_DOMAIN + principalScopeHash + ":" + replayKeyHash
+ );
+ if (!requestLock.tryLock(replayLockHash)) {
+ throw new EtlRequestException(EtlRequestError.JOB_REPLAY_IN_PROGRESS);
+ }
+
List existingReplays = jdbcTemplate.query(
SELECT_EXISTING_REPLAY_SQL,
(resultSet, rowNumber) -> {
From f62f319cc2535623ee959c54d7d7de39ffe3d874 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 20:28:29 +0900
Subject: [PATCH 24/32] test(etl): bound replay payload before lock acquisition
---
.../etl/job/EtlJobReplayBoundaryTest.java | 60 +++++++++++++++++++
1 file changed, 60 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
index 3fd5c706..516b88e2 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -132,6 +132,66 @@ void validatesIdentityKeyPrincipalAndPayloadBeforeDatabaseWork() {
verifyNoInteractions(jdbcTemplate);
}
+ @Test
+ void rejectsOversizedReplayPayloadBeforeLockOrDatabaseWork() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ EtlBatchProperties properties = new EtlBatchProperties();
+ properties.setMaxPayloadBytes(4);
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ new ObjectMapper(),
+ properties,
+ hash -> false
+ );
+
+ assertErrorCode(
+ "etl_payload_too_large",
+ () -> service.replayOwned(SOURCE_ID, PAYLOAD, REPLAY_KEY, "tenant_alpha")
+ );
+ verifyNoInteractions(jdbcTemplate);
+ }
+
+ @Test
+ void rejectsOversizedReplayBatchBeforeLockOrDatabaseWork() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ EtlBatchProperties properties = new EtlBatchProperties();
+ properties.setMaxBatchRecords(1);
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ new ObjectMapper(),
+ properties,
+ hash -> false
+ );
+
+ assertErrorCode(
+ "etl_batch_too_large",
+ () -> service.replayOwned(
+ SOURCE_ID,
+ "[{\"id\":\"record_alpha\"},{\"id\":\"record_beta\"}]",
+ REPLAY_KEY,
+ "tenant_alpha"
+ )
+ );
+ verifyNoInteractions(jdbcTemplate);
+ }
+
+ @Test
+ void rejectsInvalidReplayRecordBeforeLockOrDatabaseWork() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ EtlJobReplayService service = new EtlJobReplayService(
+ jdbcTemplate,
+ new ObjectMapper(),
+ new EtlBatchProperties(),
+ hash -> false
+ );
+
+ assertErrorCode(
+ "etl_invalid_record",
+ () -> service.replayOwned(SOURCE_ID, "[{}]", REPLAY_KEY, "tenant_alpha")
+ );
+ verifyNoInteractions(jdbcTemplate);
+ }
+
@Test
void rejectsCompetingReplayKeyBeforeDatabaseWork() {
JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
From 3a4d6ee0d56c66ce8594dccbf62539df6b6b69f2 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 22:06:27 +0900
Subject: [PATCH 25/32] fix(etl): validate replay payload before lock
acquisition
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 25 ++++++++++++++-----
1 file changed, 19 insertions(+), 6 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index 9dba4a45..52f23847 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -16,6 +16,7 @@
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
+import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Objects;
import java.util.UUID;
@@ -250,17 +251,29 @@ private static String validatePrincipalScope(@Nullable String principalName) {
}
private JsonNode validateReplayPayload(@Nullable String requestBody) {
- if (requestBody == null || requestBody.isEmpty()) {
+ if (requestBody == null) {
throw new EtlRequestException(EtlRequestError.INVALID_JSON);
}
+ int payloadBytes = requestBody.getBytes(StandardCharsets.UTF_8).length;
+ if (payloadBytes > batchProperties.getMaxPayloadBytes()) {
+ throw new EtlRequestException(EtlRequestError.PAYLOAD_TOO_LARGE);
+ }
+
+ final JsonNode root;
try {
- JsonNode root = objectMapper.readTree(requestBody);
- if (root == null || !root.isArray()) {
- throw new EtlRequestException(EtlRequestError.INVALID_JSON);
- }
- return root;
+ root = objectMapper.readTree(requestBody);
} catch (JsonProcessingException exception) {
throw new EtlRequestException(EtlRequestError.INVALID_JSON, exception);
}
+ if (root == null || root.isNull() || !root.isArray()) {
+ throw new EtlRequestException(EtlRequestError.INVALID_JSON);
+ }
+ if (root.size() > batchProperties.getMaxBatchRecords()) {
+ throw new EtlRequestException(EtlRequestError.BATCH_TOO_LARGE);
+ }
+ for (JsonNode record : root) {
+ EtlJobService.validateRecord(record);
+ }
+ return root;
}
}
From fd071ea1741b35de4be22a2d74fe838ddf09b12d Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 22:13:42 +0900
Subject: [PATCH 26/32] test(etl): classify replay source failures
---
...layStateClassificationIntegrationTest.java | 209 ++++++++++++++++++
1 file changed, 209 insertions(+)
create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
new file mode 100644
index 00000000..95f0b46e
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
@@ -0,0 +1,209 @@
+package com.xtrmetl.etl.job;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.xtrmetl.etl.service.EtlBatchProperties;
+import com.xtrmetl.etl.service.EtlRequestException;
+import com.xtrmetl.etl.service.EtlRequestLock;
+import com.xtrmetl.etl.service.Sha256Digest;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.jdbc.datasource.DataSourceTransactionManager;
+import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseBuilder;
+import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType;
+import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.annotation.EnableTransactionManagement;
+
+import javax.sql.DataSource;
+import java.util.UUID;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * Proves stable owner-safe source-state and payload classifications before replay insertion.
+ */
+@SpringJUnitConfig(EtlJobReplayStateClassificationIntegrationTest.TestConfiguration.class)
+class EtlJobReplayStateClassificationIntegrationTest {
+
+ private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]";
+ private static final String OTHER_PAYLOAD = "[{\"id\":\"record_beta\"}]";
+ private static final String REPLAY_KEY = "1e05bdca-447c-4ad3-882c-e33963ce517c";
+ private static final String PRINCIPAL = "tenant_alpha";
+ private static final String FOREIGN_PRINCIPAL = "tenant_beta";
+
+ private final EtlJobReplayService replayService;
+ private final JdbcTemplate jdbcTemplate;
+
+ @Autowired
+ EtlJobReplayStateClassificationIntegrationTest(
+ EtlJobReplayService replayService,
+ JdbcTemplate jdbcTemplate
+ ) {
+ this.replayService = replayService;
+ this.jdbcTemplate = jdbcTemplate;
+ }
+
+ @BeforeEach
+ void createDurableJobTable() {
+ jdbcTemplate.execute("DROP TABLE IF EXISTS etl_job_records");
+ jdbcTemplate.execute("""
+ CREATE TABLE etl_job_records (
+ job_record_id UUID PRIMARY KEY,
+ principal_scope_hash CHAR(64) NOT NULL,
+ submission_key_hash CHAR(64) NOT NULL,
+ request_digest CHAR(64) NOT NULL,
+ request_payload CLOB,
+ job_status VARCHAR(32) NOT NULL,
+ attempt_count INTEGER NOT NULL DEFAULT 0,
+ replay_source_job_record_id UUID,
+ replay_root_job_record_id UUID,
+ replay_generation_count INTEGER,
+ created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ updated_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ CONSTRAINT etl_job_submission_scope_unique
+ UNIQUE (principal_scope_hash, submission_key_hash)
+ )
+ """);
+ }
+
+ @Test
+ void missingAndForeignSourcesShareOwnerSafeNotFoundClassification() {
+ assertReplayError(
+ "etl_job_not_found",
+ UUID.randomUUID(),
+ PAYLOAD,
+ PRINCIPAL
+ );
+
+ UUID foreignSource = insertRootSource("FAILED", FOREIGN_PRINCIPAL, PAYLOAD);
+ assertReplayError("etl_job_not_found", foreignSource, PAYLOAD, PRINCIPAL);
+ }
+
+ @Test
+ void pendingAndRunningSourcesUseStableActiveConflict() {
+ UUID pendingSource = insertRootSource("PENDING", PRINCIPAL, PAYLOAD);
+ assertReplayError("etl_job_replay_source_active", pendingSource, PAYLOAD, PRINCIPAL);
+
+ UUID runningSource = insertRootSource("RUNNING", PRINCIPAL, PAYLOAD);
+ assertReplayError("etl_job_replay_source_active", runningSource, PAYLOAD, PRINCIPAL);
+ }
+
+ @Test
+ void succeededSourceUsesStableSucceededConflict() {
+ UUID succeededSource = insertRootSource("SUCCEEDED", PRINCIPAL, PAYLOAD);
+
+ assertReplayError(
+ "etl_job_replay_source_succeeded",
+ succeededSource,
+ PAYLOAD,
+ PRINCIPAL
+ );
+ }
+
+ @Test
+ void failedSourceWithDifferentResuppliedPayloadUsesStableMismatch() {
+ UUID failedSource = insertRootSource("FAILED", PRINCIPAL, PAYLOAD);
+
+ assertReplayError(
+ "etl_job_replay_payload_mismatch",
+ failedSource,
+ OTHER_PAYLOAD,
+ PRINCIPAL
+ );
+ }
+
+ private UUID insertRootSource(String status, String principal, String digestPayload) {
+ UUID sourceJobRecordId = UUID.randomUUID();
+ jdbcTemplate.update(
+ """
+ INSERT INTO etl_job_records (
+ job_record_id, principal_scope_hash, submission_key_hash,
+ request_digest, request_payload, job_status, attempt_count
+ ) VALUES (?, ?, ?, ?, NULL, ?, 1)
+ """,
+ sourceJobRecordId,
+ Sha256Digest.digest(principal),
+ Sha256Digest.digest("submission-" + sourceJobRecordId),
+ Sha256Digest.digest(digestPayload),
+ status
+ );
+ return sourceJobRecordId;
+ }
+
+ private void assertReplayError(
+ String expectedCode,
+ UUID sourceJobRecordId,
+ String payload,
+ String principal
+ ) {
+ EtlRequestException exception = assertThrows(
+ EtlRequestException.class,
+ () -> replayService.replayOwned(
+ sourceJobRecordId,
+ payload,
+ REPLAY_KEY,
+ principal
+ )
+ );
+ assertEquals(expectedCode, exception.getMessage());
+ }
+
+ /** Minimal transaction-enabled Spring context for source-classification verification. */
+ @Configuration
+ @EnableTransactionManagement
+ static class TestConfiguration {
+
+ @Bean
+ DataSource dataSource() {
+ return new EmbeddedDatabaseBuilder()
+ .generateUniqueName(true)
+ .setType(EmbeddedDatabaseType.H2)
+ .build();
+ }
+
+ @Bean
+ JdbcTemplate jdbcTemplate(DataSource dataSource) {
+ return new JdbcTemplate(dataSource);
+ }
+
+ @Bean
+ PlatformTransactionManager transactionManager(DataSource dataSource) {
+ return new DataSourceTransactionManager(dataSource);
+ }
+
+ @Bean
+ ObjectMapper objectMapper() {
+ return new ObjectMapper();
+ }
+
+ @Bean
+ EtlBatchProperties etlBatchProperties() {
+ return new EtlBatchProperties();
+ }
+
+ @Bean
+ EtlRequestLock etlRequestLock() {
+ return lockHash -> true;
+ }
+
+ @Bean
+ EtlJobReplayService replayService(
+ JdbcTemplate jdbcTemplate,
+ ObjectMapper objectMapper,
+ EtlBatchProperties batchProperties,
+ EtlRequestLock requestLock
+ ) {
+ return new EtlJobReplayService(
+ jdbcTemplate,
+ objectMapper,
+ batchProperties,
+ requestLock
+ );
+ }
+ }
+}
From e935d5ebb4c9cca7e77fff5889b759eff00d4fcb Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 22:15:12 +0900
Subject: [PATCH 27/32] feat(etl): classify replay source state safely
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 72 +++++++++++++------
.../xtrmetl/etl/service/EtlRequestError.java | 27 +++++++
2 files changed, 77 insertions(+), 22 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index 52f23847..af608c05 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -25,13 +25,11 @@
/**
* Admits owner-scoped durable-job replay requests on the repaired stack.
*
- * The service validates replay identity, authenticated principal scope, and JSON shape before
- * database work. The first persistence increments admit an owner's first-generation replay of a
- * {@code FAILED} or {@code CANCELLED} source, serialize creation by principal-scoped replay key,
- * return the already-created job for an identical owner/source/key/payload replay, and reject a
- * committed replay key that identifies a different source or payload. Later test-first increments
- * add source-state classification and generation-depth policies before the HTTP replay resource is
- * exposed.
+ * The service validates replay identity, authenticated principal scope, and the complete
+ * bounded JSON request before database work. Replay creation is serialized by principal-scoped
+ * replay identity, returns a committed identical replay when present, locks the owner-matched
+ * first-generation source before classifying its state and immutable payload digest, and inserts
+ * a fresh {@code PENDING} child without mutating terminal source evidence.
*/
@Service
public class EtlJobReplayService {
@@ -52,13 +50,11 @@ public class EtlJobReplayService {
WHERE principal_scope_hash = ?
AND submission_key_hash = ?
""";
- private static final String SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL = """
- SELECT job_record_id
+ private static final String SELECT_FIRST_GENERATION_SOURCE_SQL = """
+ SELECT job_record_id, request_digest, job_status
FROM etl_job_records
WHERE job_record_id = ?
AND principal_scope_hash = ?
- AND request_digest = ?
- AND job_status IN ('FAILED', 'CANCELLED')
AND replay_source_job_record_id IS NULL
AND replay_root_job_record_id IS NULL
AND replay_generation_count IS NULL
@@ -142,16 +138,18 @@ public EtlJobReplayService(
* The principal-scoped replay identity is serialized before table access. An identical
* committed replay is then looked up by owner/key/source/payload identity and returned without
* a second insert. A committed key bound to another source or payload is rejected with
- * {@link EtlRequestError#JOB_REPLAY_KEY_REUSED}. A new replay locks the immutable terminal
- * source before inserting a fresh {@code PENDING} child.
+ * {@link EtlRequestError#JOB_REPLAY_KEY_REUSED}. For a new replay, source selection first
+ * preserves owner-safe not-found behavior and then classifies active, succeeded, eligible
+ * terminal, and immutable-payload-mismatch outcomes with stable errors while holding the row
+ * lock.
*
- * @param sourceJobRecordId terminal durable job whose intent is being replayed
+ * @param sourceJobRecordId durable root job whose immutable intent is being replayed
* @param requestBody bounded JSON-array payload resupplied by the authenticated owner
* @param rawReplayKey principal-scoped idempotency key for this replay request
* @param principalName authenticated principal namespace
* @return newly created or previously committed replay job
* @throws NullPointerException when the source job identifier is absent
- * @throws EtlRequestException when validation fails, the key is busy, or intent conflicts
+ * @throws EtlRequestException when validation, ownership, state, payload, key, or lock fails
*/
@Transactional
public EtlJobReplay replayOwned(
@@ -205,13 +203,35 @@ public EtlJobReplay replayOwned(
return existingReplays.getFirst();
}
- UUID lockedSourceJobRecordId = jdbcTemplate.queryForObject(
- SELECT_FIRST_GENERATION_TERMINAL_SOURCE_SQL,
- UUID.class,
+ List replaySources = jdbcTemplate.query(
+ SELECT_FIRST_GENERATION_SOURCE_SQL,
+ (resultSet, rowNumber) -> new ReplaySource(
+ resultSet.getObject("job_record_id", UUID.class),
+ resultSet.getString("request_digest"),
+ EtlJobStatus.valueOf(resultSet.getString("job_status"))
+ ),
validatedSourceJobRecordId,
- principalScopeHash,
- requestDigest
+ principalScopeHash
);
+ if (replaySources.isEmpty()) {
+ throw new EtlRequestException(EtlRequestError.JOB_NOT_FOUND);
+ }
+ ReplaySource source = replaySources.getFirst();
+ switch (source.jobStatus()) {
+ case PENDING, RUNNING -> throw new EtlRequestException(
+ EtlRequestError.JOB_REPLAY_SOURCE_ACTIVE
+ );
+ case SUCCEEDED -> throw new EtlRequestException(
+ EtlRequestError.JOB_REPLAY_SOURCE_SUCCEEDED
+ );
+ case FAILED, CANCELLED -> {
+ // These terminal outcomes are the bounded first-generation replay sources.
+ }
+ }
+ if (!requestDigest.equals(source.requestDigest())) {
+ throw new EtlRequestException(EtlRequestError.JOB_REPLAY_PAYLOAD_MISMATCH);
+ }
+
UUID replayJobRecordId = UUID.randomUUID();
jdbcTemplate.update(
INSERT_FIRST_GENERATION_REPLAY_SQL,
@@ -220,8 +240,8 @@ public EtlJobReplay replayOwned(
replayKeyHash,
requestDigest,
requestBody,
- lockedSourceJobRecordId,
- lockedSourceJobRecordId
+ source.jobRecordId(),
+ source.jobRecordId()
);
return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false);
}
@@ -276,4 +296,12 @@ private JsonNode validateReplayPayload(@Nullable String requestBody) {
}
return root;
}
+
+ private record ReplaySource(UUID jobRecordId, String requestDigest, EtlJobStatus jobStatus) {
+ private ReplaySource {
+ Objects.requireNonNull(jobRecordId, "jobRecordId must not be null");
+ Objects.requireNonNull(requestDigest, "requestDigest must not be null");
+ Objects.requireNonNull(jobStatus, "jobStatus must not be null");
+ }
+ }
}
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
index 0e554c85..45b441b3 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
@@ -140,6 +140,15 @@ public enum EtlRequestError {
"Replay requires a supported principal-scoped Idempotency-Key."
),
+ /** The resupplied payload does not match the immutable terminal source digest. */
+ JOB_REPLAY_PAYLOAD_MISMATCH(
+ HttpStatus.UNPROCESSABLE_ENTITY,
+ "etl_job_replay_payload_mismatch",
+ "urn:mightyetl:problem:etl-job-replay-payload-mismatch",
+ "ETL job replay payload mismatch",
+ "The replay payload does not match the immutable source job payload digest."
+ ),
+
/** A replay idempotency key already identifies a different source or payload. */
JOB_REPLAY_KEY_REUSED(
HttpStatus.UNPROCESSABLE_ENTITY,
@@ -158,6 +167,24 @@ public enum EtlRequestError {
"A durable job replay with the same principal-scoped Idempotency-Key is being created."
),
+ /** A pending or running source is still active and cannot be replayed. */
+ JOB_REPLAY_SOURCE_ACTIVE(
+ HttpStatus.CONFLICT,
+ "etl_job_replay_source_active",
+ "urn:mightyetl:problem:etl-job-replay-source-active",
+ "ETL job replay source active",
+ "Pending or running durable jobs cannot be replayed."
+ ),
+
+ /** A succeeded source is excluded to prevent silent duplicate target effects. */
+ JOB_REPLAY_SOURCE_SUCCEEDED(
+ HttpStatus.CONFLICT,
+ "etl_job_replay_source_succeeded",
+ "urn:mightyetl:problem:etl-job-replay-source-succeeded",
+ "ETL job replay source succeeded",
+ "A succeeded durable job cannot be replayed through this endpoint."
+ ),
+
/** The durable job was already cancelled with a different cancellation key. */
JOB_CANCELLATION_KEY_REUSED(
HttpStatus.UNPROCESSABLE_ENTITY,
From 93b225ee6829672e8dac79dbbde434647a730393 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Mon, 10 Aug 2026 21:22:03 +0900
Subject: [PATCH 28/32] test(etl): fail closed on unknown replay source state
---
...lJobReplayStateClassificationIntegrationTest.java | 12 ++++++++++++
1 file changed, 12 insertions(+)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
index 95f0b46e..a06bcfc6 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
@@ -105,6 +105,18 @@ void succeededSourceUsesStableSucceededConflict() {
);
}
+ @Test
+ void unrecognizedStoredSourceStatusFailsClosedWithStableUnsupportedConflict() {
+ UUID unsupportedSource = insertRootSource("PAUSED", PRINCIPAL, PAYLOAD);
+
+ assertReplayError(
+ "etl_job_replay_source_unsupported",
+ unsupportedSource,
+ PAYLOAD,
+ PRINCIPAL
+ );
+ }
+
@Test
void failedSourceWithDifferentResuppliedPayloadUsesStableMismatch() {
UUID failedSource = insertRootSource("FAILED", PRINCIPAL, PAYLOAD);
From 847ce838a80e903c59cbad50382f3aa5619d1e55 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Mon, 10 Aug 2026 21:25:28 +0900
Subject: [PATCH 29/32] fix(etl): reject unsupported replay source states
---
.../com/xtrmetl/etl/job/EtlJobReplayService.java | 16 +++++++++++++++-
.../com/xtrmetl/etl/service/EtlRequestError.java | 9 +++++++++
2 files changed, 24 insertions(+), 1 deletion(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index af608c05..d87ac569 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -208,7 +208,7 @@ public EtlJobReplay replayOwned(
(resultSet, rowNumber) -> new ReplaySource(
resultSet.getObject("job_record_id", UUID.class),
resultSet.getString("request_digest"),
- EtlJobStatus.valueOf(resultSet.getString("job_status"))
+ parseReplaySourceStatus(resultSet.getString("job_status"))
),
validatedSourceJobRecordId,
principalScopeHash
@@ -227,6 +227,9 @@ public EtlJobReplay replayOwned(
case FAILED, CANCELLED -> {
// These terminal outcomes are the bounded first-generation replay sources.
}
+ default -> throw new EtlRequestException(
+ EtlRequestError.JOB_REPLAY_SOURCE_UNSUPPORTED
+ );
}
if (!requestDigest.equals(source.requestDigest())) {
throw new EtlRequestException(EtlRequestError.JOB_REPLAY_PAYLOAD_MISMATCH);
@@ -246,6 +249,17 @@ public EtlJobReplay replayOwned(
return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false);
}
+ private static EtlJobStatus parseReplaySourceStatus(String storedJobStatus) {
+ try {
+ return EtlJobStatus.valueOf(storedJobStatus);
+ } catch (IllegalArgumentException exception) {
+ throw new EtlRequestException(
+ EtlRequestError.JOB_REPLAY_SOURCE_UNSUPPORTED,
+ exception
+ );
+ }
+ }
+
private static String validateReplayKey(@Nullable String rawReplayKey) {
if (rawReplayKey == null) {
throw new EtlRequestException(EtlRequestError.JOB_REPLAY_KEY_REQUIRED);
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
index 45b441b3..03ae5625 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlRequestError.java
@@ -185,6 +185,15 @@ public enum EtlRequestError {
"A succeeded durable job cannot be replayed through this endpoint."
),
+ /** An unrecognized or future source state is rejected until replay semantics are defined. */
+ JOB_REPLAY_SOURCE_UNSUPPORTED(
+ HttpStatus.CONFLICT,
+ "etl_job_replay_source_unsupported",
+ "urn:mightyetl:problem:etl-job-replay-source-unsupported",
+ "ETL job replay source state unsupported",
+ "The durable job source state is not recognized as replay-eligible."
+ ),
+
/** The durable job was already cancelled with a different cancellation key. */
JOB_CANCELLATION_KEY_REUSED(
HttpStatus.UNPROCESSABLE_ENTITY,
From 11472d9ad3017fb607f8b47ebb6a92cd0c4894e7 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Mon, 10 Aug 2026 21:29:22 +0900
Subject: [PATCH 30/32] fix(etl): classify persisted replay state fail closed
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 21 +++++--------------
1 file changed, 5 insertions(+), 16 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index d87ac569..0458e96d 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -208,7 +208,7 @@ public EtlJobReplay replayOwned(
(resultSet, rowNumber) -> new ReplaySource(
resultSet.getObject("job_record_id", UUID.class),
resultSet.getString("request_digest"),
- parseReplaySourceStatus(resultSet.getString("job_status"))
+ resultSet.getString("job_status")
),
validatedSourceJobRecordId,
principalScopeHash
@@ -218,13 +218,13 @@ public EtlJobReplay replayOwned(
}
ReplaySource source = replaySources.getFirst();
switch (source.jobStatus()) {
- case PENDING, RUNNING -> throw new EtlRequestException(
+ case "PENDING", "RUNNING" -> throw new EtlRequestException(
EtlRequestError.JOB_REPLAY_SOURCE_ACTIVE
);
- case SUCCEEDED -> throw new EtlRequestException(
+ case "SUCCEEDED" -> throw new EtlRequestException(
EtlRequestError.JOB_REPLAY_SOURCE_SUCCEEDED
);
- case FAILED, CANCELLED -> {
+ case "FAILED", "CANCELLED" -> {
// These terminal outcomes are the bounded first-generation replay sources.
}
default -> throw new EtlRequestException(
@@ -249,17 +249,6 @@ public EtlJobReplay replayOwned(
return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false);
}
- private static EtlJobStatus parseReplaySourceStatus(String storedJobStatus) {
- try {
- return EtlJobStatus.valueOf(storedJobStatus);
- } catch (IllegalArgumentException exception) {
- throw new EtlRequestException(
- EtlRequestError.JOB_REPLAY_SOURCE_UNSUPPORTED,
- exception
- );
- }
- }
-
private static String validateReplayKey(@Nullable String rawReplayKey) {
if (rawReplayKey == null) {
throw new EtlRequestException(EtlRequestError.JOB_REPLAY_KEY_REQUIRED);
@@ -311,7 +300,7 @@ private JsonNode validateReplayPayload(@Nullable String requestBody) {
return root;
}
- private record ReplaySource(UUID jobRecordId, String requestDigest, EtlJobStatus jobStatus) {
+ private record ReplaySource(UUID jobRecordId, String requestDigest, String jobStatus) {
private ReplaySource {
Objects.requireNonNull(jobRecordId, "jobRecordId must not be null");
Objects.requireNonNull(requestDigest, "requestDigest must not be null");
From 72f889a356a6e9afc028cb08a03275b1131dcf5d Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Mon, 10 Aug 2026 21:33:13 +0900
Subject: [PATCH 31/32] test(etl): require replay lineage generation advance
---
...tlJobReplayPersistenceIntegrationTest.java | 44 ++++++++++++++++++-
1 file changed, 43 insertions(+), 1 deletion(-)
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
index d147fa87..233956e4 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -31,7 +31,7 @@
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
- * Proves the first database-owned durable-job replay transition on the repaired stack.
+ * Proves durable-job replay persistence, idempotency, and immutable lineage on the repaired stack.
*/
@SpringJUnitConfig(EtlJobReplayPersistenceIntegrationTest.TestConfiguration.class)
class EtlJobReplayPersistenceIntegrationTest {
@@ -261,6 +261,48 @@ void createsFirstGenerationPendingReplayFromCancelledSourceWithoutMutatingCancel
assertEquals(sourceUpdatedAt, instantColumn(sourceJobRecordId, "updated_at"));
}
+ @Test
+ void replayOfReplayPreservesRootAndIncrementsGeneration() {
+ UUID rootJobRecordId = insertFailedSource("lineage-root-submission-key");
+ EtlJobReplay firstReplay = replayService.replayOwned(
+ rootJobRecordId,
+ PAYLOAD,
+ REPLAY_KEY,
+ PRINCIPAL
+ );
+ jdbcTemplate.update(
+ """
+ UPDATE etl_job_records
+ SET job_status = 'FAILED',
+ request_payload = NULL,
+ failure_code = 'etl_target_failure'
+ WHERE job_record_id = ?
+ """,
+ firstReplay.jobRecordId()
+ );
+
+ EtlJobReplay secondReplay = replayService.replayOwned(
+ firstReplay.jobRecordId(),
+ PAYLOAD,
+ OTHER_REPLAY_KEY,
+ PRINCIPAL
+ );
+
+ assertFalse(secondReplay.replayed());
+ assertEquals(EtlJobStatus.PENDING, secondReplay.jobStatus());
+ assertEquals(
+ firstReplay.jobRecordId(),
+ uuidColumn(secondReplay.jobRecordId(), "replay_source_job_record_id")
+ );
+ assertEquals(
+ rootJobRecordId,
+ uuidColumn(secondReplay.jobRecordId(), "replay_root_job_record_id")
+ );
+ assertEquals(2, integerColumn(secondReplay.jobRecordId(), "replay_generation_count"));
+ assertEquals("FAILED", textColumn(firstReplay.jobRecordId(), "job_status"));
+ assertNull(textColumn(firstReplay.jobRecordId(), "request_payload"));
+ }
+
private UUID insertFailedSource(String submissionKey) {
UUID sourceJobRecordId = UUID.randomUUID();
jdbcTemplate.update(
From 962844d4108ee930acc85ff04df1b1d3a648a894 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Tue, 11 Aug 2026 01:17:38 +0900
Subject: [PATCH 32/32] fix(etl): advance replay lineage generation
---
.../xtrmetl/etl/job/EtlJobReplayService.java | 52 ++++++++++++-------
1 file changed, 34 insertions(+), 18 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
index 0458e96d..ac57ddaf 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -28,8 +28,9 @@
* The service validates replay identity, authenticated principal scope, and the complete
* bounded JSON request before database work. Replay creation is serialized by principal-scoped
* replay identity, returns a committed identical replay when present, locks the owner-matched
- * first-generation source before classifying its state and immutable payload digest, and inserts
- * a fresh {@code PENDING} child without mutating terminal source evidence.
+ * source before classifying its state and immutable payload digest, and inserts a fresh
+ * {@code PENDING} child that preserves the first root and advances lineage by one generation
+ * without mutating terminal source evidence.
*/
@Service
public class EtlJobReplayService {
@@ -50,17 +51,15 @@ public class EtlJobReplayService {
WHERE principal_scope_hash = ?
AND submission_key_hash = ?
""";
- private static final String SELECT_FIRST_GENERATION_SOURCE_SQL = """
- SELECT job_record_id, request_digest, job_status
+ private static final String SELECT_REPLAY_SOURCE_SQL = """
+ SELECT job_record_id, request_digest, job_status,
+ replay_root_job_record_id, replay_generation_count
FROM etl_job_records
WHERE job_record_id = ?
AND principal_scope_hash = ?
- AND replay_source_job_record_id IS NULL
- AND replay_root_job_record_id IS NULL
- AND replay_generation_count IS NULL
FOR UPDATE
""";
- private static final String INSERT_FIRST_GENERATION_REPLAY_SQL = """
+ private static final String INSERT_REPLAY_SQL = """
INSERT INTO etl_job_records (
job_record_id,
principal_scope_hash,
@@ -72,7 +71,7 @@ INSERT INTO etl_job_records (
replay_source_job_record_id,
replay_root_job_record_id,
replay_generation_count
- ) VALUES (?, ?, ?, ?, ?, 'PENDING', 0, ?, ?, 1)
+ ) VALUES (?, ?, ?, ?, ?, 'PENDING', 0, ?, ?, ?)
""";
private final JdbcTemplate jdbcTemplate;
@@ -133,7 +132,7 @@ public EtlJobReplayService(
}
/**
- * Creates or returns one first-generation replay of an owner-scoped terminal durable job.
+ * Creates or returns one replay of an owner-scoped terminal durable job.
*
* The principal-scoped replay identity is serialized before table access. An identical
* committed replay is then looked up by owner/key/source/payload identity and returned without
@@ -141,9 +140,10 @@ public EtlJobReplayService(
* {@link EtlRequestError#JOB_REPLAY_KEY_REUSED}. For a new replay, source selection first
* preserves owner-safe not-found behavior and then classifies active, succeeded, eligible
* terminal, and immutable-payload-mismatch outcomes with stable errors while holding the row
- * lock.
+ * lock. A root source creates generation one; a replay source passes its persisted first root
+ * to the next child and advances the generation by exactly one.
*
- * @param sourceJobRecordId durable root job whose immutable intent is being replayed
+ * @param sourceJobRecordId durable immediate source whose immutable intent is being replayed
* @param requestBody bounded JSON-array payload resupplied by the authenticated owner
* @param rawReplayKey principal-scoped idempotency key for this replay request
* @param principalName authenticated principal namespace
@@ -204,11 +204,13 @@ public EtlJobReplay replayOwned(
}
List replaySources = jdbcTemplate.query(
- SELECT_FIRST_GENERATION_SOURCE_SQL,
+ SELECT_REPLAY_SOURCE_SQL,
(resultSet, rowNumber) -> new ReplaySource(
resultSet.getObject("job_record_id", UUID.class),
resultSet.getString("request_digest"),
- resultSet.getString("job_status")
+ resultSet.getString("job_status"),
+ resultSet.getObject("replay_root_job_record_id", UUID.class),
+ resultSet.getObject("replay_generation_count", Integer.class)
),
validatedSourceJobRecordId,
principalScopeHash
@@ -225,7 +227,7 @@ public EtlJobReplay replayOwned(
EtlRequestError.JOB_REPLAY_SOURCE_SUCCEEDED
);
case "FAILED", "CANCELLED" -> {
- // These terminal outcomes are the bounded first-generation replay sources.
+ // These terminal outcomes are the bounded replay sources.
}
default -> throw new EtlRequestException(
EtlRequestError.JOB_REPLAY_SOURCE_UNSUPPORTED
@@ -235,16 +237,24 @@ public EtlJobReplay replayOwned(
throw new EtlRequestException(EtlRequestError.JOB_REPLAY_PAYLOAD_MISMATCH);
}
+ UUID replayRootJobRecordId = source.replayRootJobRecordId() == null
+ ? source.jobRecordId()
+ : source.replayRootJobRecordId();
+ int replayGenerationCount = source.replayGenerationCount() == null
+ ? 1
+ : Math.addExact(source.replayGenerationCount(), 1);
+
UUID replayJobRecordId = UUID.randomUUID();
jdbcTemplate.update(
- INSERT_FIRST_GENERATION_REPLAY_SQL,
+ INSERT_REPLAY_SQL,
replayJobRecordId,
principalScopeHash,
replayKeyHash,
requestDigest,
requestBody,
source.jobRecordId(),
- source.jobRecordId()
+ replayRootJobRecordId,
+ replayGenerationCount
);
return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false);
}
@@ -300,7 +310,13 @@ private JsonNode validateReplayPayload(@Nullable String requestBody) {
return root;
}
- private record ReplaySource(UUID jobRecordId, String requestDigest, String jobStatus) {
+ private record ReplaySource(
+ UUID jobRecordId,
+ String requestDigest,
+ String jobStatus,
+ @Nullable UUID replayRootJobRecordId,
+ @Nullable Integer replayGenerationCount
+ ) {
private ReplaySource {
Objects.requireNonNull(jobRecordId, "jobRecordId must not be null");
Objects.requireNonNull(requestDigest, "requestDigest must not be null");