From 96e1e9c39102ab0ca11811f193e4e543d03f19e4 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:11:20 +0900 Subject: [PATCH 1/8] test(etl): reproduce durable lease model contract --- .../xtrmetl/etl/job/EtlJobLeaseModelTest.java | 147 ++++++++++++++++++ 1 file changed, 147 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseModelTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseModelTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseModelTest.java new file mode 100644 index 00000000..b221d95c --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseModelTest.java @@ -0,0 +1,147 @@ +package com.xtrmetl.etl.job; + +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +/** + * Specifies the immutable value contract carried from a database claim into execution. + */ +class EtlJobLeaseModelTest { + + private static final UUID JOB_RECORD_ID = UUID.randomUUID(); + private static final UUID LEASE_CLAIM_ID = UUID.randomUUID(); + private static final String OWNER_ID = "worker-alpha"; + private static final String PRINCIPAL_SCOPE_HASH = "a".repeat(64); + private static final String SUBMISSION_KEY_HASH = "b".repeat(64); + private static final String REQUEST_DIGEST = "c".repeat(64); + private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]"; + private static final Instant EXPIRY = Instant.parse("2026-08-05T01:00:00Z"); + + @Test + void retainsEveryValidatedClaimField() { + EtlJobLease lease = validLease(); + + assertEquals(JOB_RECORD_ID, lease.jobRecordId()); + assertEquals(LEASE_CLAIM_ID, lease.leaseClaimId()); + assertEquals(OWNER_ID, lease.leaseOwnerId()); + assertEquals(PRINCIPAL_SCOPE_HASH, lease.principalScopeHash()); + assertEquals(SUBMISSION_KEY_HASH, lease.submissionKeyHash()); + assertEquals(REQUEST_DIGEST, lease.requestDigest()); + assertEquals(PAYLOAD, lease.requestPayload()); + assertEquals(2, lease.attemptCount()); + assertEquals(EXPIRY, lease.leaseExpiresAt()); + } + + @Test + void rejectsMissingUnsafeOrImpossibleFields() { + assertThrows( + NullPointerException.class, + () -> lease(null, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, null, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, null, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + IllegalArgumentException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, "unsafe owner", + PRINCIPAL_SCOPE_HASH, SUBMISSION_KEY_HASH, REQUEST_DIGEST, + PAYLOAD, 1, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, null, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + IllegalArgumentException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, "A".repeat(64), + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + null, REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + IllegalArgumentException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + "short", REQUEST_DIGEST, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, null, PAYLOAD, 1, EXPIRY) + ); + assertThrows( + IllegalArgumentException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, "g".repeat(64), PAYLOAD, 1, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, null, 1, EXPIRY) + ); + assertThrows( + IllegalArgumentException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 0, EXPIRY) + ); + assertThrows( + NullPointerException.class, + () -> lease(JOB_RECORD_ID, LEASE_CLAIM_ID, OWNER_ID, PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, REQUEST_DIGEST, PAYLOAD, 1, null) + ); + } + + private static EtlJobLease validLease() { + return lease( + JOB_RECORD_ID, + LEASE_CLAIM_ID, + OWNER_ID, + PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, + REQUEST_DIGEST, + PAYLOAD, + 2, + EXPIRY + ); + } + + private static EtlJobLease lease( + UUID jobRecordId, + UUID leaseClaimId, + String leaseOwnerId, + String principalScopeHash, + String submissionKeyHash, + String requestDigest, + String requestPayload, + int attemptCount, + Instant leaseExpiresAt + ) { + return new EtlJobLease( + jobRecordId, + leaseClaimId, + leaseOwnerId, + principalScopeHash, + submissionKeyHash, + requestDigest, + requestPayload, + attemptCount, + leaseExpiresAt + ); + } +} From 026ee7b1527b699ebc0c98d829c1fbfeba44003b Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:13:18 +0900 Subject: [PATCH 2/8] test(etl): require durable worker config aliases --- ...nfigAliasEnvironmentPostProcessorTest.java | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/etl-service/src/test/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessorTest.java b/etl-service/src/test/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessorTest.java index 285b24d7..389dda7f 100644 --- a/etl-service/src/test/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessorTest.java +++ b/etl-service/src/test/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessorTest.java @@ -45,6 +45,38 @@ void mirrorsModernDurableIntakeFlagToTheLegacyControllerCondition() { assertEquals("true", aliases.get("xtrmetl.etl.jobs.intake-enabled")); } + @Test + void mirrorsEveryModernDurableWorkerSettingToLegacyConsumers() { + MockEnvironment env = new MockEnvironment(); + env.setProperty("mightyetl.etl.jobs.worker.enabled", "true"); + env.setProperty("mightyetl.etl.jobs.worker.fixed-delay-milliseconds", "2500"); + env.setProperty("mightyetl.etl.jobs.worker.initial-delay-milliseconds", "1000"); + env.setProperty("mightyetl.etl.jobs.worker.lease-duration-seconds", "120"); + env.setProperty("mightyetl.etl.jobs.worker.max-attempts", "5"); + env.setProperty("mightyetl.etl.jobs.worker.lease-owner-id", "worker-primary"); + + Map aliases = MightyEtlConfigAliasEnvironmentPostProcessor.buildAliases(env); + + assertEquals("true", aliases.get("xtrmetl.etl.jobs.worker.enabled")); + assertEquals( + "2500", + aliases.get("xtrmetl.etl.jobs.worker.fixed-delay-milliseconds") + ); + assertEquals( + "1000", + aliases.get("xtrmetl.etl.jobs.worker.initial-delay-milliseconds") + ); + assertEquals( + "120", + aliases.get("xtrmetl.etl.jobs.worker.lease-duration-seconds") + ); + assertEquals("5", aliases.get("xtrmetl.etl.jobs.worker.max-attempts")); + assertEquals( + "worker-primary", + aliases.get("xtrmetl.etl.jobs.worker.lease-owner-id") + ); + } + @Test void mirrorsLegacyBatchLimitForModernTooling() { MockEnvironment env = new MockEnvironment(); From 68be146447958a15697b98b560086887f79f3089 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:14:33 +0900 Subject: [PATCH 3/8] test(etl): specify nonblocking claim-index rollout --- .../job/EtlJobClaimIndexMigrationTest.java | 111 ++++++++++++++++++ 1 file changed, 111 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobClaimIndexMigrationTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobClaimIndexMigrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobClaimIndexMigrationTest.java new file mode 100644 index 00000000..a3bc6479 --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobClaimIndexMigrationTest.java @@ -0,0 +1,111 @@ +package com.xtrmetl.etl.job; + +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Guards the production rollout contract for the durable-job claim eligibility index. + * + *

