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");
+ }
+}
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..9ed754b4
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -0,0 +1,254 @@
+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 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.List;
+import java.util.Objects;
+import java.util.UUID;
+import java.util.regex.Pattern;
+
+/**
+ * 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 current 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. It classifies owner-safe absence, active and successful sources, unsupported
+ * persisted states, and payload mismatch before insertion. Later test-first increments add
+ * idempotent replay lookup, concurrency, and generation-depth policies before the HTTP replay
+ * resource is exposed.
+ */
+@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 static final String REPLAY_KEY_HASH_DOMAIN = "mightyetl:durable-job-replay-key:v1:";
+ 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 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;
+ 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 reserved for the idempotency increment
+ */
+ @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");
+ }
+
+ /**
+ * Creates one first-generation replay of an owner-scoped failed or cancelled durable job.
+ *
+ * Authentication scope identifies the source row before its persisted state is classified,
+ * preserving the same not-found result for missing and foreign-owned jobs. Only failed or
+ * cancelled root jobs whose immutable request digest matches the resupplied payload proceed to
+ * insertion. The source row is locked and never mutated. Replay-key normalization remains
+ * domain-separated before reuse of the existing durable submission-key column.
+ *
+ * @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 pending replay job
+ * @throws NullPointerException when the source job identifier is absent
+ * @throws EtlRequestException when key, principal, JSON, ownership, state, or payload validation fails
+ */
+ @Transactional
+ public EtlJobReplay replayOwned(
+ UUID sourceJobRecordId,
+ @Nullable String requestBody,
+ @Nullable String rawReplayKey,
+ @Nullable String principalName
+ ) {
+ UUID validatedSourceJobRecordId = Objects.requireNonNull(
+ sourceJobRecordId,
+ "sourceJobRecordId must not be null"
+ );
+ String validatedReplayKey = validateReplayKey(rawReplayKey);
+ String validatedPrincipalName = validatePrincipalScope(principalName);
+ validateReplayPayload(requestBody);
+
+ String principalScopeHash = Sha256Digest.digest(validatedPrincipalName);
+ String requestDigest = Sha256Digest.digest(requestBody);
+ List replaySources = jdbcTemplate.query(
+ SELECT_FIRST_GENERATION_SOURCE_SQL,
+ (resultSet, rowNumber) -> new ReplaySource(
+ resultSet.getObject("job_record_id", UUID.class),
+ resultSet.getString("request_digest"),
+ resultSet.getString("job_status")
+ ),
+ validatedSourceJobRecordId,
+ 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.
+ }
+ default -> throw new EtlRequestException(
+ EtlRequestError.JOB_REPLAY_SOURCE_UNSUPPORTED
+ );
+ }
+ if (!requestDigest.equals(source.requestDigest())) {
+ throw new EtlRequestException(EtlRequestError.JOB_REPLAY_PAYLOAD_MISMATCH);
+ }
+
+ 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,
+ source.jobRecordId(),
+ source.jobRecordId()
+ );
+ return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false);
+ }
+
+ 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);
+ }
+ }
+
+ private record ReplaySource(UUID jobRecordId, String requestDigest, String jobStatus) {
+ }
+}
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..6cba3714 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,51 @@ 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 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 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."
+ ),
+
+ /** 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,
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..f766d97d
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -0,0 +1,227 @@
+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;
+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;
+import static org.mockito.Mockito.verifyNoInteractions;
+
+/**
+ * Captures immutable replay-result, construction, and pre-database 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() {
+ 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)
+ );
+ }
+
+ @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);
+ }
+
+ @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, " ", 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);
+ }
+
+ @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 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);
+ }
+
+ 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
new file mode 100644
index 00000000..04308af2
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -0,0 +1,290 @@
+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.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;
+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 OTHER_REPLAY_KEY = "519bc126-1398-4b4e-a4e3-1fb18a00f19b";
+ 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 = insertFailedSource("original-submission-key");
+ 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"));
+ }
+
+ @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")
+ );
+ }
+
+ @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(
+ """
+ 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 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=?",
+ 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
+ );
+ }
+ }
+}
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..a06bcfc6
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayStateClassificationIntegrationTest.java
@@ -0,0 +1,221 @@
+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 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);
+
+ 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
+ );
+ }
+ }
+}