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.
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.
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.
The three lowercase SHA-256 values are non-reversible persistence identifiers copied from the
+ * accepted job row. They let the worker reuse the durable response ledger without retaining or
+ * reconstructing raw authenticated principals or raw client idempotency keys.
+ *
+ * @param jobRecordId durable job identifier
+ * @param leaseClaimId unique token generated for this exact claim or reclaim
+ * @param leaseOwnerId non-sensitive process-lifetime worker identifier
+ * @param principalScopeHash SHA-256 hash of the authenticated principal namespace
+ * @param submissionKeyHash SHA-256 hash of the normalized durable submission key
+ * @param requestDigest SHA-256 digest of the exact retained request payload
+ * @param requestPayload validated JSON payload retained while the job is non-terminal
+ * @param attemptCount one-based claim attempt count after this claim was persisted
+ * @param leaseExpiresAt database-derived instant after which this claim is stale
+ */
+public record EtlJobLease(
+ UUID jobRecordId,
+ UUID leaseClaimId,
+ String leaseOwnerId,
+ String principalScopeHash,
+ String submissionKeyHash,
+ String requestDigest,
+ String requestPayload,
+ int attemptCount,
+ Instant leaseExpiresAt
+) {
+
+ private static final Pattern SAFE_LEASE_OWNER_PATTERN = Pattern.compile(
+ "[A-Za-z0-9._:-]{8,128}"
+ );
+ private static final Pattern SHA256_HEX_PATTERN = Pattern.compile("[0-9a-f]{64}");
+
+ /**
+ * Validates every field needed for exact lease fencing and deterministic execution.
+ */
+ public EtlJobLease {
+ Objects.requireNonNull(jobRecordId, "jobRecordId must not be null");
+ Objects.requireNonNull(leaseClaimId, "leaseClaimId must not be null");
+ String requiredOwnerId = Objects.requireNonNull(
+ leaseOwnerId,
+ "leaseOwnerId must not be null"
+ );
+ if (!SAFE_LEASE_OWNER_PATTERN.matcher(requiredOwnerId).matches()) {
+ throw new IllegalArgumentException(
+ "leaseOwnerId must match [A-Za-z0-9._:-]{8,128}"
+ );
+ }
+ requireSha256Hex(principalScopeHash, "principalScopeHash");
+ requireSha256Hex(submissionKeyHash, "submissionKeyHash");
+ requireSha256Hex(requestDigest, "requestDigest");
+ Objects.requireNonNull(requestPayload, "requestPayload must not be null");
+ if (attemptCount < 1) {
+ throw new IllegalArgumentException("attemptCount must be positive");
+ }
+ Objects.requireNonNull(leaseExpiresAt, "leaseExpiresAt must not be null");
+ }
+
+ private static void requireSha256Hex(String value, String fieldName) {
+ String requiredValue = Objects.requireNonNull(value, fieldName + " must not be null");
+ if (!SHA256_HEX_PATTERN.matcher(requiredValue).matches()) {
+ throw new IllegalArgumentException(
+ fieldName + " must be lowercase 64-character SHA-256 hex"
+ );
+ }
+ }
+}
From ad5683ffca1be041c0dbcb009eaebf809e457da1 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:10:02 +0900
Subject: [PATCH 10/32] feat(etl): add durable-job integrity failure
---
.../etl/job/EtlJobIntegrityException.java | 30 +++++++++++++++++++
1 file changed, 30 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIntegrityException.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIntegrityException.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIntegrityException.java
new file mode 100644
index 00000000..7aa04390
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIntegrityException.java
@@ -0,0 +1,30 @@
+package com.xtrmetl.etl.job;
+
+/**
+ * Signals that persisted durable-job execution identity no longer matches retained or ledger data.
+ *
+ *
The exception exposes only one stable machine-readable failure code and a non-sensitive
+ * message. It deliberately omits payloads, hashes, identifiers, SQL, timestamps, and stored
+ * response bodies so accidental logging does not disclose customer or operational data.
+ */
+public class EtlJobIntegrityException extends RuntimeException {
+
+ /** Stable terminal failure code for payload or response-ledger integrity mismatches. */
+ public static final String FAILURE_CODE = "etl_job_integrity_failure";
+
+ /**
+ * Creates a non-sensitive integrity failure signal.
+ */
+ public EtlJobIntegrityException() {
+ super("Durable ETL job execution identity failed integrity validation");
+ }
+
+ /**
+ * Returns the stable machine-readable terminal failure classification.
+ *
+ * @return {@value #FAILURE_CODE}
+ */
+ public String failureCode() {
+ return FAILURE_CODE;
+ }
+}
From 5c6f87b3154c73867e740daaae501db58ab9e8ba Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:10:22 +0900
Subject: [PATCH 11/32] feat(etl): add stale durable-job lease signal
---
.../etl/job/StaleEtlJobLeaseException.java | 17 +++++++++++++++++
1 file changed, 17 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/StaleEtlJobLeaseException.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/StaleEtlJobLeaseException.java b/etl-service/src/main/java/com/xtrmetl/etl/job/StaleEtlJobLeaseException.java
new file mode 100644
index 00000000..d7320870
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/StaleEtlJobLeaseException.java
@@ -0,0 +1,17 @@
+package com.xtrmetl.etl.job;
+
+/**
+ * Signals that a worker no longer owns the exact live lease required for a state transition.
+ *
+ *
The exception intentionally carries no job, claim, owner, payload, SQL, or timestamp values so
+ * accidental logging cannot disclose operational identifiers or retained customer data.
The worker is disabled unless an operator explicitly enables it. Polling delays and lease
+ * durations are capped at one day so a malformed environment value cannot create effectively
+ * permanent scheduling gaps, arithmetic overflow, or a lease that prevents timely crash recovery.
+ * A process-lifetime lease owner identifier is generated when no external value is supplied. The
+ * identifier is deliberately restricted to a short safe ASCII profile because it is persisted as
+ * operational metadata and must never become a free-form log or database injection surface.
+ */
+@ConfigurationProperties(prefix = "xtrmetl.etl.jobs.worker")
+public class EtlJobWorkerProperties {
+
+ /** Maximum supported fixed or initial scheduler delay: one day in milliseconds. */
+ public static final long MAXIMUM_SCHEDULER_DELAY_MILLISECONDS = 86_400_000L;
+
+ /** Maximum supported durable-job lease duration: one day in seconds. */
+ public static final long MAXIMUM_LEASE_DURATION_SECONDS = 86_400L;
+
+ private static final Pattern SAFE_LEASE_OWNER_PATTERN = Pattern.compile(
+ "[A-Za-z0-9._:-]{8,128}"
+ );
+
+ private boolean enabled;
+ private long fixedDelayMilliseconds = 5_000L;
+ private long initialDelayMilliseconds = 5_000L;
+ private long leaseDurationSeconds = 300L;
+ private int maxAttempts = 3;
+ private String leaseOwnerId = "worker-" + UUID.randomUUID();
+
+ /**
+ * Creates disabled worker configuration with bounded production-safe defaults.
+ */
+ public EtlJobWorkerProperties() {
+ // Spring Boot binds through the public setters while preserving generated defaults.
+ }
+
+ /**
+ * Reports whether scheduled durable-job execution is explicitly enabled.
+ *
+ * @return {@code true} only when an operator enabled the worker
+ */
+ public boolean isEnabled() {
+ return enabled;
+ }
+
+ /**
+ * Enables or disables scheduled durable-job execution.
+ *
+ * @param enabled whether the worker should run
+ */
+ public void setEnabled(boolean enabled) {
+ this.enabled = enabled;
+ }
+
+ /**
+ * Returns the delay measured after one polling invocation completes.
+ *
+ * @return fixed delay from one millisecond through one day
+ */
+ public long getFixedDelayMilliseconds() {
+ return fixedDelayMilliseconds;
+ }
+
+ /**
+ * Sets the delay measured after one polling invocation completes.
+ *
+ * @param fixedDelayMilliseconds delay from one millisecond through one day
+ * @throws IllegalArgumentException when the delay is outside the supported range
+ */
+ public void setFixedDelayMilliseconds(long fixedDelayMilliseconds) {
+ if (fixedDelayMilliseconds < 1L
+ || fixedDelayMilliseconds > MAXIMUM_SCHEDULER_DELAY_MILLISECONDS) {
+ throw new IllegalArgumentException(
+ "fixedDelayMilliseconds must be between 1 and "
+ + MAXIMUM_SCHEDULER_DELAY_MILLISECONDS
+ );
+ }
+ this.fixedDelayMilliseconds = fixedDelayMilliseconds;
+ }
+
+ /**
+ * Returns the delay before the first polling invocation after application startup.
+ *
+ * @return initial delay from zero milliseconds through one day
+ */
+ public long getInitialDelayMilliseconds() {
+ return initialDelayMilliseconds;
+ }
+
+ /**
+ * Sets the delay before the first polling invocation after application startup.
+ *
+ * @param initialDelayMilliseconds delay from zero milliseconds through one day
+ * @throws IllegalArgumentException when the delay is outside the supported range
+ */
+ public void setInitialDelayMilliseconds(long initialDelayMilliseconds) {
+ if (initialDelayMilliseconds < 0L
+ || initialDelayMilliseconds > MAXIMUM_SCHEDULER_DELAY_MILLISECONDS) {
+ throw new IllegalArgumentException(
+ "initialDelayMilliseconds must be between 0 and "
+ + MAXIMUM_SCHEDULER_DELAY_MILLISECONDS
+ );
+ }
+ this.initialDelayMilliseconds = initialDelayMilliseconds;
+ }
+
+ /**
+ * Returns how long one database claim remains valid without renewal.
+ *
+ * @return lease duration from one second through one day
+ */
+ public long getLeaseDurationSeconds() {
+ return leaseDurationSeconds;
+ }
+
+ /**
+ * Sets how long one database claim remains valid without renewal.
+ *
+ * @param leaseDurationSeconds duration from one second through one day
+ * @throws IllegalArgumentException when the duration is outside the supported range
+ */
+ public void setLeaseDurationSeconds(long leaseDurationSeconds) {
+ if (leaseDurationSeconds < 1L
+ || leaseDurationSeconds > MAXIMUM_LEASE_DURATION_SECONDS) {
+ throw new IllegalArgumentException(
+ "leaseDurationSeconds must be between 1 and "
+ + MAXIMUM_LEASE_DURATION_SECONDS
+ );
+ }
+ this.leaseDurationSeconds = leaseDurationSeconds;
+ }
+
+ /**
+ * Returns the maximum number of claims permitted before terminal failure.
+ *
+ * @return maximum attempt count from 1 through 100
+ */
+ public int getMaxAttempts() {
+ return maxAttempts;
+ }
+
+ /**
+ * Sets the maximum number of claims permitted before terminal failure.
+ *
+ * @param maxAttempts maximum attempt count from 1 through 100
+ * @throws IllegalArgumentException when the value is outside the supported range
+ */
+ public void setMaxAttempts(int maxAttempts) {
+ if (maxAttempts < 1 || maxAttempts > 100) {
+ throw new IllegalArgumentException("maxAttempts must be between 1 and 100");
+ }
+ this.maxAttempts = maxAttempts;
+ }
+
+ /**
+ * Returns the non-sensitive process identifier persisted on active leases.
+ *
+ * @return safe process-lifetime lease owner identifier
+ */
+ public String getLeaseOwnerId() {
+ return leaseOwnerId;
+ }
+
+ /**
+ * Sets the non-sensitive process identifier persisted on active leases.
+ *
+ * @param leaseOwnerId 8-to-128-character safe ASCII process identifier
+ * @throws NullPointerException when the identifier is {@code null}
+ * @throws IllegalArgumentException when the identifier is too short, too long, or unsafe
+ */
+ public void setLeaseOwnerId(String leaseOwnerId) {
+ String requiredOwnerId = Objects.requireNonNull(
+ leaseOwnerId,
+ "leaseOwnerId must not be null"
+ );
+ if (!SAFE_LEASE_OWNER_PATTERN.matcher(requiredOwnerId).matches()) {
+ throw new IllegalArgumentException(
+ "leaseOwnerId must match [A-Za-z0-9._:-]{8,128}"
+ );
+ }
+ this.leaseOwnerId = requiredOwnerId;
+ }
+}
From ba83462851c6256a2e4898dcc801e6ee08d38a75 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:11:36 +0900
Subject: [PATCH 13/32] feat(etl): add fenced durable-job claim repository
---
.../etl/job/EtlJobLeaseRepository.java | 387 ++++++++++++++++++
1 file changed, 387 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobLeaseRepository.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobLeaseRepository.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobLeaseRepository.java
new file mode 100644
index 00000000..f88f50b6
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobLeaseRepository.java
@@ -0,0 +1,387 @@
+package com.xtrmetl.etl.job;
+
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.stereotype.Repository;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.support.TransactionSynchronizationManager;
+import org.springframework.transaction.support.TransactionTemplate;
+
+import java.time.Duration;
+import java.time.Instant;
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.UUID;
+import java.util.regex.Pattern;
+
+/**
+ * Owns PostgreSQL-backed durable-job claims and exact lease-fenced state transitions.
+ *
+ *
Every claim transaction first terminalizes eligible exhausted rows, then locks at most one
+ * oldest eligible row with {@code FOR UPDATE SKIP LOCKED}, and finally writes a fresh claim token,
+ * owner, expiry, and incremented attempt count before commit. State transitions repeat the exact
+ * claim token, owner, running status, and database-time expiry predicates so stale workers cannot
+ * mutate lifecycle state. Public callers cannot create leases longer than the worker's one-day
+ * operational ceiling, even when they bypass Spring configuration binding.
+ *
+ *
Terminal success is intentionally stricter than retry or failure bookkeeping: it is accepted
+ * only inside the caller-owned transaction that also contains the target effects and durable
+ * response-ledger write. This prevents a direct repository call from publishing false success
+ * independently of the data it claims to have committed.
+ */
+@Repository
+public class EtlJobLeaseRepository {
+
+ /** Stable terminal code assigned when no additional claim is permitted. */
+ public static final String ATTEMPTS_EXHAUSTED_FAILURE_CODE =
+ "etl_worker_attempts_exhausted";
+
+ private static final Pattern SAFE_OWNER_PATTERN = Pattern.compile(
+ "[A-Za-z0-9._:-]{8,128}"
+ );
+ private static final Pattern SAFE_FAILURE_CODE_PATTERN = Pattern.compile(
+ "[a-z][a-z0-9_]{2,127}"
+ );
+
+ private static final String TERMINALIZE_EXHAUSTED_SQL = """
+ UPDATE etl_job_records
+ SET job_status = 'FAILED',
+ request_payload = NULL,
+ failure_code = ?,
+ lease_claim_id = NULL,
+ lease_owner_id = NULL,
+ lease_expires_at = NULL,
+ updated_at = CURRENT_TIMESTAMP
+ WHERE attempt_count >= ?
+ AND (
+ job_status = 'PENDING'
+ OR (
+ job_status = 'RUNNING'
+ AND lease_expires_at <= CURRENT_TIMESTAMP
+ )
+ )
+ """;
+
+ private static final String SELECT_CANDIDATE_SQL = """
+ SELECT job_record_id,
+ principal_scope_hash,
+ submission_key_hash,
+ request_digest,
+ request_payload,
+ attempt_count,
+ CURRENT_TIMESTAMP AS database_now
+ FROM etl_job_records
+ WHERE attempt_count < ?
+ AND (
+ job_status = 'PENDING'
+ OR (
+ job_status = 'RUNNING'
+ AND lease_expires_at <= CURRENT_TIMESTAMP
+ )
+ )
+ ORDER BY created_at, job_record_id
+ FETCH FIRST 1 ROW ONLY
+ FOR UPDATE SKIP LOCKED
+ """;
+
+ private static final String CLAIM_CANDIDATE_SQL = """
+ UPDATE etl_job_records
+ SET job_status = 'RUNNING',
+ attempt_count = attempt_count + 1,
+ failure_code = NULL,
+ lease_claim_id = ?,
+ lease_owner_id = ?,
+ lease_expires_at = ?,
+ updated_at = CURRENT_TIMESTAMP
+ WHERE job_record_id = ?
+ AND attempt_count = ?
+ AND (
+ job_status = 'PENDING'
+ OR (
+ job_status = 'RUNNING'
+ AND lease_expires_at <= CURRENT_TIMESTAMP
+ )
+ )
+ """;
+
+ private static final String MARK_SUCCEEDED_SQL = """
+ UPDATE etl_job_records
+ SET job_status = 'SUCCEEDED',
+ request_payload = NULL,
+ failure_code = NULL,
+ lease_claim_id = NULL,
+ lease_owner_id = NULL,
+ lease_expires_at = NULL,
+ updated_at = CURRENT_TIMESTAMP
+ WHERE job_record_id = ?
+ AND job_status = 'RUNNING'
+ AND lease_claim_id = ?
+ AND lease_owner_id = ?
+ AND lease_expires_at > CURRENT_TIMESTAMP
+ """;
+
+ private static final String RELEASE_FOR_RETRY_SQL = """
+ UPDATE etl_job_records
+ SET job_status = 'PENDING',
+ failure_code = NULL,
+ lease_claim_id = NULL,
+ lease_owner_id = NULL,
+ lease_expires_at = NULL,
+ updated_at = CURRENT_TIMESTAMP
+ WHERE job_record_id = ?
+ AND job_status = 'RUNNING'
+ AND lease_claim_id = ?
+ AND lease_owner_id = ?
+ AND lease_expires_at > CURRENT_TIMESTAMP
+ AND attempt_count < ?
+ """;
+
+ private static final String MARK_FAILED_SQL = """
+ UPDATE etl_job_records
+ SET job_status = 'FAILED',
+ request_payload = NULL,
+ failure_code = ?,
+ lease_claim_id = NULL,
+ lease_owner_id = NULL,
+ lease_expires_at = NULL,
+ updated_at = CURRENT_TIMESTAMP
+ WHERE job_record_id = ?
+ AND job_status = 'RUNNING'
+ AND lease_claim_id = ?
+ AND lease_owner_id = ?
+ AND lease_expires_at > CURRENT_TIMESTAMP
+ """;
+
+ private final JdbcTemplate jdbcTemplate;
+ private final TransactionTemplate transactionTemplate;
+
+ /**
+ * Creates lease persistence using one JDBC adapter and one transaction authority.
+ *
+ * @param jdbcTemplate JDBC operations for the durable-job table
+ * @param transactionManager transaction manager that owns row locks and claim commits
+ */
+ public EtlJobLeaseRepository(
+ JdbcTemplate jdbcTemplate,
+ PlatformTransactionManager transactionManager
+ ) {
+ this.jdbcTemplate = Objects.requireNonNull(
+ jdbcTemplate,
+ "jdbcTemplate must not be null"
+ );
+ this.transactionTemplate = new TransactionTemplate(Objects.requireNonNull(
+ transactionManager,
+ "transactionManager must not be null"
+ ));
+ }
+
+ /**
+ * Claims at most one oldest eligible job for one worker process.
+ *
+ * @param leaseOwnerId safe non-sensitive process identifier
+ * @param leaseDuration duration from one second through one day
+ * @param maxAttempts maximum permitted claim count from 1 through 100
+ * @return a fresh claim, or an empty result when no row is eligible
+ * @throws NullPointerException when an argument is {@code null}
+ * @throws IllegalArgumentException when an argument violates its bounded contract
+ * @throws IllegalStateException when a locked candidate unexpectedly cannot be claimed
+ */
+ public Optional claimNext(
+ String leaseOwnerId,
+ Duration leaseDuration,
+ int maxAttempts
+ ) {
+ String validatedOwnerId = requireSafeOwnerId(leaseOwnerId);
+ Duration validatedDuration = requirePositiveDuration(leaseDuration);
+ int validatedMaxAttempts = requireMaxAttempts(maxAttempts);
+
+ return Objects.requireNonNull(transactionTemplate.execute(transactionStatus -> {
+ jdbcTemplate.update(
+ TERMINALIZE_EXHAUSTED_SQL,
+ ATTEMPTS_EXHAUSTED_FAILURE_CODE,
+ validatedMaxAttempts
+ );
+ List candidates = jdbcTemplate.query(
+ SELECT_CANDIDATE_SQL,
+ (resultSet, rowNumber) -> new ClaimCandidate(
+ resultSet.getObject("job_record_id", UUID.class),
+ resultSet.getString("principal_scope_hash"),
+ resultSet.getString("submission_key_hash"),
+ resultSet.getString("request_digest"),
+ resultSet.getString("request_payload"),
+ resultSet.getInt("attempt_count"),
+ resultSet.getObject("database_now", OffsetDateTime.class).toInstant()
+ ),
+ validatedMaxAttempts
+ );
+ if (candidates.isEmpty()) {
+ return Optional.empty();
+ }
+
+ ClaimCandidate candidate = candidates.getFirst();
+ UUID leaseClaimId = UUID.randomUUID();
+ Instant leaseExpiresAt = candidate.databaseNow().plus(validatedDuration);
+ int updatedRows = jdbcTemplate.update(
+ CLAIM_CANDIDATE_SQL,
+ leaseClaimId,
+ validatedOwnerId,
+ OffsetDateTime.ofInstant(leaseExpiresAt, ZoneOffset.UTC),
+ candidate.jobRecordId(),
+ candidate.attemptCount()
+ );
+ if (updatedRows != 1) {
+ throw new IllegalStateException("Locked ETL job candidate could not be claimed");
+ }
+ return Optional.of(new EtlJobLease(
+ candidate.jobRecordId(),
+ leaseClaimId,
+ validatedOwnerId,
+ candidate.principalScopeHash(),
+ candidate.submissionKeyHash(),
+ candidate.requestDigest(),
+ candidate.requestPayload(),
+ candidate.attemptCount() + 1,
+ leaseExpiresAt
+ ));
+ }), "claim transaction must return a result");
+ }
+
+ /**
+ * Commits terminal success only for the exact live claim in the atomic execution transaction.
+ *
+ * @param lease exact claim whose target effects completed in the same transaction
+ * @throws NullPointerException when the lease is {@code null}
+ * @throws IllegalStateException when no actual Spring transaction is active
+ * @throws StaleEtlJobLeaseException when the claim is expired or no longer authoritative
+ */
+ public void markSucceeded(EtlJobLease lease) {
+ EtlJobLease requiredLease = Objects.requireNonNull(lease, "lease must not be null");
+ requireActiveSuccessTransaction();
+ requireTransition(jdbcTemplate.update(
+ MARK_SUCCEEDED_SQL,
+ requiredLease.jobRecordId(),
+ requiredLease.leaseClaimId(),
+ requiredLease.leaseOwnerId()
+ ));
+ }
+
+ /**
+ * Returns a failed execution to pending only while attempts remain and the claim is exact.
+ *
+ * @param lease exact live claim to release
+ * @param maxAttempts maximum permitted claim count from 1 through 100
+ * @throws NullPointerException when the lease is {@code null}
+ * @throws IllegalArgumentException when the maximum is outside the supported range
+ * @throws StaleEtlJobLeaseException when the claim is stale or no retry remains
+ */
+ public void releaseForRetry(EtlJobLease lease, int maxAttempts) {
+ EtlJobLease requiredLease = Objects.requireNonNull(lease, "lease must not be null");
+ int validatedMaxAttempts = requireMaxAttempts(maxAttempts);
+ requireTransition(jdbcTemplate.update(
+ RELEASE_FOR_RETRY_SQL,
+ requiredLease.jobRecordId(),
+ requiredLease.leaseClaimId(),
+ requiredLease.leaseOwnerId(),
+ validatedMaxAttempts
+ ));
+ }
+
+ /**
+ * Commits terminal failure and clears the retained payload for the exact live claim.
+ *
+ * @param lease exact live claim to fail
+ * @param failureCode stable non-sensitive machine-readable failure classification
+ * @throws NullPointerException when an argument is {@code null}
+ * @throws IllegalArgumentException when the failure code is unsafe
+ * @throws StaleEtlJobLeaseException when the claim is expired or no longer authoritative
+ */
+ public void markFailed(EtlJobLease lease, String failureCode) {
+ EtlJobLease requiredLease = Objects.requireNonNull(lease, "lease must not be null");
+ String validatedFailureCode = requireSafeFailureCode(failureCode);
+ requireTransition(jdbcTemplate.update(
+ MARK_FAILED_SQL,
+ validatedFailureCode,
+ requiredLease.jobRecordId(),
+ requiredLease.leaseClaimId(),
+ requiredLease.leaseOwnerId()
+ ));
+ }
+
+ private static String requireSafeOwnerId(String leaseOwnerId) {
+ String requiredOwnerId = Objects.requireNonNull(
+ leaseOwnerId,
+ "leaseOwnerId must not be null"
+ );
+ if (!SAFE_OWNER_PATTERN.matcher(requiredOwnerId).matches()) {
+ throw new IllegalArgumentException(
+ "leaseOwnerId must match [A-Za-z0-9._:-]{8,128}"
+ );
+ }
+ return requiredOwnerId;
+ }
+
+ private static Duration requirePositiveDuration(Duration leaseDuration) {
+ Duration requiredDuration = Objects.requireNonNull(
+ leaseDuration,
+ "leaseDuration must not be null"
+ );
+ Duration maximumDuration = Duration.ofSeconds(
+ EtlJobWorkerProperties.MAXIMUM_LEASE_DURATION_SECONDS
+ );
+ if (requiredDuration.isZero()
+ || requiredDuration.isNegative()
+ || requiredDuration.compareTo(maximumDuration) > 0) {
+ throw new IllegalArgumentException(
+ "leaseDuration must be between one second and one day"
+ );
+ }
+ return requiredDuration;
+ }
+
+ private static int requireMaxAttempts(int maxAttempts) {
+ if (maxAttempts < 1 || maxAttempts > 100) {
+ throw new IllegalArgumentException("maxAttempts must be between 1 and 100");
+ }
+ return maxAttempts;
+ }
+
+ private static String requireSafeFailureCode(String failureCode) {
+ String requiredFailureCode = Objects.requireNonNull(
+ failureCode,
+ "failureCode must not be null"
+ );
+ if (!SAFE_FAILURE_CODE_PATTERN.matcher(requiredFailureCode).matches()) {
+ throw new IllegalArgumentException(
+ "failureCode must match [a-z][a-z0-9_]{2,127}"
+ );
+ }
+ return requiredFailureCode;
+ }
+
+ private static void requireActiveSuccessTransaction() {
+ if (!TransactionSynchronizationManager.isActualTransactionActive()) {
+ throw new IllegalStateException(
+ "Durable ETL success requires an active transaction"
+ );
+ }
+ }
+
+ private static void requireTransition(int updatedRows) {
+ if (updatedRows != 1) {
+ throw new StaleEtlJobLeaseException();
+ }
+ }
+
+ private record ClaimCandidate(
+ UUID jobRecordId,
+ String principalScopeHash,
+ String submissionKeyHash,
+ String requestDigest,
+ String requestPayload,
+ int attemptCount,
+ Instant databaseNow
+ ) {
+ }
+}
From 017c65ca2cfb0ee3b47560ebc7861e065b01030e Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:12:06 +0900
Subject: [PATCH 14/32] feat(etl): reuse durable idempotency ledger for jobs
---
.../etl/job/EtlJobIdempotencyService.java | 148 ++++++++++++++++++
1 file changed, 148 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIdempotencyService.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIdempotencyService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIdempotencyService.java
new file mode 100644
index 00000000..cc9d3965
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobIdempotencyService.java
@@ -0,0 +1,148 @@
+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.springframework.dao.CannotAcquireLockException;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.support.TransactionSynchronizationManager;
+
+import java.util.List;
+import java.util.Objects;
+
+/**
+ * Reuses the durable ETL response ledger with hashed job identity only.
+ *
+ *
The accepted job stores independent hashes of its authenticated principal and normalized
+ * submission key. This service domain-separates and hashes those values into one response-ledger
+ * key, verifies the exact retained payload digest, serializes execution with the existing
+ * transaction-level request lock, replays an existing matching response, or writes the target and
+ * response ledger in the caller-owned lease-fenced transaction. Raw principals and client keys are
+ * neither required nor reconstructed.
+ *
+ *
One invocation represents exactly one persisted durable attempt. This service deliberately
+ * creates no transaction and performs no in-process retry. It joins the transaction owned by
+ * {@link EtlJobExecutionService}, calls the non-retrying ETL entry point once, and lets transient
+ * failures escape to {@link EtlJobWorker}, which owns bounded retry accounting in
+ * {@code attempt_count}. Direct Spring-proxy invocation without an existing transaction fails
+ * before lock or JDBC access.
+ */
+@Service
+public class EtlJobIdempotencyService {
+
+ private static final String LEDGER_KEY_DOMAIN = "mightyetl:durable-job:v1:";
+ private static final String SELECT_LEDGER_SQL = """
+ SELECT request_digest, response_body
+ FROM etl_idempotency_records
+ WHERE idempotency_key_hash = ?
+ """;
+ private static final String INSERT_LEDGER_SQL = """
+ INSERT INTO etl_idempotency_records (
+ idempotency_key_hash,
+ request_digest,
+ response_body
+ ) VALUES (?, ?, ?)
+ """;
+
+ private final JdbcTemplate jdbcTemplate;
+ private final EtlService etlService;
+ private final EtlRequestLock requestLock;
+
+ /**
+ * Creates the hashed durable-job response-ledger adapter.
+ *
+ * @param jdbcTemplate parameterized response-ledger database access
+ * @param etlService validated ETL target writer
+ * @param requestLock transaction-lifetime response-ledger lock
+ */
+ public EtlJobIdempotencyService(
+ JdbcTemplate jdbcTemplate,
+ EtlService etlService,
+ EtlRequestLock requestLock
+ ) {
+ this.jdbcTemplate = Objects.requireNonNull(
+ jdbcTemplate,
+ "jdbcTemplate must not be null"
+ );
+ this.etlService = Objects.requireNonNull(etlService, "etlService must not be null");
+ this.requestLock = Objects.requireNonNull(requestLock, "requestLock must not be null");
+ }
+
+ /**
+ * Executes or replays one durable job inside the caller's exact lease-fenced transaction.
+ *
+ *
The method has neither transaction-creation nor retry advice. The caller must establish
+ * the transaction that also contains terminal success fencing; otherwise execution fails
+ * before any request lock, target write, or response-ledger access. This prevents a direct
+ * proxy caller from committing durable effects without the lease-success predicate.
+ *
+ * @param lease exact live claim carrying hashed execution identity and retained payload
+ * @return newly generated or replayed stable response body
+ * @throws NullPointerException when the lease is {@code null}
+ * @throws IllegalStateException when invoked without an actual transaction
+ * @throws EtlJobIntegrityException when retained payload or ledger identity conflicts
+ * @throws CannotAcquireLockException when another transaction owns the execution ledger key
+ */
+ public String process(EtlJobLease lease) {
+ EtlJobLease requiredLease = Objects.requireNonNull(lease, "lease must not be null");
+ requireActiveTransaction();
+ if (!Sha256Digest.digest(requiredLease.requestPayload()).equals(
+ requiredLease.requestDigest()
+ )) {
+ throw new EtlJobIntegrityException();
+ }
+
+ String ledgerKeyHash = Sha256Digest.digest(
+ LEDGER_KEY_DOMAIN
+ + requiredLease.principalScopeHash()
+ + ':'
+ + requiredLease.submissionKeyHash()
+ );
+ if (!requestLock.tryLock(ledgerKeyHash)) {
+ throw new CannotAcquireLockException("Durable ETL execution ledger is busy");
+ }
+
+ List storedResponses = jdbcTemplate.query(
+ SELECT_LEDGER_SQL,
+ (resultSet, rowNumber) -> new StoredResponse(
+ resultSet.getString("request_digest"),
+ resultSet.getString("response_body")
+ ),
+ ledgerKeyHash
+ );
+ if (!storedResponses.isEmpty()) {
+ StoredResponse storedResponse = storedResponses.getFirst();
+ if (!storedResponse.requestDigest().equals(requiredLease.requestDigest())) {
+ throw new EtlJobIntegrityException();
+ }
+ return storedResponse.responseBody();
+ }
+
+ String responseBody = etlService.processDataInExistingTransaction(
+ requiredLease.requestPayload()
+ );
+ jdbcTemplate.update(
+ INSERT_LEDGER_SQL,
+ ledgerKeyHash,
+ requiredLease.requestDigest(),
+ responseBody
+ );
+ return responseBody;
+ }
+
+ private static void requireActiveTransaction() {
+ if (!TransactionSynchronizationManager.isActualTransactionActive()) {
+ throw new IllegalStateException(
+ "Durable ETL job execution requires an active transaction"
+ );
+ }
+ }
+
+ private record StoredResponse(String requestDigest, String responseBody) {
+ private StoredResponse {
+ Objects.requireNonNull(requestDigest, "requestDigest must not be null");
+ Objects.requireNonNull(responseBody, "responseBody must not be null");
+ }
+ }
+}
From 8dedad47e683a712ddcd1c028673440cb66a12a7 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:12:31 +0900
Subject: [PATCH 15/32] feat(etl): add atomic durable-job execution boundary
---
.../etl/job/EtlJobExecutionService.java | 59 +++++++++++++++++++
1 file changed, 59 insertions(+)
create mode 100644 etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobExecutionService.java
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobExecutionService.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobExecutionService.java
new file mode 100644
index 00000000..74b0192a
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobExecutionService.java
@@ -0,0 +1,59 @@
+package com.xtrmetl.etl.job;
+
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.Objects;
+
+/**
+ * Executes one claimed ETL payload and commits ledger, target, and terminal success atomically.
+ *
+ *
{@link EtlJobIdempotencyService} verifies the retained payload identity, acquires the durable
+ * execution-ledger lock, replays or writes the response ledger, and writes target rows. The
+ * subsequent conditional success transition must match the exact unexpired claim. If that
+ * transition reports a stale lease, {@link StaleEtlJobLeaseException} escapes and Spring rolls back
+ * every target and response-ledger write made by the same transaction.
+ */
+@Service
+public class EtlJobExecutionService {
+
+ private final EtlJobIdempotencyService idempotencyService;
+ private final EtlJobLeaseRepository leaseRepository;
+
+ /**
+ * Creates the atomic durable-job execution boundary.
+ *
+ * @param idempotencyService hashed response-ledger and target execution service
+ * @param leaseRepository exact lease-fenced lifecycle persistence
+ */
+ public EtlJobExecutionService(
+ EtlJobIdempotencyService idempotencyService,
+ EtlJobLeaseRepository leaseRepository
+ ) {
+ this.idempotencyService = Objects.requireNonNull(
+ idempotencyService,
+ "idempotencyService must not be null"
+ );
+ this.leaseRepository = Objects.requireNonNull(
+ leaseRepository,
+ "leaseRepository must not be null"
+ );
+ }
+
+ /**
+ * Processes or replays the retained job and marks the exact live claim successful atomically.
+ *
+ * @param lease exact database claim to execute
+ * @throws NullPointerException when the lease is {@code null}
+ * @throws EtlJobIntegrityException when persisted job or ledger identity conflicts
+ * @throws com.xtrmetl.etl.service.EtlRequestException when the retained request is invalid
+ * @throws org.springframework.dao.DataAccessException when locking or a database write fails
+ * @throws StaleEtlJobLeaseException when the claim expires or is superseded before success
+ */
+ @Transactional
+ public void execute(EtlJobLease lease) {
+ EtlJobLease requiredLease = Objects.requireNonNull(lease, "lease must not be null");
+ idempotencyService.process(requiredLease);
+ leaseRepository.markSucceeded(requiredLease);
+ }
+}
From fddaa2893027f7900851c69c61f1b31743ff4847 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:14:36 +0900
Subject: [PATCH 16/32] feat(etl): add non-retrying durable transaction entry
point
---
.../com/xtrmetl/etl/service/EtlService.java | 37 +++++++++++++++----
1 file changed, 30 insertions(+), 7 deletions(-)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlService.java b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlService.java
index 6bb32a15..28b86bb2 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/service/EtlService.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/service/EtlService.java
@@ -38,6 +38,11 @@
* batch back rather than leaving committed prefix records. The service intentionally avoids the
* JVM common pool and one-task-per-record fan-out.
*
+ *
Synchronous request entry points own Spring Retry outside their transaction boundaries. A
+ * durable worker instead uses {@link #processDataInExistingTransaction(String)} exactly once per
+ * persisted attempt so the lease, target rows, response ledger, and terminal state remain in one
+ * database transaction and retry accounting stays in the durable job record.
+ *
*
Callers may optionally use {@link #processDataIdempotently(String, String, String)}. That
* method attempts a principal-scoped transaction lock without waiting, replays a prior successful
* response, and commits the ETL rows and durable ledger entry in the same transaction.
@@ -142,8 +147,8 @@ public EtlService(
/**
* Processes one JSON-array request as a prevalidated transaction-scoped batch.
*
- *
Only transient Spring data-access failures are retried. Typed input failures and
- * deterministic target constraints fail immediately instead of repeating the same work.
+ *
Only transient Spring data-access failures are retried. Retry advice wraps the transaction
+ * advice so every synchronous request attempt receives a fresh transaction.
*
* @param data UTF-8 JSON array payload
* @return one {@code Processed: } line per record, in input order
@@ -160,6 +165,26 @@ public String processData(@Nullable String data) {
return processDataInCurrentTransaction(data);
}
+ /**
+ * Processes one durable-job payload exactly once inside the caller's existing transaction.
+ *
+ *
This method deliberately has neither {@link Retryable} nor {@link Transactional}. The
+ * durable worker owns its persisted retry count and supplies the transaction that also contains
+ * its lease-fenced terminal transition and response-ledger write. An in-process retry inside
+ * that outer transaction could reuse an already failed transaction and would not represent a
+ * new durable attempt.
+ *
+ * @param data retained UTF-8 JSON array payload
+ * @return one {@code Processed: } line per record, in input order
+ * @throws IllegalStateException when the caller did not establish an actual transaction
+ * @throws EtlRequestException when the retained request violates a deterministic contract
+ * @throws org.springframework.dao.DataAccessException when the target database rejects work
+ */
+ public String processDataInExistingTransaction(@Nullable String data) {
+ requireActiveTransaction("Durable ETL execution requires an active transaction");
+ return processDataInCurrentTransaction(data);
+ }
+
/**
* Processes or replays one principal-scoped idempotent ETL request.
*
@@ -200,7 +225,7 @@ public EtlIdempotencyResult processDataIdempotently(
if (data == null) {
throw new EtlRequestException(EtlRequestError.INVALID_JSON);
}
- requireActiveTransaction();
+ requireActiveTransaction("Idempotent ETL processing requires an active transaction");
enforcePayloadLimit(data);
String idempotencyKeyHash = sha256(
@@ -288,11 +313,9 @@ private static String validatePrincipalScope(@Nullable String principalScope) {
return principalScope;
}
- private static void requireActiveTransaction() {
+ private static void requireActiveTransaction(String failureMessage) {
if (!TransactionSynchronizationManager.isActualTransactionActive()) {
- throw new IllegalStateException(
- "Idempotent ETL processing requires an active transaction"
- );
+ throw new IllegalStateException(failureMessage);
}
}
From b81ff06326fab30057f12ee0404880a4822888e5 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:15:05 +0900
Subject: [PATCH 17/32] feat(etl): add durable-job lease fencing migration
---
.../V3__add_etl_job_lease_fencing.sql | 43 +++++++++++++++++++
1 file changed, 43 insertions(+)
create mode 100644 etl-service/src/main/resources/db/migration/V3__add_etl_job_lease_fencing.sql
diff --git a/etl-service/src/main/resources/db/migration/V3__add_etl_job_lease_fencing.sql b/etl-service/src/main/resources/db/migration/V3__add_etl_job_lease_fencing.sql
new file mode 100644
index 00000000..2469acc7
--- /dev/null
+++ b/etl-service/src/main/resources/db/migration/V3__add_etl_job_lease_fencing.sql
@@ -0,0 +1,43 @@
+ALTER TABLE etl_job_records
+ ADD COLUMN lease_claim_id UUID,
+ ADD COLUMN lease_owner_id VARCHAR(128),
+ ADD COLUMN lease_expires_at TIMESTAMPTZ;
+
+-- Repair legacy rows before enforcing the stronger failure lifecycle invariant.
+UPDATE etl_job_records
+SET failure_code = 'etl_legacy_failure'
+WHERE job_status = 'FAILED'
+ AND failure_code IS NULL;
+
+UPDATE etl_job_records
+SET failure_code = NULL
+WHERE job_status <> 'FAILED'
+ AND failure_code IS NOT NULL;
+
+ALTER TABLE etl_job_records
+ ADD CONSTRAINT etl_job_lease_lifecycle_check CHECK (
+ (
+ job_status = 'RUNNING'
+ AND lease_claim_id IS NOT NULL
+ AND lease_owner_id IS NOT NULL
+ AND lease_expires_at IS NOT NULL
+ )
+ OR
+ (
+ job_status <> 'RUNNING'
+ AND lease_claim_id IS NULL
+ AND lease_owner_id IS NULL
+ AND lease_expires_at IS NULL
+ )
+ ),
+ ADD CONSTRAINT etl_job_failure_lifecycle_check CHECK (
+ (
+ job_status = 'FAILED'
+ AND failure_code IS NOT NULL
+ )
+ OR
+ (
+ job_status <> 'FAILED'
+ AND failure_code IS NULL
+ )
+ );
From 1af9a2ec924234a7c9a55eb507936e6442a3cffa Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:15:24 +0900
Subject: [PATCH 18/32] perf(etl): index durable-job claim eligibility
---
.../V4__add_etl_job_claim_eligibility_index.sql | 11 +++++++++++
1 file changed, 11 insertions(+)
create mode 100644 etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql
diff --git a/etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql b/etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql
new file mode 100644
index 00000000..a7939f34
--- /dev/null
+++ b/etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql
@@ -0,0 +1,11 @@
+-- Support oldest-first claim selection across pending and expired-running durable jobs.
+-- CONCURRENTLY preserves inserts, updates, and deletes while PostgreSQL builds the index.
+-- The companion .sql.conf disables Flyway's per-migration transaction because PostgreSQL
+-- rejects CREATE INDEX CONCURRENTLY inside a transaction block.
+CREATE INDEX CONCURRENTLY etl_job_claim_eligibility_index
+ ON etl_job_records (
+ job_status,
+ lease_expires_at,
+ created_at,
+ job_record_id
+ );
From a7e639983a0891eaa65711faaafd639be881492a Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:15:42 +0900
Subject: [PATCH 19/32] build(etl): run claim index migration outside
transaction
---
.../migration/V4__add_etl_job_claim_eligibility_index.sql.conf | 1 +
1 file changed, 1 insertion(+)
create mode 100644 etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql.conf
diff --git a/etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql.conf b/etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql.conf
new file mode 100644
index 00000000..73bd53a1
--- /dev/null
+++ b/etl-service/src/main/resources/db/migration/V4__add_etl_job_claim_eligibility_index.sql.conf
@@ -0,0 +1 @@
+executeInTransaction=false
From 975e7d9a1666bee192808ae78ead0a184f8879cd Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:18:26 +0900
Subject: [PATCH 20/32] test(etl): fail cleanly on missing claim-index rollout
artifacts
---
.../job/EtlJobClaimIndexMigrationTest.java | 31 +++++++++++++------
1 file changed, 22 insertions(+), 9 deletions(-)
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
index a3bc6479..cb09c0fd 100644
--- a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobClaimIndexMigrationTest.java
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobClaimIndexMigrationTest.java
@@ -26,6 +26,10 @@ class EtlJobClaimIndexMigrationTest {
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";
+ private static final String APPLICATION_PROPERTIES =
+ "etl-service/src/main/resources/application.properties";
+ private static final String ROLLOUT_RUNBOOK =
+ "docs/operations/durable-job-claim-index-rollout.md";
@Test
void keepsTransactionalLeaseSchemaSeparateFromConcurrentIndexBuild() throws IOException {
@@ -48,28 +52,35 @@ void keepsTransactionalLeaseSchemaSeparateFromConcurrentIndexBuild() throws IOEx
void disablesFlywayTransactionsAndPostgresqlTransactionalLocksForConcurrentDdl()
throws IOException {
Path configurationPath = projectRoot().resolve(V4_CONFIGURATION);
- String applicationProperties = read(
- "etl-service/src/main/resources/application.properties"
- );
+ Path applicationPropertiesPath = projectRoot().resolve(APPLICATION_PROPERTIES);
assertTrue(
Files.exists(configurationPath),
"the concurrent migration requires a matching Flyway script configuration"
);
+ assertTrue(
+ Files.exists(applicationPropertiesPath),
+ "concurrent PostgreSQL Flyway DDL requires explicit non-transactional locking config"
+ );
assertTrue(
Files.readString(configurationPath, StandardCharsets.UTF_8)
.contains("executeInTransaction=false")
);
- assertTrue(applicationProperties.contains(
- "spring.flyway.postgresql.transactional-lock=false"
- ));
+ assertTrue(
+ Files.readString(applicationPropertiesPath, StandardCharsets.UTF_8).contains(
+ "spring.flyway.postgresql.transactional-lock=false"
+ )
+ );
}
@Test
void runbookDocumentsConcurrentFailureRecoveryAndRollback() throws IOException {
- String runbook = normalize(read(
- "docs/operations/durable-job-claim-index-rollout.md"
- ));
+ Path runbookPath = projectRoot().resolve(ROLLOUT_RUNBOOK);
+ assertTrue(
+ Files.exists(runbookPath),
+ "concurrent index rollout requires an operator recovery and rollback runbook"
+ );
+ String runbook = normalize(Files.readString(runbookPath, StandardCharsets.UTF_8));
assertTrue(runbook.contains("V4__add_etl_job_claim_eligibility_index.sql"));
assertTrue(runbook.contains("CREATE INDEX CONCURRENTLY"));
@@ -90,6 +101,8 @@ private static String normalize(String value) {
/**
* Finds the reactor root from either repository-root or module-local Maven execution.
+ *
+ * @return absolute repository root containing the source under test
*/
private static Path projectRoot() {
Path current = Paths.get(System.getProperty("user.dir")).toAbsolutePath();
From 04d05fc230979be58a91fe840f142891ade2e6e6 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:19:26 +0900
Subject: [PATCH 21/32] feat(etl): mirror durable worker config aliases
---
.../MightyEtlConfigAliasEnvironmentPostProcessor.java | 6 ++++++
1 file changed, 6 insertions(+)
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessor.java b/etl-service/src/main/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessor.java
index 5462db4b..16a9ad00 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessor.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/config/MightyEtlConfigAliasEnvironmentPostProcessor.java
@@ -28,6 +28,12 @@ public class MightyEtlConfigAliasEnvironmentPostProcessor implements Environment
"etl.max-payload-bytes",
"etl.max-batch-records",
"etl.jobs.intake-enabled",
+ "etl.jobs.worker.enabled",
+ "etl.jobs.worker.fixed-delay-milliseconds",
+ "etl.jobs.worker.initial-delay-milliseconds",
+ "etl.jobs.worker.lease-duration-seconds",
+ "etl.jobs.worker.max-attempts",
+ "etl.jobs.worker.lease-owner-id",
"connectors.databricks.enabled",
"connectors.snowflake.enabled",
"connectors.qlik-sense.enabled"
From 7f6854aea462fa7dd09666dc893bfdc3f145928b Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:21:02 +0900
Subject: [PATCH 22/32] build(etl): configure non-transactional Flyway
PostgreSQL lock
---
etl-service/src/main/resources/application.properties | 2 ++
1 file changed, 2 insertions(+)
create mode 100644 etl-service/src/main/resources/application.properties
diff --git a/etl-service/src/main/resources/application.properties b/etl-service/src/main/resources/application.properties
new file mode 100644
index 00000000..918dd4ad
--- /dev/null
+++ b/etl-service/src/main/resources/application.properties
@@ -0,0 +1,2 @@
+# PostgreSQL concurrent index migrations cannot use Flyway's transactional advisory lock.
+spring.flyway.postgresql.transactional-lock=false
From a9e58fbfd8a3e89cf3ab27fd4b2441c11e2c34e2 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:21:34 +0900
Subject: [PATCH 23/32] docs(etl): add concurrent claim-index rollout runbook
---
.../durable-job-claim-index-rollout.md | 95 +++++++++++++++++++
1 file changed, 95 insertions(+)
create mode 100644 docs/operations/durable-job-claim-index-rollout.md
diff --git a/docs/operations/durable-job-claim-index-rollout.md b/docs/operations/durable-job-claim-index-rollout.md
new file mode 100644
index 00000000..79adde34
--- /dev/null
+++ b/docs/operations/durable-job-claim-index-rollout.md
@@ -0,0 +1,95 @@
+# Durable-job claim index rollout
+
+## Purpose
+
+The durable worker queries `etl_job_records` for the oldest eligible `PENDING` row or expired
+`RUNNING` row. The descriptive `etl_job_claim_eligibility_index` supports that queue-like access
+path without changing the lease-fencing state machine.
+
+The table can already receive job submissions when this index is introduced. PostgreSQL's ordinary
+`CREATE INDEX` permits reads but blocks `INSERT`, `UPDATE`, and `DELETE` until the build completes.
+For a production ETL control plane that write outage is not acceptable. Migration
+`V4__add_etl_job_claim_eligibility_index.sql` therefore uses `CREATE INDEX CONCURRENTLY`.
+
+## Flyway execution boundary
+
+PostgreSQL rejects `CREATE INDEX CONCURRENTLY` inside a transaction block. The companion script
+configuration file
+`V4__add_etl_job_claim_eligibility_index.sql.conf` contains:
+
+```properties
+executeInTransaction=false
+```
+
+Flyway's PostgreSQL transactional advisory lock is also disabled with:
+
+```properties
+spring.flyway.postgresql.transactional-lock=false
+```
+
+This causes Flyway to use the PostgreSQL integration's non-transactional lock mode required for
+concurrent index DDL. The transactional V3 migration remains responsible only for lease columns,
+legacy-data repair, and lifecycle constraints. Isolating the index in V4 prevents an index-build
+failure from partially committing those schema invariants.
+
+## Deployment procedure
+
+1. Keep durable-job intake and worker execution disabled while validating the migration package.
+2. Confirm no other concurrent index build or schema migration is active on `etl_job_records`.
+3. Apply V3 and verify the lease columns and lifecycle constraints.
+4. Apply V4 and monitor `pg_stat_progress_create_index`, database I/O, and transaction latency.
+5. Verify `pg_index.indisvalid` is true for `etl_job_claim_eligibility_index`.
+6. Run the claim-selection plan against production-equivalent data and confirm the expected index is
+ available without forcing it through planner settings.
+7. Enable one worker canary only after schema history, catalog state, and application health agree.
+
+`CREATE INDEX CONCURRENTLY` performs more work and can take longer than a regular build. It preserves
+normal writes, but it still adds CPU, memory, and I/O load and allows only one concurrent index build
+per table.
+
+## Failed migration and invalid-index recovery
+
+A failed concurrent build can leave an **invalid index** in the PostgreSQL catalog. Do not mark the
+Flyway migration successful merely because an index name exists.
+
+Recovery is fail-closed:
+
+1. Keep worker execution disabled and inspect the failed Flyway record plus `pg_index.indisvalid`.
+2. Preserve database and migration logs under incident-response controls.
+3. Correct the underlying resource, permission, duplicate-build, or transaction-mode cause.
+4. Remove an unusable index without blocking normal table access:
+
+ ```sql
+ DROP INDEX CONCURRENTLY etl_job_claim_eligibility_index;
+ ```
+
+5. Use Flyway repair only after catalog inspection and operator approval, then rerun the unchanged,
+ checksum-verified migration.
+6. Reconfirm index validity and query-plan evidence before enabling workers.
+
+Do not edit an applied versioned migration or create a same-name replacement with different SQL.
+
+## Rollback boundary
+
+Application rollback does not require removing the index; an unused valid index is compatible with
+older binaries, although it adds write-maintenance overhead. Remove it only after every deployed
+worker version no longer relies on it and a controlled change has verified the performance impact:
+
+```sql
+DROP INDEX CONCURRENTLY etl_job_claim_eligibility_index;
+```
+
+`DROP INDEX CONCURRENTLY` must also run outside a transaction block. Lease columns and constraints
+require a later forward compensating migration after all compatible binaries have been removed; they
+must not be rolled back by editing V3.
+
+## Standards and primary documentation
+
+PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation: Building indexes
+concurrently*. https://www.postgresql.org/docs/18/sql-createindex.html#SQL-CREATEINDEX-CONCURRENTLY
+
+Redgate Software. (2026). *Flyway script configuration*.
+https://documentation.red-gate.com/flyway/reference/script-configuration
+
+Redgate Software. (2026). *Flyway PostgreSQL transactional lock setting*.
+https://documentation.red-gate.com/fd/flyway-postgresql-transactional-lock-setting-277579114.html
From 1b8a480cfa86649cb1a09d09cb214b70838c5622 Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:23:24 +0900
Subject: [PATCH 24/32] test(etl): cover durable worker configuration bounds
---
.../etl/job/EtlJobWorkerPropertiesTest.java | 119 ++++++++++++++++++
1 file changed, 119 insertions(+)
create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobWorkerPropertiesTest.java
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobWorkerPropertiesTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobWorkerPropertiesTest.java
new file mode 100644
index 00000000..8a5239f3
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobWorkerPropertiesTest.java
@@ -0,0 +1,119 @@
+package com.xtrmetl.etl.job;
+
+import org.junit.jupiter.api.Test;
+
+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.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Specifies fail-closed activation and bounded durable-job worker configuration.
+ */
+class EtlJobWorkerPropertiesTest {
+
+ private static final long MAXIMUM_SCHEDULER_DELAY_MILLISECONDS = 86_400_000L;
+ private static final long MAXIMUM_LEASE_DURATION_SECONDS = 86_400L;
+
+ @Test
+ void defaultsToDisabledBoundedPollingWithGeneratedSafeOwner() {
+ EtlJobWorkerProperties properties = new EtlJobWorkerProperties();
+
+ assertFalse(properties.isEnabled());
+ assertEquals(5_000L, properties.getFixedDelayMilliseconds());
+ assertEquals(5_000L, properties.getInitialDelayMilliseconds());
+ assertEquals(300L, properties.getLeaseDurationSeconds());
+ assertEquals(3, properties.getMaxAttempts());
+ assertTrue(properties.getLeaseOwnerId().matches("[A-Za-z0-9._:-]{8,128}"));
+
+ EtlJobWorkerProperties another = new EtlJobWorkerProperties();
+ assertNotEquals(properties.getLeaseOwnerId(), another.getLeaseOwnerId());
+ }
+
+ @Test
+ void acceptsEverySupportedBoundary() {
+ EtlJobWorkerProperties properties = new EtlJobWorkerProperties();
+
+ properties.setEnabled(true);
+ properties.setFixedDelayMilliseconds(1L);
+ properties.setInitialDelayMilliseconds(0L);
+ properties.setLeaseDurationSeconds(1L);
+ properties.setMaxAttempts(1);
+ properties.setLeaseOwnerId("worker-01");
+
+ assertTrue(properties.isEnabled());
+ assertEquals(1L, properties.getFixedDelayMilliseconds());
+ assertEquals(0L, properties.getInitialDelayMilliseconds());
+ assertEquals(1L, properties.getLeaseDurationSeconds());
+ assertEquals(1, properties.getMaxAttempts());
+ assertEquals("worker-01", properties.getLeaseOwnerId());
+
+ properties.setFixedDelayMilliseconds(MAXIMUM_SCHEDULER_DELAY_MILLISECONDS);
+ properties.setInitialDelayMilliseconds(MAXIMUM_SCHEDULER_DELAY_MILLISECONDS);
+ properties.setLeaseDurationSeconds(MAXIMUM_LEASE_DURATION_SECONDS);
+ properties.setMaxAttempts(100);
+ properties.setLeaseOwnerId("w".repeat(128));
+ assertEquals(MAXIMUM_SCHEDULER_DELAY_MILLISECONDS, properties.getFixedDelayMilliseconds());
+ assertEquals(MAXIMUM_SCHEDULER_DELAY_MILLISECONDS, properties.getInitialDelayMilliseconds());
+ assertEquals(MAXIMUM_LEASE_DURATION_SECONDS, properties.getLeaseDurationSeconds());
+ assertEquals(100, properties.getMaxAttempts());
+ assertEquals(128, properties.getLeaseOwnerId().length());
+ }
+
+ @Test
+ void rejectsUnsafeNumericConfiguration() {
+ EtlJobWorkerProperties properties = new EtlJobWorkerProperties();
+
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setFixedDelayMilliseconds(0L)
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setFixedDelayMilliseconds(
+ MAXIMUM_SCHEDULER_DELAY_MILLISECONDS + 1L
+ )
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setInitialDelayMilliseconds(-1L)
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setInitialDelayMilliseconds(
+ MAXIMUM_SCHEDULER_DELAY_MILLISECONDS + 1L
+ )
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setLeaseDurationSeconds(0L)
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setLeaseDurationSeconds(MAXIMUM_LEASE_DURATION_SECONDS + 1L)
+ );
+ assertThrows(IllegalArgumentException.class, () -> properties.setMaxAttempts(0));
+ assertThrows(IllegalArgumentException.class, () -> properties.setMaxAttempts(101));
+ }
+
+ @Test
+ void rejectsMissingShortLongOrUnsafeOwnerIdentifiers() {
+ EtlJobWorkerProperties properties = new EtlJobWorkerProperties();
+
+ assertThrows(NullPointerException.class, () -> properties.setLeaseOwnerId(null));
+ assertThrows(IllegalArgumentException.class, () -> properties.setLeaseOwnerId("short"));
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setLeaseOwnerId("w".repeat(129))
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setLeaseOwnerId("worker identifier")
+ );
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> properties.setLeaseOwnerId("worker/identifier")
+ );
+ }
+}
From ecb0856be477c839261719fa63aade4b1e7a554e Mon Sep 17 00:00:00 2001
From: Seongho Bae
Date: Sun, 9 Aug 2026 11:23:52 +0900
Subject: [PATCH 25/32] test(etl): cover durable lease repository validation
---
.../EtlJobLeaseRepositoryValidationTest.java | 92 +++++++++++++++++++
1 file changed, 92 insertions(+)
create mode 100644 etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseRepositoryValidationTest.java
diff --git a/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseRepositoryValidationTest.java b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseRepositoryValidationTest.java
new file mode 100644
index 00000000..b56f85ae
--- /dev/null
+++ b/etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseRepositoryValidationTest.java
@@ -0,0 +1,92 @@
+package com.xtrmetl.etl.job;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.jdbc.core.RowMapper;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.TransactionDefinition;
+import org.springframework.transaction.TransactionStatus;
+
+import java.sql.ResultSet;
+import java.time.Duration;
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.util.List;
+import java.util.UUID;
+
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+/**
+ * Proves that the lease repository rejects unsafe arguments and impossible claim transitions.
+ */
+class EtlJobLeaseRepositoryValidationTest {
+
+ @Test
+ void rejectsLeaseDurationsAboveTheOperationalSafetyCeilingBeforeDatabaseAccess() {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ PlatformTransactionManager transactionManager = mock(PlatformTransactionManager.class);
+ EtlJobLeaseRepository repository = new EtlJobLeaseRepository(
+ jdbcTemplate,
+ transactionManager
+ );
+
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> repository.claimNext(
+ "worker-alpha",
+ Duration.ofSeconds(
+ EtlJobWorkerProperties.MAXIMUM_LEASE_DURATION_SECONDS + 1L
+ ),
+ 3
+ )
+ );
+ verifyNoInteractions(jdbcTemplate, transactionManager);
+ }
+
+ @Test
+ void failsClosedWhenALockedCandidateCannotBeUpdated() throws Exception {
+ JdbcTemplate jdbcTemplate = mock(JdbcTemplate.class);
+ PlatformTransactionManager transactionManager = mock(PlatformTransactionManager.class);
+ TransactionStatus transactionStatus = mock(TransactionStatus.class);
+ ResultSet resultSet = mock(ResultSet.class);
+ UUID jobRecordId = UUID.randomUUID();
+
+ when(transactionManager.getTransaction(any(TransactionDefinition.class)))
+ .thenReturn(transactionStatus);
+ when(jdbcTemplate.update(anyString(), any(Object[].class))).thenReturn(0);
+ when(resultSet.getObject("job_record_id", UUID.class)).thenReturn(jobRecordId);
+ when(resultSet.getString("principal_scope_hash")).thenReturn("a".repeat(64));
+ when(resultSet.getString("submission_key_hash")).thenReturn("b".repeat(64));
+ when(resultSet.getString("request_digest")).thenReturn("c".repeat(64));
+ when(resultSet.getString("request_payload"))
+ .thenReturn("[{\"id\":\"record_alpha\"}]");
+ when(resultSet.getInt("attempt_count")).thenReturn(0);
+ when(resultSet.getObject("database_now", OffsetDateTime.class))
+ .thenReturn(OffsetDateTime.of(2026, 8, 5, 0, 0, 0, 0, ZoneOffset.UTC));
+ doAnswer(invocation -> {
+ @SuppressWarnings("unchecked")
+ RowMapper