From 121654bfa463f90e2d64222f9ca60f1750a33e0e Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 09:44:18 +0900 Subject: [PATCH 1/6] test(etl): rebuild replay boundary on cancellation repair --- .../etl/job/EtlJobReplayBoundaryTest.java | 167 ++++++++++++++++++ 1 file changed, 167 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..c0bac809 --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobReplayBoundaryTest.java @@ -0,0 +1,167 @@ +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 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()); + } +} From fd9e38da4676c2caa876810be227302f4f16ad7f Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 09:47:36 +0900 Subject: [PATCH 2/6] feat(etl): restore replay admission boundary --- .../com/xtrmetl/etl/job/EtlJobReplay.java | 29 +++ .../xtrmetl/etl/job/EtlJobReplayService.java | 223 ++++++++++++++++++ .../xtrmetl/etl/service/EtlRequestError.java | 9 + 3 files changed, 261 insertions(+) create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplay.java 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/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..e5e859fc --- /dev/null +++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobReplayService.java @@ -0,0 +1,223 @@ +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.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 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.

+ */ +@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_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 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 + 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 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 newly created pending replay job + * @throws NullPointerException when the source job identifier is absent + * @throws EtlRequestException when key, principal, or JSON 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); + UUID lockedSourceJobRecordId = jdbcTemplate.queryForObject( + SELECT_FIRST_GENERATION_TERMINAL_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) { + 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); + } + } +} 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 cd5946380dfeb39a07477b32033629ab6f64bb13 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 10:07:26 +0900 Subject: [PATCH 3/6] test(etl): restore replay persistence coverage --- ...tlJobReplayPersistenceIntegrationTest.java | 290 ++++++++++++++++++ 1 file changed, 290 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..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 + ); + } + } +} From 6161c003b679fc5e9716771e0502a4b9a805d631 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 10:11:59 +0900 Subject: [PATCH 4/6] test(etl): classify replay source states --- ...layStateClassificationIntegrationTest.java | 221 ++++++++++++++++++ 1 file changed, 221 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..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 + ); + } + } +} From e6c7acd63559c45f3a7312e7457d2aa363bb4541 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 10:20:17 +0900 Subject: [PATCH 5/6] feat(etl): classify replay source states --- .../xtrmetl/etl/job/EtlJobReplayService.java | 73 +++++++++++++------ .../xtrmetl/etl/service/EtlRequestError.java | 36 +++++++++ 2 files changed, 88 insertions(+), 21 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..9ed754b4 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,12 @@ * 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 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 { @@ -41,13 +44,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_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 @@ -128,18 +129,19 @@ public EtlJobReplayService( /** * 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. - * 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.

+ *

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 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 pending replay job * @throws NullPointerException when the source job identifier is absent - * @throws EtlRequestException when key, principal, or JSON validation fails + * @throws EtlRequestException when key, principal, JSON, ownership, state, or payload validation fails */ @Transactional public EtlJobReplay replayOwned( @@ -158,13 +160,39 @@ public EtlJobReplay replayOwned( String principalScopeHash = Sha256Digest.digest(validatedPrincipalName); String requestDigest = Sha256Digest.digest(requestBody); - 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"), + 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. + } + 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 @@ -176,8 +204,8 @@ public EtlJobReplay replayOwned( replayKeyHash, requestDigest, requestBody, - lockedSourceJobRecordId, - lockedSourceJobRecordId + source.jobRecordId(), + source.jobRecordId() ); return new EtlJobReplay(replayJobRecordId, EtlJobStatus.PENDING, false); } @@ -220,4 +248,7 @@ private JsonNode validateReplayPayload(@Nullable String requestBody) { 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 e62e89f5..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 @@ -140,6 +140,42 @@ 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 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, From 8c818aadabcc4f924b9d050ea2542ceee5fab5e6 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 10:23:35 +0900 Subject: [PATCH 6/6] test(etl): enforce bounded replay payloads --- .../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 c0bac809..f766d97d 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 rejectsAnAbsentParsedRootBeforeDatabaseWork() { JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);