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..ac57ddaf
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java
@@ -0,0 +1,326 @@
+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.nio.charset.StandardCharsets;
+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 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
+ * 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 {
+
+ 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 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
+ WHERE principal_scope_hash = ?
+ AND submission_key_hash = ?
+ """;
+ 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 = ?
+ FOR UPDATE
+ """;
+ private static final String INSERT_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, ?, ?, ?)
+ """;
+
+ 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 principal-scoped 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");
+ }
+
+ /**
+ * 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
+ * a second insert. A committed key bound to another source or payload is rejected with
+ * {@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. 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 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
+ * @return newly created or previously committed replay job
+ * @throws NullPointerException when the source job identifier is absent
+ * @throws EtlRequestException when validation, ownership, state, payload, key, or lock 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);
+ 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) -> {
+ 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
+ );
+ if (!existingReplays.isEmpty()) {
+ return existingReplays.getFirst();
+ }
+
+ List replaySources = jdbcTemplate.query(
+ SELECT_REPLAY_SOURCE_SQL,
+ (resultSet, rowNumber) -> new ReplaySource(
+ resultSet.getObject("job_record_id", UUID.class),
+ resultSet.getString("request_digest"),
+ resultSet.getString("job_status"),
+ resultSet.getObject("replay_root_job_record_id", UUID.class),
+ resultSet.getObject("replay_generation_count", Integer.class)
+ ),
+ 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 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 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_REPLAY_SQL,
+ replayJobRecordId,
+ principalScopeHash,
+ replayKeyHash,
+ requestDigest,
+ requestBody,
+ source.jobRecordId(),
+ replayRootJobRecordId,
+ replayGenerationCount
+ );
+ 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) {
+ 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 {
+ 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;
+ }
+
+ 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");
+ 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 d28aa33d..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
@@ -131,6 +131,69 @@ 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 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."
+ ),
+
+ /** 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."
+ ),
+
+ /** 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..516b88e2
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java
@@ -0,0 +1,244 @@
+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 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);
+ 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..233956e4
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayPersistenceIntegrationTest.java
@@ -0,0 +1,430 @@
+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.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;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Proves durable-job replay persistence, idempotency, and immutable lineage on the repaired stack.
+ */
+@SpringJUnitConfig(EtlJobReplayPersistenceIntegrationTest.TestConfiguration.class)
+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";
+
+ 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 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 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 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");
+
+ 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"));
+ }
+
+ @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(
+ """
+ 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
+ );
+ }
+ }
+}