The table can already contain accepted work when lease fencing is introduced. PostgreSQL's + * regular index build blocks inserts, updates, and deletes, so the claim index must be isolated in + * a non-transactional concurrent migration while the lease columns and constraints remain in the + * transactional V3 migration.

+ */ +class EtlJobClaimIndexMigrationTest { + + private static final String V3_MIGRATION = + "etl-service/src/main/resources/db/migration/V3__add_etl_job_lease_fencing.sql"; + private static final String V4_MIGRATION = + "etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql"; + private static final String V4_CONFIGURATION = V4_MIGRATION + ".conf"; + + @Test + void keepsTransactionalLeaseSchemaSeparateFromConcurrentIndexBuild() throws IOException { + String leaseMigration = normalize(read(V3_MIGRATION)); + String indexMigration = normalize(read(V4_MIGRATION)); + + assertFalse( + leaseMigration.contains("CREATE INDEX"), + "the transactional lease migration must not contain a production index build" + ); + assertTrue(indexMigration.contains( + "CREATE INDEX CONCURRENTLY etl_job_claim_eligibility_index" + )); + assertTrue(indexMigration.contains( + "ON etl_job_records ( job_status, lease_expires_at, created_at, job_record_id )" + )); + } + + @Test + void disablesFlywayTransactionsAndPostgresqlTransactionalLocksForConcurrentDdl() + throws IOException { + Path configurationPath = projectRoot().resolve(V4_CONFIGURATION); + String applicationProperties = read( + "etl-service/src/main/resources/application.properties" + ); + + assertTrue( + Files.exists(configurationPath), + "the concurrent migration requires a matching Flyway script configuration" + ); + assertTrue( + Files.readString(configurationPath, StandardCharsets.UTF_8) + .contains("executeInTransaction=false") + ); + assertTrue(applicationProperties.contains( + "spring.flyway.postgresql.transactional-lock=false" + )); + } + + @Test + void runbookDocumentsConcurrentFailureRecoveryAndRollback() throws IOException { + String runbook = normalize(read( + "docs/operations/durable-job-claim-index-rollout.md" + )); + + assertTrue(runbook.contains("V4__add_etl_job_claim_eligibility_index.sql")); + assertTrue(runbook.contains("CREATE INDEX CONCURRENTLY")); + assertTrue(runbook.contains("invalid index")); + assertTrue(runbook.contains( + "DROP INDEX CONCURRENTLY etl_job_claim_eligibility_index" + )); + assertTrue(runbook.contains("executeInTransaction=false")); + } + + private static String read(String relativePath) throws IOException { + return Files.readString(projectRoot().resolve(relativePath), StandardCharsets.UTF_8); + } + + private static String normalize(String value) { + return value.replaceAll("\\s+", " ").trim(); + } + + /** + * Finds the reactor root from either repository-root or module-local Maven execution. + */ + private static Path projectRoot() { + Path current = Paths.get(System.getProperty("user.dir")).toAbsolutePath(); + Path lastPomParent = null; + while (current != null) { + if (Files.exists(current.resolve(".git"))) { + return current; + } + if (Files.exists(current.resolve("pom.xml"))) { + lastPomParent = current; + } + current = current.getParent(); + } + if (lastPomParent != null) { + return lastPomParent; + } + throw new IllegalStateException("Could not find project root"); + } +} From 1fd0f1c1a48a47a439f50cc010ad173f1184b774 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:17:45 +0900 Subject: [PATCH 4/8] test(etl): specify atomic durable execution --- ...EtlJobExecutionServiceIntegrationTest.java | 303 ++++++++++++++++++ 1 file changed, 303 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobExecutionServiceIntegrationTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobExecutionServiceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobExecutionServiceIntegrationTest.java new file mode 100644 index 00000000..9fd231e1 --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobExecutionServiceIntegrationTest.java @@ -0,0 +1,303 @@ +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.EtlService; +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.Duration; +import java.time.Instant; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +/** + * Proves that ledger, target writes, and exact-live-lease success commit or roll back together. + */ +@SpringJUnitConfig(EtlJobExecutionServiceIntegrationTest.TestConfiguration.class) +class EtlJobExecutionServiceIntegrationTest { + + private static final String OWNER_ID = "worker-alpha"; + private static final String PAYLOAD = """ + [{"id":"record_alpha","name":"accepted","email":"USER@EXAMPLE.COM"}] + """; + + private final EtlJobExecutionService executionService; + private final EtlJobLeaseRepository leaseRepository; + private final EtlJobIdempotencyService idempotencyService; + private final JdbcTemplate jdbcTemplate; + + @Autowired + EtlJobExecutionServiceIntegrationTest( + EtlJobExecutionService executionService, + EtlJobLeaseRepository leaseRepository, + EtlJobIdempotencyService idempotencyService, + JdbcTemplate jdbcTemplate + ) { + this.executionService = executionService; + this.leaseRepository = leaseRepository; + this.idempotencyService = idempotencyService; + this.jdbcTemplate = jdbcTemplate; + } + + @BeforeEach + void createTables() { + jdbcTemplate.execute("DROP TABLE IF EXISTS processed_data"); + jdbcTemplate.execute("DROP TABLE IF EXISTS etl_idempotency_records"); + 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, + created_at TIMESTAMP WITH TIME ZONE NOT NULL, + updated_at TIMESTAMP WITH TIME ZONE NOT NULL + ) + """); + jdbcTemplate.execute(""" + CREATE TABLE processed_data ( + processed_record_id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, + data VARCHAR(8192) NOT NULL + ) + """); + jdbcTemplate.execute(""" + CREATE TABLE etl_idempotency_records ( + idempotency_key_hash CHAR(64) PRIMARY KEY, + request_digest CHAR(64) NOT NULL, + response_body CLOB NOT NULL, + created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + """); + } + + @Test + void commitsLedgerTargetRowsAndTerminalSuccessInOneTransaction() { + UUID jobRecordId = insertPendingJob(); + EtlJobLease lease = leaseRepository.claimNext( + OWNER_ID, + Duration.ofMinutes(5), + 3 + ).orElseThrow(); + + executionService.execute(lease); + + assertEquals(jobRecordId, lease.jobRecordId()); + assertEquals(1, tableCount("processed_data")); + assertEquals(1, tableCount("etl_idempotency_records")); + assertEquals( + "ID:record_alpha,NAME:ACCEPTED,EMAIL:user@example.com,", + jdbcTemplate.queryForObject("SELECT data FROM processed_data", String.class) + ); + assertEquals("SUCCEEDED", jobStatus(jobRecordId)); + assertEquals(0, retainedPayloadCount(jobRecordId)); + } + + @Test + void rollsBackLedgerAndTargetRowsWhenTheClaimWasSuperseded() { + UUID jobRecordId = insertPendingJob(); + EtlJobLease lease = leaseRepository.claimNext( + OWNER_ID, + Duration.ofMinutes(5), + 3 + ).orElseThrow(); + jdbcTemplate.update( + "UPDATE etl_job_records SET lease_claim_id = ? WHERE job_record_id = ?", + UUID.randomUUID(), + jobRecordId + ); + + assertThrows(StaleEtlJobLeaseException.class, () -> executionService.execute(lease)); + + assertUncommittedDurableEffects(jobRecordId); + } + + @Test + void rollsBackLedgerAndTargetRowsWhenTheLeaseExpiredBeforeExecution() { + UUID jobRecordId = insertPendingJob(); + EtlJobLease lease = leaseRepository.claimNext( + OWNER_ID, + Duration.ofMinutes(5), + 3 + ).orElseThrow(); + jdbcTemplate.update( + "UPDATE etl_job_records SET lease_expires_at = ? WHERE job_record_id = ?", + Instant.now().minusSeconds(1), + jobRecordId + ); + + assertThrows(StaleEtlJobLeaseException.class, () -> executionService.execute(lease)); + + assertUncommittedDurableEffects(jobRecordId); + } + + @Test + void rejectsMissingCollaboratorsOrLease() { + assertThrows( + NullPointerException.class, + () -> new EtlJobExecutionService(null, leaseRepository) + ); + assertThrows( + NullPointerException.class, + () -> new EtlJobExecutionService(idempotencyService, null) + ); + assertThrows(NullPointerException.class, () -> executionService.execute(null)); + } + + private void assertUncommittedDurableEffects(UUID jobRecordId) { + assertEquals(0, tableCount("processed_data")); + assertEquals(0, tableCount("etl_idempotency_records")); + assertEquals("RUNNING", jobStatus(jobRecordId)); + assertEquals(1, retainedPayloadCount(jobRecordId)); + } + + private UUID insertPendingJob() { + UUID jobRecordId = UUID.randomUUID(); + Instant now = Instant.now(); + jdbcTemplate.update( + """ + INSERT INTO etl_job_records ( + job_record_id, principal_scope_hash, submission_key_hash, + request_digest, request_payload, job_status, attempt_count, + created_at, updated_at + ) VALUES (?, ?, ?, ?, ?, 'PENDING', 0, ?, ?) + """, + jobRecordId, + "a".repeat(64), + "b".repeat(64), + Sha256Digest.digest(PAYLOAD), + PAYLOAD, + now, + now + ); + return jobRecordId; + } + + private int tableCount(String tableName) { + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM " + tableName, + Integer.class + ); + return count == null ? 0 : count; + } + + private String jobStatus(UUID jobRecordId) { + return jdbcTemplate.queryForObject( + "SELECT job_status FROM etl_job_records WHERE job_record_id = ?", + String.class, + jobRecordId + ); + } + + private int retainedPayloadCount(UUID jobRecordId) { + Integer count = jdbcTemplate.queryForObject( + """ + SELECT COUNT(*) + FROM etl_job_records + WHERE job_record_id = ? + AND request_payload IS NOT NULL + """, + Integer.class, + jobRecordId + ); + return count == null ? 0 : count; + } + + /** + * Transaction-enabled execution context using one database for job, ledger, and target effects. + */ + @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 + EtlBatchProperties etlBatchProperties() { + return new EtlBatchProperties(); + } + + @Bean + ObjectMapper objectMapper() { + return new ObjectMapper(); + } + + @Bean + EtlRequestLock etlRequestLock() { + return idempotencyKeyHash -> true; + } + + @Bean + EtlService etlService( + JdbcTemplate jdbcTemplate, + ObjectMapper objectMapper, + EtlBatchProperties properties, + EtlRequestLock requestLock + ) { + return new EtlService(jdbcTemplate, objectMapper, properties, requestLock); + } + + @Bean + EtlJobIdempotencyService etlJobIdempotencyService( + JdbcTemplate jdbcTemplate, + EtlService etlService, + EtlRequestLock requestLock + ) { + return new EtlJobIdempotencyService(jdbcTemplate, etlService, requestLock); + } + + @Bean + EtlJobLeaseRepository etlJobLeaseRepository( + JdbcTemplate jdbcTemplate, + PlatformTransactionManager transactionManager + ) { + return new EtlJobLeaseRepository(jdbcTemplate, transactionManager); + } + + @Bean + EtlJobExecutionService etlJobExecutionService( + EtlJobIdempotencyService idempotencyService, + EtlJobLeaseRepository leaseRepository + ) { + return new EtlJobExecutionService(idempotencyService, leaseRepository); + } + } +} From d6f2b9597dc4bd18c5912fc181a5c4c3f5d80339 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:18:13 +0900 Subject: [PATCH 5/8] test(etl): guard durable retry boundary --- .../EtlJobIdempotencyRetryBoundaryTest.java | 100 ++++++++++++++++++ 1 file changed, 100 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyRetryBoundaryTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyRetryBoundaryTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyRetryBoundaryTest.java new file mode 100644 index 00000000..cc266ede --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyRetryBoundaryTest.java @@ -0,0 +1,100 @@ +package com.xtrmetl.etl.job; + +import com.xtrmetl.etl.service.EtlRequestLock; +import com.xtrmetl.etl.service.EtlService; +import com.xtrmetl.etl.service.Sha256Digest; +import org.junit.jupiter.api.Test; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.DataSourceTransactionManager; +import org.springframework.jdbc.datasource.embedded.EmbeddedDatabase; +import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseBuilder; +import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType; +import org.springframework.transaction.support.TransactionTemplate; + +import java.time.Instant; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * Guards the retry boundary between durable database attempts and synchronous request retries. + * + *

A durable worker already owns bounded, persisted retries through {@code attempt_count}. It + * must therefore invoke an ETL entry point that joins the current lease transaction exactly once, + * rather than the synchronous {@code @Retryable} API whose advice is designed to wrap a fresh + * transaction for each HTTP-request attempt.

+ */ +class EtlJobIdempotencyRetryBoundaryTest { + + private static final String PAYLOAD = "[{\"id\":\"record_alpha\"}]"; + private static final String RESPONSE = "Processed: record_alpha"; + + /** + * Proves one durable attempt uses only the non-retrying current-transaction ETL entry point. + */ + @Test + void durableAttemptDoesNotInvokeTheSynchronousRetryableEntryPoint() { + EmbeddedDatabase database = new EmbeddedDatabaseBuilder() + .generateUniqueName(true) + .setType(EmbeddedDatabaseType.H2) + .build(); + try { + JdbcTemplate jdbcTemplate = new JdbcTemplate(database); + jdbcTemplate.execute(""" + CREATE TABLE etl_idempotency_records ( + idempotency_key_hash CHAR(64) PRIMARY KEY, + request_digest CHAR(64) NOT NULL, + response_body CLOB NOT NULL, + created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + """); + + EtlService etlService = mock(EtlService.class); + EtlRequestLock requestLock = mock(EtlRequestLock.class); + when(requestLock.tryLock(anyString())).thenReturn(true); + when(etlService.processDataInExistingTransaction(PAYLOAD)).thenReturn(RESPONSE); + EtlJobIdempotencyService service = new EtlJobIdempotencyService( + jdbcTemplate, + etlService, + requestLock + ); + TransactionTemplate transactionTemplate = new TransactionTemplate( + new DataSourceTransactionManager(database) + ); + + String result = transactionTemplate.execute(status -> service.process(lease())); + + assertEquals(RESPONSE, result); + verify(etlService).processDataInExistingTransaction(PAYLOAD); + verify(etlService, never()).processData(anyString()); + assertEquals( + 1, + jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM etl_idempotency_records", + Integer.class + ) + ); + } finally { + database.shutdown(); + } + } + + private static EtlJobLease lease() { + return new EtlJobLease( + UUID.randomUUID(), + UUID.randomUUID(), + "worker-alpha", + "a".repeat(64), + "b".repeat(64), + Sha256Digest.digest(PAYLOAD), + PAYLOAD, + 1, + Instant.now().plusSeconds(300) + ); + } +} From c23a1848c299d2c2b2fa2a21aa78526fb7c97b10 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:18:57 +0900 Subject: [PATCH 6/8] test(etl): specify durable idempotency ledger behavior --- ...lJobIdempotencyServiceIntegrationTest.java | 270 ++++++++++++++++++ 1 file changed, 270 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyServiceIntegrationTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyServiceIntegrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyServiceIntegrationTest.java new file mode 100644 index 00000000..adeb0867 --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyServiceIntegrationTest.java @@ -0,0 +1,270 @@ +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.EtlService; +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.dao.CannotAcquireLockException; +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 org.springframework.transaction.annotation.Transactional; + +import javax.sql.DataSource; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.Instant; +import java.util.HexFormat; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.reset; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * Proves that durable jobs reuse the response ledger without retaining raw identity values. + */ +@SpringJUnitConfig(EtlJobIdempotencyServiceIntegrationTest.TestConfiguration.class) +class EtlJobIdempotencyServiceIntegrationTest { + + private static final String OWNER_ID = "worker-alpha"; + private static final String PRINCIPAL_SCOPE_HASH = "a".repeat(64); + private static final String SUBMISSION_KEY_HASH = "b".repeat(64); + private static final String PAYLOAD = """ + [{"id":"record_alpha","name":"accepted","email":"USER@EXAMPLE.COM"}] + """; + + private final EtlJobIdempotencyService idempotencyService; + private final EtlService etlService; + private final EtlRequestLock requestLock; + private final JdbcTemplate jdbcTemplate; + + @Autowired + EtlJobIdempotencyServiceIntegrationTest( + EtlJobIdempotencyService idempotencyService, + EtlService etlService, + EtlRequestLock requestLock, + JdbcTemplate jdbcTemplate + ) { + this.idempotencyService = idempotencyService; + this.etlService = etlService; + this.requestLock = requestLock; + this.jdbcTemplate = jdbcTemplate; + } + + @BeforeEach + void createTables() { + reset(requestLock); + when(requestLock.tryLock(anyString())).thenReturn(true); + jdbcTemplate.execute("DROP TABLE IF EXISTS processed_data"); + jdbcTemplate.execute("DROP TABLE IF EXISTS etl_idempotency_records"); + jdbcTemplate.execute(""" + CREATE TABLE processed_data ( + processed_record_id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, + data VARCHAR(8192) NOT NULL + ) + """); + jdbcTemplate.execute(""" + CREATE TABLE etl_idempotency_records ( + idempotency_key_hash CHAR(64) PRIMARY KEY, + request_digest CHAR(64) NOT NULL, + response_body CLOB NOT NULL, + created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + """); + } + + @Test + @Transactional + void writesTargetAndLedgerThenReplaysWithoutASecondTargetWrite() { + EtlJobLease lease = lease(PAYLOAD, sha256(PAYLOAD)); + + String firstResponse = idempotencyService.process(lease); + String replayedResponse = idempotencyService.process(lease); + + assertEquals("Processed: record_alpha", firstResponse); + assertEquals(firstResponse, replayedResponse); + assertEquals(1, count("processed_data")); + assertEquals(1, count("etl_idempotency_records")); + verify(requestLock, org.mockito.Mockito.times(2)).tryLock(anyString()); + } + + @Test + @Transactional + void rejectsPayloadDigestMismatchBeforeLockOrWrites() { + EtlJobLease lease = lease(PAYLOAD, "c".repeat(64)); + + EtlJobIntegrityException exception = assertThrows( + EtlJobIntegrityException.class, + () -> idempotencyService.process(lease) + ); + + assertEquals("etl_job_integrity_failure", exception.failureCode()); + verify(requestLock, never()).tryLock(anyString()); + assertEquals(0, count("processed_data")); + assertEquals(0, count("etl_idempotency_records")); + } + + @Test + @Transactional + void rejectsConflictingStoredDigestWithoutAnotherTargetWrite() { + EtlJobLease lease = lease(PAYLOAD, sha256(PAYLOAD)); + idempotencyService.process(lease); + jdbcTemplate.update( + "UPDATE etl_idempotency_records SET request_digest = ?", + "d".repeat(64) + ); + + assertThrows(EtlJobIntegrityException.class, () -> idempotencyService.process(lease)); + + assertEquals(1, count("processed_data")); + assertEquals(1, count("etl_idempotency_records")); + } + + @Test + @Transactional + void reportsBusyLedgerAsTransientWithoutWrites() { + when(requestLock.tryLock(anyString())).thenReturn(false); + EtlJobLease lease = lease(PAYLOAD, sha256(PAYLOAD)); + + assertThrows(CannotAcquireLockException.class, () -> idempotencyService.process(lease)); + + assertEquals(0, count("processed_data")); + assertEquals(0, count("etl_idempotency_records")); + } + + @Test + void failsClosedWithoutARealTransaction() { + EtlJobIdempotencyService directService = new EtlJobIdempotencyService( + jdbcTemplate, + etlService, + requestLock + ); + + assertThrows(IllegalStateException.class, () -> directService.process( + lease(PAYLOAD, sha256(PAYLOAD)) + )); + verify(requestLock, never()).tryLock(anyString()); + } + + @Test + void rejectsMissingCollaboratorsAndLease() { + assertThrows( + NullPointerException.class, + () -> new EtlJobIdempotencyService(null, etlService, requestLock) + ); + assertThrows( + NullPointerException.class, + () -> new EtlJobIdempotencyService(jdbcTemplate, null, requestLock) + ); + assertThrows( + NullPointerException.class, + () -> new EtlJobIdempotencyService(jdbcTemplate, etlService, null) + ); + assertThrows(NullPointerException.class, () -> idempotencyService.process(null)); + } + + private int count(String tableName) { + Integer count = jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM " + tableName, + Integer.class + ); + return count == null ? 0 : count; + } + + private static EtlJobLease lease(String payload, String requestDigest) { + return new EtlJobLease( + UUID.randomUUID(), + UUID.randomUUID(), + OWNER_ID, + PRINCIPAL_SCOPE_HASH, + SUBMISSION_KEY_HASH, + requestDigest, + payload, + 1, + Instant.now().plusSeconds(300) + ); + } + + private static String sha256(String value) { + try { + MessageDigest digest = MessageDigest.getInstance("SHA-256"); + return HexFormat.of().formatHex(digest.digest(value.getBytes(StandardCharsets.UTF_8))); + } catch (NoSuchAlgorithmException exception) { + throw new AssertionError(exception); + } + } + + /** Transaction-enabled service context backed by an isolated H2 database. */ + @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 + EtlBatchProperties etlBatchProperties() { + return new EtlBatchProperties(); + } + + @Bean + ObjectMapper objectMapper() { + return new ObjectMapper(); + } + + @Bean + EtlRequestLock etlRequestLock() { + return mock(EtlRequestLock.class); + } + + @Bean + EtlService etlService( + JdbcTemplate jdbcTemplate, + ObjectMapper objectMapper, + EtlBatchProperties properties, + EtlRequestLock requestLock + ) { + return new EtlService(jdbcTemplate, objectMapper, properties, requestLock); + } + + @Bean + EtlJobIdempotencyService etlJobIdempotencyService( + JdbcTemplate jdbcTemplate, + EtlService etlService, + EtlRequestLock requestLock + ) { + return new EtlJobIdempotencyService(jdbcTemplate, etlService, requestLock); + } + } +} From 75743dd2d95e0eac84c3654281625aa04eb989e1 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:19:15 +0900 Subject: [PATCH 7/8] test(etl): bind durable idempotency to caller transaction --- ...bIdempotencyTransactionAnnotationTest.java | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyTransactionAnnotationTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyTransactionAnnotationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyTransactionAnnotationTest.java new file mode 100644 index 00000000..42fe752a --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobIdempotencyTransactionAnnotationTest.java @@ -0,0 +1,36 @@ +package com.xtrmetl.etl.job; + +import org.junit.jupiter.api.Test; +import org.springframework.transaction.annotation.Transactional; + +import java.lang.reflect.Method; + +import static org.junit.jupiter.api.Assertions.assertNull; + +/** + * Guards the durable response-ledger service from creating its own execution transaction. + * + *

The lease-fenced execution service owns the transaction that contains target effects, the + * response ledger, and terminal success. A transactional annotation on the public ledger method + * would let a direct Spring-proxy caller commit target and ledger effects without the success fence.

+ */ +class EtlJobIdempotencyTransactionAnnotationTest { + + /** + * Requires the public processing method to join, rather than create, the caller transaction. + * + * @throws NoSuchMethodException when the public durable processing contract is missing + */ + @Test + void processDoesNotCreateAStandaloneTransaction() throws NoSuchMethodException { + Method processMethod = EtlJobIdempotencyService.class.getMethod( + "process", + EtlJobLease.class + ); + + assertNull( + processMethod.getAnnotation(Transactional.class), + "Durable job idempotency must not own a transaction" + ); + } +} From 778956039f6e41d9d8f264097c6ebf14cd3306f3 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Fri, 7 Aug 2026 21:20:00 +0900 Subject: [PATCH 8/8] test(etl): specify lease-fencing migration contract --- .../etl/job/EtlJobLeaseMigrationTest.java | 87 +++++++++++++++++++ 1 file changed, 87 insertions(+) create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseMigrationTest.java diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseMigrationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseMigrationTest.java new file mode 100644 index 00000000..6491616c --- /dev/null +++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseMigrationTest.java @@ -0,0 +1,87 @@ +package com.xtrmetl.etl.job; + +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Specifies the transactional Flyway contract for exact durable-job lease fencing. + */ +class EtlJobLeaseMigrationTest { + + @Test + void addsDescriptiveLeaseColumnsWithoutNonTransactionalIndexDdl() throws IOException { + String migration = readMigration(); + + assertTrue(migration.contains("ADD COLUMN lease_claim_id UUID")); + assertTrue(migration.contains("ADD COLUMN lease_owner_id VARCHAR(128)")); + assertTrue(migration.contains("ADD COLUMN lease_expires_at TIMESTAMPTZ")); + assertFalse(migration.contains("CREATE INDEX")); + assertFalse(migration.contains(" ADD COLUMN owner ")); + assertFalse(migration.contains(" ADD COLUMN lease ")); + } + + @Test + void requiresLeaseFieldsOnlyForRunningRows() throws IOException { + String migration = normalize(readMigration()); + + assertTrue(migration.contains("CONSTRAINT etl_job_lease_lifecycle_check")); + assertTrue(migration.contains( + "job_status = 'RUNNING' AND lease_claim_id IS NOT NULL AND lease_owner_id IS NOT NULL AND lease_expires_at IS NOT NULL" + )); + assertTrue(migration.contains( + "job_status <> 'RUNNING' AND lease_claim_id IS NULL AND lease_owner_id IS NULL AND lease_expires_at IS NULL" + )); + } + + @Test + void requiresFailureCodesOnlyForFailedRows() throws IOException { + String migration = normalize(readMigration()); + + assertTrue(migration.contains("CONSTRAINT etl_job_failure_lifecycle_check")); + assertTrue(migration.contains("job_status = 'FAILED' AND failure_code IS NOT NULL")); + assertTrue(migration.contains("job_status <> 'FAILED' AND failure_code IS NULL")); + } + + private static String readMigration() throws IOException { + return Files.readString( + projectRoot().resolve( + "etl-service/src/main/resources/db/migration/" + + "V3__add_etl_job_lease_fencing.sql" + ), + StandardCharsets.UTF_8 + ); + } + + private static String normalize(String value) { + return value.replaceAll("\\s+", " ").trim(); + } + + /** + * Finds the reactor root from either repository-root or module-local Maven execution. + */ + private static Path projectRoot() { + Path current = Paths.get(System.getProperty("user.dir")).toAbsolutePath(); + Path lastPomParent = null; + while (current != null) { + if (Files.exists(current.resolve(".git"))) { + return current; + } + if (Files.exists(current.resolve("pom.xml"))) { + lastPomParent = current; + } + current = current.getParent(); + } + if (lastPomParent != null) { + return lastPomParent; + } + throw new IllegalStateException("Could not find project root"); + } +}