diff --git a/CHANGELOG.md b/CHANGELOG.md
index 5e527f50..397c6e1d 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -10,13 +10,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed
- Pull-request CI and CycloneDX SBOM jobs now check out the literal current source head, immediately assert `git rev-parse HEAD` against `github.event.pull_request.head.sha`, and disable checkout credential persistence; generated merge revisions remain useful compatibility previews but no longer masquerade as direct exact-head source evidence.
+- The durable-job claim eligibility index now builds in a separate PostgreSQL `CREATE INDEX CONCURRENTLY` migration with Flyway non-transactional script configuration and session-level PostgreSQL migration locking, preserving normal job writes during rollout while keeping lease columns and constraints transactional.
+- Durable asynchronous ETL jobs now progress from `PENDING` through lease-fenced execution to `SUCCEEDED` or `FAILED`; PostgreSQL owns cross-replica claiming, stale workers cannot commit target or lifecycle effects, and intake and execution remain independently fail-closed.
+- Durable-worker observability now records one terminal outcome counter and one matching duration sample for every completed poll, including idle polls and database failures while persisting retry or terminal transitions.
- The hourly pull-request disposition loop now requires at least one non-author approval anchored to the exact current head SHA; stale approvals, comment-only reviews, and the mere absence of requested changes cannot authorize unattended merge.
- The hourly OpenCode workflow now scopes repository write permissions to its sole maintenance job, replaces the npm installation command with the immutable OpenCode 1.18.13 Linux release archive plus pinned SHA-256 validation, requires exactly one regular-file archive member before private-directory extraction, rejects non-regular or symbolic-link output, and uses a removable repository-local GitHub CLI credential helper instead of storing an encoded authorization header while retaining `persist-credentials: false`.
- The hourly OpenCode workflow now snapshots same-repository `develop` pull-request heads before the agent runs and uses job-scoped Actions write authority only to authorize approval-required workflow runs for an unchanged exact head; `.github/**` and `CODEOWNERS` changes remain human-authorized, and no review or merge authority is added.
- Updated existing pull-request candidates now carry their captured pre-agent head into the deterministic publisher, which rejects destructive ancestry, more than 50 agent-introduced files, and any agent-introduced `.github/**` or `CODEOWNERS` change before exposing the updated pull request or authorizing checks.
- The hourly OpenCode workflow now uses the current free NVIDIA `deepseek-ai/deepseek-v4-pro` endpoint for long-context coding and agentic tool use instead of the deprecated Qwen3 Coder free endpoint; model or endpoint rejection fails visibly without a non-NVIDIA, partner-only, or automatic fallback.
- The managed Jackson component set now uses the patched 2.21.5 BOM, closing CVE-2026-54515, CVE-2026-59889, and GHSA-mhm7-754m-9p8w while keeping core, annotations, datatype, and module artifacts aligned.
-- Durable `POST /api/etl/jobs` submissions now return RFC 9110 `202 Accepted`, a stable pending-job representation, `Location` status-monitor metadata, and explicit replay metadata without changing the synchronous `/api/etl/process` contract. The incomplete intake controller is fail-closed and requires explicit `xtrmetl.etl.jobs.intake-enabled=true` operator opt-in until worker execution and terminal payload clearing are implemented.
+- Durable `POST /api/etl/jobs` submissions now return RFC 9110 `202 Accepted`, a stable job representation, `Location` status-monitor metadata, and explicit replay metadata without changing the synchronous `/api/etl/process` contract.
- Concurrent requests using the same authenticated-principal-scoped semantic idempotency key now return immediate RFC 9457 `409 etl_idempotency_request_in_progress` responses through PostgreSQL `pg_try_advisory_xact_lock`; retries after completion still replay the committed response.
- `POST /api/etl/process` now supports optional authenticated-principal-scoped `Idempotency-Key` retries with atomic target writes, durable response replay, payload-conflict rejection, and explicit replay response metadata.
- `Idempotency-Key` now prefers the quoted RFC 9651 Structured Field String representation while retaining and normalizing the legacy raw representation to the same durable ledger key.
@@ -30,11 +33,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Added
- Permanent fail-first exact-head workflow contracts and authoritative evidence in `docs/doctoring/exact-head-source-workflow-evidence.md`, including the observed synthetic merge checkout, cross-platform source identity assertions, least-privilege boundary, stack invalidation rule, rollback prohibition, and APA 7th GitHub references.
+- A production rollout and invalid-index recovery runbook for the nonblocking durable-job claim index: `docs/operations/durable-job-claim-index-rollout.md`.
+- PostgreSQL `FOR UPDATE SKIP LOCKED` durable-job claiming, per-process and per-claim lease fencing, expiry reclaim, bounded attempts, exact-live-lease transitions, terminal payload clearing, stable failure codes, and finite-cardinality worker metrics.
+- Hashed durable execution identity and domain-separated reuse of `etl_idempotency_records`, coupling response replay or creation, target writes, and terminal `SUCCEEDED` in one transaction without retaining or reconstructing raw principals or raw client idempotency keys.
+- Deterministic migration, concurrency, expiry, exhaustion, response-replay, integrity, stale-lease rollback, privacy, configuration-boundary, and operator-recovery tests plus `docs/operations/durable-job-worker.md`.
- A separate fail-closed hourly OpenCode maintenance workflow pinned to OpenCode 1.18.13 and `nvidia/deepseek-ai/deepseek-v4-pro`, using only the existing `NVIDIA_NIM_API_KEY` through OpenCode's `NVIDIA_API_KEY` provider variable while preserving the independent review agent and deterministic merge-disposition workflow.
- Exact-head workflow-run authorization doctoring evidence for the repository-token recursion boundary, before/after SHA snapshots, policy-path exclusion, time-of-check/time-of-use validation, least privilege, test-first regression evidence, and rollback in `docs/doctoring/github-token-exact-head-check-authorization-evidence.md`.
- Supply-chain doctoring evidence for checksum binding, exact archive-member and entry-type validation, private extraction, post-extraction file checks, test-first regression evidence, and rollback in `docs/doctoring/opencode-archive-extraction-evidence.md`.
- NVIDIA model-selection doctoring evidence for endpoint availability, deprecated-endpoint rejection, capability and context evidence, no-fallback semantics, test-first regression evidence, and replacement procedure in `docs/doctoring/nvidia-opencode-model-selection-evidence.md`.
-- Principal-scoped durable asynchronous ETL job intake and owner-scoped status resources, Flyway `etl_job_records` migration, deterministic replay/conflict coverage, and the explicit worker boundary in `docs/etl/durable-job-intake.md`.
+- Principal-scoped durable asynchronous ETL job intake and owner-scoped status resources, Flyway `etl_job_records` migration, deterministic replay/conflict coverage, and the authoritative lifecycle contract in `docs/etl/durable-job-intake.md`.
- Durable idempotency ledger migration, PostgreSQL transaction advisory-lock adapter, deterministic concurrency/rollback coverage, and the operator/client contract `docs/etl/idempotent-retries.md`.
- ETL problem-details client and operator contract: `docs/api/problem-details.md`.
- Operator-configurable ETL admission limits under `mightyetl.etl.*` / `xtrmetl.etl.*`, backed by `ETL_MAX_PAYLOAD_BYTES` and `ETL_MAX_BATCH_RECORDS` environment variables with hard safety ceilings.
@@ -60,6 +67,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Root POM `mightyETL` (artifactId remains `xtrmETL`).
- README honest “Supported today” matrix; compose file product-name header.
+### Security
+
+- Durable-worker metrics and ordinary logs exclude payloads, raw principals, raw idempotency keys, hashes, job and lease identifiers, SQL, exception messages, and unbounded exception labels.
+- Retained payload or response-ledger identity conflicts fail closed with `etl_job_integrity_failure`; an expired or superseded lease rolls back target, ledger, and terminal-state effects.
+
### Added (historical)
- Comprehensive documentation suite (2026-01-08)
diff --git a/docs/doctoring/durable-job-retry-boundary-evidence.md b/docs/doctoring/durable-job-retry-boundary-evidence.md
new file mode 100644
index 00000000..e40ddfdf
--- /dev/null
+++ b/docs/doctoring/durable-job-retry-boundary-evidence.md
@@ -0,0 +1,101 @@
+# Durable-job retry and success-boundary doctoring evidence
+
+## Decision
+
+One durable worker claim represents exactly one persisted execution attempt. The worker must call
+`EtlService.processDataInExistingTransaction` rather than the synchronous `processData` entry point.
+
+The synchronous endpoint keeps `@Retryable` outside `@Transactional` so each request retry can own a
+fresh transaction. The durable worker already owns bounded retries through `attempt_count`, claim
+renewal by a later poll, and lease-fenced lifecycle transitions. Its target writes, response-ledger
+write, and exact-live-lease success transition must remain inside one caller-owned transaction.
+
+`EtlJobLeaseRepository.markSucceeded` therefore rejects direct invocation when no actual Spring
+transaction is active. Retry and failure transitions remain independently persistable after an
+execution exception, but success can never be published separately from the effects it certifies.
+
+## Failure modes prevented
+
+### Retry advice inside an existing durable transaction
+
+Calling the synchronous `@Retryable` method from an already active durable execution transaction can
+place retry advice inside an outer transaction that it did not create. After a transactional database
+failure, another in-process invocation can reuse a rollback-only or otherwise failed transaction
+instead of starting the fresh transaction assumed by Spring Retry. It also performs retries that are
+not represented by the durable job's `attempt_count`.
+
+### False success outside the atomic execution transaction
+
+A public success-transition method that can run in autocommit mode allows accidental callers to mark
+a job `SUCCEEDED` without the target rows and response ledger being committed in the same unit of
+work. Even an exact lease predicate cannot prove those effects exist. Requiring an actual transaction
+before the success SQL executes makes the repository fail closed at its public boundary.
+
+The complete contract is:
+
+```text
+one database claim
+→ one existing execution transaction
+→ one non-retrying ETL invocation
+→ target rows + response ledger + exact-live-lease success
+→ one atomic commit
+
+or
+
+one escaped failure
+→ worker-owned durable retry / terminal-failure decision
+```
+
+No in-process retry happens inside the lease transaction. A transient exception escapes to
+`EtlJobWorker`, which either returns the exact live lease to `PENDING` while attempts remain or records
+a stable terminal failure after the configured maximum.
+
+## Test-first evidence
+
+The regression contracts were added before their production behavior:
+
+- `EtlJobIdempotencyRetryBoundaryTest` required durable execution to call
+ `processDataInExistingTransaction` exactly once and never call `processData`;
+- `EtlServiceIdempotencyTransactionBoundaryTest` required the durable ETL entry point to reject
+ direct use without an actual transaction before JDBC or request-lock access;
+- `EtlJobLeaseSuccessTransactionBoundaryTest` required `markSucceeded` to reject use without an
+ actual transaction before JDBC access;
+- `EtlJobLeaseRepositoryIntegrationTest` executes successful lease fencing inside a real Spring test
+ transaction and still verifies expiry and supersession rejection.
+
+Production then added the non-retrying, transaction-requiring ETL entry point, changed
+`EtlJobIdempotencyService` to use it, and added the active-transaction guard to the public success
+transition. Existing synchronous processing retains its retry behavior.
+
+## Review and operational evidence
+
+Reviewers should confirm all of the following on the exact current head:
+
+1. `processDataInExistingTransaction` has neither `@Retryable` nor `@Transactional`;
+2. it fails closed when no actual Spring transaction is active;
+3. `EtlJobIdempotencyService` invokes only that entry point for a newly executed job;
+4. response replay does not invoke ETL target writes;
+5. `EtlJobLeaseRepository.markSucceeded` fails before JDBC without an actual transaction;
+6. `EtlJobExecutionService.execute` owns the transaction containing ETL, ledger, and success;
+7. transient exceptions escape to `EtlJobWorker` and affect durable attempt accounting once;
+8. retry and terminal-failure transitions remain exact-lease fenced;
+9. target rows, response ledger, and `SUCCEEDED` roll back together when success fencing fails;
+10. statement and branch coverage gates remain at 100% for the configured production scope.
+
+## Rollback
+
+Disable the durable worker before reverting either boundary. Do not restore synchronous retry advice
+inside the lease transaction, and do not allow `markSucceeded` to run in autocommit mode. A safe
+replacement must preserve one persisted attempt per claim and prove that target effects, the response
+ledger, and terminal success commit or roll back together.
+
+## References — APA 7th
+
+Spring Retry Authors. (2026). *EnableRetry.java* [Source code]. GitHub.
+https://github.com/spring-projects/spring-retry/blob/main/src/main/java/org/springframework/retry/annotation/EnableRetry.java
+
+Spring Retry Authors. (2026). *RetryOperationsInterceptor.java* [Source code]. GitHub.
+https://github.com/spring-projects/spring-retry/blob/main/src/main/java/org/springframework/retry/interceptor/RetryOperationsInterceptor.java
+
+Spring Framework Authors. (2026). *TransactionSynchronizationManager.java* [Source code]. GitHub.
+https://github.com/spring-projects/spring-framework/blob/main/spring-tx/src/main/java/org/springframework/transaction/support/TransactionSynchronizationManager.java
diff --git a/docs/etl/durable-job-intake.md b/docs/etl/durable-job-intake.md
index 7e29a2d2..dc54f365 100644
--- a/docs/etl/durable-job-intake.md
+++ b/docs/etl/durable-job-intake.md
@@ -1,24 +1,25 @@
-# Durable asynchronous ETL job intake
+# Durable asynchronous ETL jobs
-## Scope
+## Scope and activation
-`POST /api/etl/jobs` creates a durable, authenticated-principal-scoped ETL job resource. This
-bounded intake slice persists accepted work and exposes its status monitor; it does not execute jobs
-yet. Execution, PostgreSQL `FOR UPDATE SKIP LOCKED` claiming, lease fencing, bounded attempts,
-terminal payload clearing, and crash recovery belong to the following worker and lease-fencing slice.
+`POST /api/etl/jobs` creates a durable, authenticated-principal-scoped ETL job resource. A separate
+lease-fenced worker claims accepted jobs across replicas, replays or writes the durable response
+ledger, writes target rows, and commits terminal state atomically.
-The incomplete intake surface is disabled by default. It is absent unless an operator explicitly
-sets the preferred `mightyetl.etl.jobs.intake-enabled=true` property, its supported legacy alias
-`xtrmetl.etl.jobs.intake-enabled=true`, or `ETL_JOB_INTAKE_ENABLED=true`. When both full namespaces
-are supplied, `mightyetl.*` wins. Enabling intake accepts the temporary boundary that submitted
-payloads remain retained in `PENDING` jobs until the worker and terminal payload-clearing slice is
-implemented. Deployments that cannot accept that retention boundary must leave the setting false.
+Both capabilities are fail-closed:
-The existing synchronous `POST /api/etl/process` endpoint remains unchanged.
+```text
+mightyetl.etl.jobs.intake-enabled=false
+mightyetl.etl.jobs.worker.enabled=false
+```
-## Submit a job
+The supported legacy aliases use the `xtrmetl.*` namespace. Environment variables are
+`ETL_JOB_INTAKE_ENABLED` and `ETL_JOB_WORKER_ENABLED`. When both full namespaces are supplied,
+`mightyetl.*` wins. Intake may be enabled alone for a controlled retention window, and the worker may
+be enabled alone to drain accepted work. The existing synchronous `POST /api/etl/process` endpoint
+remains unchanged.
-A client sends:
+## Submit a job
```http
POST /api/etl/jobs HTTP/1.1
@@ -29,18 +30,13 @@ Idempotency-Key: "550e8400-e29b-41d4-a716-446655440000"
[{"id":"record_alpha","name":"accepted"}]
```
-The service requires:
-
-- the same authenticated principal for every retry;
-- the same semantic `Idempotency-Key`; and
-- byte-for-byte same JSON text for every retry of that key.
-
-The preferred header representation is an RFC 9651 quoted String. The legacy raw safe-ASCII profile
-remains accepted for compatibility and normalizes to the same semantic key.
+The service requires the same authenticated principal, the same semantic idempotency key, and
+byte-for-byte identical JSON text for every retry. The preferred header representation is an RFC
+9651 quoted String. The legacy raw safe-ASCII profile remains accepted and normalizes to the same
+semantic key.
-A new or replayed durable submission returns RFC 9110 `202 Accepted` because acceptance does not mean
-that processing has completed. The representation describes the current state and the `Location`
-header identifies the status monitor:
+A new or replayed submission returns RFC 9110 `202 Accepted`; acceptance does not mean processing is
+complete. The `Location` header identifies the owner-scoped status monitor:
```http
HTTP/1.1 202 Accepted
@@ -56,15 +52,10 @@ Content-Type: application/json
}
```
-All successful and covered problem responses for durable job resources include
-`Cache-Control: no-store`. These authenticated operational resources must not be retained by shared
-or private caches.
-
-A retry that resolves to the same durable resource returns the same job identifier and
-`Idempotency-Replayed: true`. Reusing the same principal-scoped key with different JSON text returns
-`422 etl_job_submission_key_reused`. A concurrent creation attempt that cannot acquire the
-transaction-level submission lock returns `409 etl_job_submission_in_progress` rather than waiting
-without a client-visible bound.
+All successful and covered problem responses include `Cache-Control: no-store`. A replay returns the
+same job identifier and `Idempotency-Replayed: true`. Reusing one principal-scoped key with different
+JSON returns `422 etl_job_submission_key_reused`. A concurrent creation attempt that cannot acquire
+the transaction-level submission lock returns `409 etl_job_submission_in_progress`.
## Read job status
@@ -73,68 +64,100 @@ GET /api/etl/jobs/{job_record_id} HTTP/1.1
Authorization: Basic
```
-The service hashes the current authenticated principal and queries by both principal scope and job
-identifier. A malformed or missing identifier and an identifier owned by another principal all
-return `404 etl_job_not_found`; callers cannot use this endpoint to probe another tenant's job
-existence.
+The query binds the current principal hash and job identifier. A malformed, missing, or foreign-owned
+identifier returns the same `404 etl_job_not_found`, preventing tenant-existence probing.
+
+The representation exposes only the opaque job identifier, stable lifecycle state, bounded attempt
+count, stable failure code where applicable, status URL, and timestamps. It excludes request payload,
+raw principal, raw submission key, internal hashes, lease identifiers, SQL, and response-ledger data.
+
+## Lifecycle and distribution
+
+Flyway migrations create descriptive multi-word `snake_case` objects:
-The response excludes the request payload, raw principal, raw submission key, and all internal
-hashes. Timestamps are explicit ISO-8601 strings. Before worker execution is implemented, newly
-accepted jobs remain `PENDING` with an `attemptCount` of zero.
+- `V2__create_etl_job_records.sql` creates `etl_job_records` and the submission uniqueness contract;
+- `V3__add_etl_job_lease_fencing.sql` adds `lease_claim_id`, `lease_owner_id`,
+ `lease_expires_at`, lifecycle constraints, and `etl_job_claim_eligibility_index`.
-## Validation and persistence
+The stable lifecycle is `PENDING`, `RUNNING`, `SUCCEEDED`, and `FAILED`.
-Before lock or table access, mightyETL enforces the same configured UTF-8 payload and record-count
-bounds used by synchronous ETL admission. The complete body must be a JSON array, duplicate JSON
-fields are rejected, every element must be an object with a safe textual `id`, and normalized field
-names must remain unique.
+Each fixed-delay poll handles at most one job. PostgreSQL, not scheduler uniqueness, distributes work:
-Flyway migration `V2__create_etl_job_records.sql` creates `etl_job_records`. All schema objects use
-descriptive multi-word `snake_case` names. The database stores:
+1. rows at the attempt limit are terminalized and their payloads are cleared;
+2. one oldest eligible `PENDING` row or expired `RUNNING` row is selected with
+ `FOR UPDATE SKIP LOCKED`;
+3. a fresh claim identifier, process owner identifier, database-derived expiry, and incremented
+ attempt count are persisted;
+4. execution verifies the retained payload digest and acquires a domain-separated response-ledger
+ lock derived only from stored hashes;
+5. an existing matching response is replayed, or target rows and `etl_idempotency_records` are written;
+6. `SUCCEEDED` is committed only for the exact unexpired claim in the same transaction.
-- an opaque UUID job identifier;
-- SHA-256 hashes of the principal scope, semantic submission key, and exact JSON text;
-- the request payload needed by the future worker while status is `PENDING` or `RUNNING`;
-- status, attempt, failure, and timestamp fields.
+If the lease is expired or superseded, the final transition fails and rolls back target and ledger
+writes. An expired row can be reclaimed with a new claim identifier. A stale worker therefore cannot
+commit duplicate target effects or terminalize a newer owner's work.
-The schema reserves the stable lifecycle vocabulary `PENDING`, `RUNNING`, `SUCCEEDED`, and `FAILED`.
-A database check requires a non-null request payload only for the two nonterminal states and requires
-that payload to be null for both terminal states. This makes terminal payload clearing an enforced
-persistence invariant rather than a documentation-only convention.
+## Retry and failure behavior
-Raw authenticated principal names and raw idempotency keys are never persisted. The request payload
-is sensitive operational data and must inherit the classification of its source records. Until the
-worker slice reaches a terminal state and clears it, operators must apply database access control,
-encryption, backup, and retention policy accordingly.
+Transient database failures return the job to `PENDING` while attempts remain. At the configured
+limit they become `FAILED` with `etl_target_unavailable`. Non-transient database failures use
+`etl_target_failure`. Persisted payload or ledger identity conflicts use
+`etl_job_integrity_failure`. Unexpected runtime failures use `etl_internal_error`. Deterministic ETL
+validation retains its existing stable `etl_*` request code. Eligible rows already at the attempt
+limit use `etl_worker_attempts_exhausted`.
-## Operational boundary
+Every retry or terminal transition repeats the exact live lease predicate. A zero-row transition is
+stale evidence and does not overwrite the authoritative owner.
-This slice deliberately does not advertise job completion or background execution. The controller is
-disabled by default; setting an activation property to `true` is an explicit operator opt-in to
-durable intake without execution. Deployments that need completed asynchronous processing must wait
-for the worker and lease-fencing slice. The next slice must claim jobs safely across replicas, fence
-stale lease owners, commit target effects and terminal success atomically, reclaim expired leases,
-bound attempts, publish stable failure codes, and clear the stored request payload at terminal state.
+## Validation, privacy, and retention
+
+Before submission lock or table access, mightyETL enforces configured UTF-8 payload and record-count
+bounds. The complete body must be a JSON array, duplicate JSON fields are rejected, every element
+must be an object with a safe textual `id`, and normalized field names must remain unique.
+
+The database stores an opaque UUID, SHA-256 hashes of principal scope, semantic submission key, and
+exact JSON text, the retained request payload while nonterminal, lifecycle and attempt fields, and
+lease metadata while running. Raw principal names and raw idempotency keys are never persisted.
+
+The request payload inherits the source records' data classification. Database constraints require a
+payload for nonterminal rows and require it to be null for terminal rows. Success, deterministic
+failure, attempts exhaustion, and non-retryable failure clear the payload in their terminal
+transition. Apply least privilege, encryption, backup, restore, and retention controls while data is
+retained.
+
+Metrics and ordinary logs must not include payloads, principals, client keys, hashes, job or lease
+identifiers, SQL, exception messages, or unbounded error classes. Operational procedures and metric
+contracts are authoritative in `docs/operations/durable-job-worker.md`.
## Standards basis
-- RFC 9110 Section 15.3.3 defines `202 Accepted` as noncommittal and recommends that the response
- describe current status and point to a status monitor.
-- RFC 9457 supplies the problem-details representation used by deterministic submission and lookup
- failures.
-- RFC 9651 defines the current Structured Fields String syntax accepted for `Idempotency-Key`.
-- The expired IETF HTTPAPI `Idempotency-Key` draft-07 is used only as work-in-progress design
- evidence for unique client keys, request fingerprints, `422` payload conflicts, and tenant-isolation
- security concerns. It expired on April 18, 2026 and is not represented as a published RFC.
+- RFC 9110 Section 15.3.3 defines `202 Accepted` as noncommittal and recommends a current-status
+ representation and status monitor.
+- RFC 9457 supplies deterministic problem-details representations.
+- RFC 9651 defines the accepted Structured Fields String syntax.
+- PostgreSQL 18 documents `SKIP LOCKED` as unsuitable for a general consistent view but useful for
+ avoiding contention among multiple consumers of a queue-like table.
+- Spring fixed-delay scheduling measures each delay from completion of the preceding invocation.
+- OpenTelemetry SQL/PostgreSQL semantic conventions define stable database telemetry fields; raw
+ query text and parameters remain privacy-sensitive opt-in data.
### References
-- Fielding, R., Nottingham, M., & Reschke, J. (2022). *HTTP semantics* (RFC 9110). RFC Editor.
- https://www.rfc-editor.org/rfc/rfc9110
-- Jena, J., & Dalal, S. (2025). *The Idempotency-Key HTTP header field*
- (draft-ietf-httpapi-idempotency-key-header-07, expired April 18, 2026). Internet Engineering Task
- Force. https://datatracker.ietf.org/doc/draft-ietf-httpapi-idempotency-key-header/
-- Nottingham, M., & Wilde, E. (2023). *Problem details for HTTP APIs* (RFC 9457). RFC Editor.
- https://www.rfc-editor.org/rfc/rfc9457
-- Nottingham, M., & Kamp, P. (2024). *Structured field values for HTTP* (RFC 9651). RFC Editor.
- https://www.rfc-editor.org/rfc/rfc9651
+Fielding, R., Nottingham, M., & Reschke, J. (2022). *HTTP semantics* (RFC 9110). RFC Editor.
+https://www.rfc-editor.org/rfc/rfc9110
+
+Nottingham, M., & Wilde, E. (2023). *Problem details for HTTP APIs* (RFC 9457). RFC Editor.
+https://www.rfc-editor.org/rfc/rfc9457
+
+Nottingham, M., & Kamp, P. (2024). *Structured field values for HTTP* (RFC 9651). RFC Editor.
+https://www.rfc-editor.org/rfc/rfc9651
+
+OpenTelemetry Authors. (2026). *OpenTelemetry semantic conventions 1.43.0: Semantic conventions for
+SQL databases client operations*. Cloud Native Computing Foundation.
+https://opentelemetry.io/docs/specs/semconv/db/sql/
+
+PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation: SELECT*.
+https://www.postgresql.org/docs/18/sql-select.html
+
+Spring Authors. (2026). *Task execution and scheduling*. Broadcom.
+https://docs.spring.io/spring-framework/reference/integration/scheduling.html
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
diff --git a/docs/operations/durable-job-worker.md b/docs/operations/durable-job-worker.md
new file mode 100644
index 00000000..34723549
--- /dev/null
+++ b/docs/operations/durable-job-worker.md
@@ -0,0 +1,197 @@
+# Durable ETL job worker operations
+
+## Purpose and safety boundary
+
+The durable worker moves accepted `etl_job_records` from `PENDING` through `RUNNING` to
+`SUCCEEDED` or `FAILED`. PostgreSQL row state is the distribution and fencing authority. Spring's
+fixed-delay scheduler only initiates polls; it does not establish exclusivity across replicas.
+
+The worker is fail-closed. Both of the following product switches must be reviewed deliberately:
+
+```text
+mightyetl.etl.jobs.intake-enabled=true
+mightyetl.etl.jobs.worker.enabled=true
+```
+
+The supported legacy aliases are `xtrmetl.etl.jobs.intake-enabled` and
+`xtrmetl.etl.jobs.worker.enabled`. Environment variables are `ETL_JOB_INTAKE_ENABLED` and
+`ETL_JOB_WORKER_ENABLED`. When both full namespaces are configured, `mightyetl.*` wins.
+
+Enable intake without the worker only for controlled maintenance windows where retained `PENDING`
+payloads are acceptable. Enable the worker without intake only to drain already accepted work.
+
+## Configuration
+
+| Preferred property | Environment variable | Default | Constraint |
+| --- | --- | ---: | --- |
+| `mightyetl.etl.jobs.worker.enabled` | `ETL_JOB_WORKER_ENABLED` | `false` | explicit opt-in |
+| `mightyetl.etl.jobs.worker.fixed-delay-milliseconds` | `ETL_JOB_WORKER_FIXED_DELAY_MILLISECONDS` | `5000` | 1 through 86,400,000 |
+| `mightyetl.etl.jobs.worker.initial-delay-milliseconds` | `ETL_JOB_WORKER_INITIAL_DELAY_MILLISECONDS` | `5000` | 0 through 86,400,000 |
+| `mightyetl.etl.jobs.worker.lease-duration-seconds` | `ETL_JOB_WORKER_LEASE_DURATION_SECONDS` | `300` | 1 through 86,400 |
+| `mightyetl.etl.jobs.worker.max-attempts` | `ETL_JOB_WORKER_MAX_ATTEMPTS` | `3` | 1 through 100 |
+| `mightyetl.etl.jobs.worker.lease-owner-id` | deployment-specific | generated | 8–128 safe ASCII characters |
+
+Scheduler delays and lease durations have a one-day safety ceiling. Configuration binding and the
+lease repository enforce the same limit, so direct repository callers cannot bypass it. Values above
+the ceiling fail application binding or claim validation rather than creating an effectively
+permanent polling pause, arithmetic overflow, or multi-day stale-work recovery delay.
+
+Set an explicit `lease-owner-id` only when the deployment platform can guarantee one stable,
+non-sensitive value per process. Never use a hostname containing customer data, a pod annotation
+containing credentials, an email address, a tenant identifier, or a raw infrastructure token.
+
+Choose a lease duration longer than the normal high-percentile execution time plus database and
+network variance. The current slice does not renew leases. A lease that expires during execution
+causes the final success transition to fail and rolls back target and response-ledger writes. If a
+normal execution can exceed one day, do not increase the ceiling silently; implement and validate
+lease renewal as a separate fenced capability first.
+
+## Claim, execution, and recovery
+
+Each poll handles at most one job:
+
+1. Eligible rows at or above `max-attempts` become terminal `FAILED`; their payload and lease fields
+ are cleared with `etl_worker_attempts_exhausted`.
+2. The worker selects the oldest `PENDING` row or expired `RUNNING` row below the attempt limit using
+ `FOR UPDATE SKIP LOCKED`.
+3. The claim writes a new `lease_claim_id`, the process `lease_owner_id`, database-derived expiry,
+ and incremented attempt count.
+4. The execution transaction verifies the retained payload digest, acquires the domain-separated
+ response-ledger lock, replays or writes `etl_idempotency_records`, writes target rows, and then
+ conditionally marks the exact live lease `SUCCEEDED`.
+5. A stale, superseded, or expired lease cannot commit target rows, response-ledger rows, or terminal
+ state. The whole execution transaction rolls back.
+
+An expired `RUNNING` job is reclaimed with a new claim identifier. The earlier worker may continue
+using CPU, but its target and lifecycle writes cannot commit after losing the exact live lease.
+
+## Stable failure codes
+
+| Failure code | Meaning | Operator response |
+| --- | --- | --- |
+| `etl_worker_attempts_exhausted` | an eligible row had no remaining claim attempt | inspect target availability and payload validity before any future replay feature |
+| `etl_target_unavailable` | transient database failures consumed the attempt limit | restore database service and retain evidence for incident review |
+| `etl_target_failure` | non-transient database write failure | inspect schema, constraints, permissions, and target compatibility |
+| `etl_job_integrity_failure` | retained payload or response-ledger identity conflicted | stop affected workers, preserve database evidence, investigate tampering or inconsistent migration |
+| `etl_internal_error` | unexpected non-database runtime failure | inspect sanitized application diagnostics and open a defect |
+| existing `etl_*` request codes | retained request failed deterministic ETL validation | correct the producer or migration source; do not blindly retry |
+
+Terminal states clear `request_payload` in the same state transition. The status API exposes only the
+stable failure code, attempt count, lifecycle state, and timestamps to the authenticated owner.
+
+## Observability and SLO evidence
+
+The worker publishes finite-cardinality metrics only:
+
+- `etl.jobs.worker.outcomes{outcome=idle|succeeded|retried|failed|stale}`;
+- `etl.jobs.execution.duration{outcome=idle|succeeded|retried|failed|stale}`.
+
+Every completed poll increments exactly one terminal outcome counter and records exactly one matching
+duration sample. Idle polls count as `idle`; a database failure while persisting a retry or terminal
+transition counts as `failed`, leaves the fenced row recoverable through lease expiry, and is never
+mislabeled as `retried`, `succeeded`, or `stale`. Claim acquisition is an internal phase, not a second
+outcome series, so summing the outcome counters yields the completed poll count without double
+counting work-bearing polls.
+
+Do not add payloads, raw principals, raw idempotency keys, hashes, job identifiers, lease identifiers,
+SQL text, exception messages, or unbounded exception classes as metric tags or log fields.
+
+Recommended initial service-level indicators are:
+
+- accepted-to-terminal latency by status;
+- oldest eligible `PENDING` age;
+- expired `RUNNING` count;
+- terminal success ratio;
+- retry and stale outcome rates;
+- exhausted-attempt and integrity-failure counts;
+- database connection-pool saturation and transaction latency.
+
+A production SLO must be calibrated from representative load and recovery tests. Do not claim a
+numerical availability or latency SLO until monitoring, alert thresholds, and retained evidence have
+been validated in the buyer's deployment topology.
+
+For OpenTelemetry database telemetry, use the stable SQL/PostgreSQL semantic conventions where the
+instrumentation supports them. Prefer low-cardinality `db.query.summary`; treat raw `db.query.text`
+and query parameters as opt-in sensitive telemetry requiring a separate privacy assessment.
+
+## Incident procedures
+
+### Backlog growth
+
+1. Confirm intake and worker switches independently.
+2. Check database connectivity, pool saturation, lock waits, and worker failure outcomes.
+3. Compare oldest eligible `PENDING` age with execution duration.
+4. Add replicas only after confirming the database can support the additional claim and target-write
+ concurrency.
+5. Do not update lifecycle fields manually while workers are active.
+
+### Repeated stale outcomes
+
+1. Compare the configured lease duration with high-percentile transaction duration.
+2. Check clock-independent database latency and long-running statements; lease decisions use database
+ time.
+3. Verify every process has a safe, distinct lease owner identifier.
+4. Increase the lease duration only within the one-day ceiling and only after confirming that crash
+ recovery delay remains acceptable; implement lease renewal instead of exceeding the ceiling.
+
+### Integrity failure
+
+1. Disable the worker while preserving intake only if continued payload retention is acceptable.
+2. Snapshot the affected database under incident-response controls.
+3. Compare the job's stored request digest with a digest of the retained payload and compare the
+ domain-separated response-ledger row.
+4. Review migration, restore, replication, and unauthorized-write evidence.
+5. Do not disclose hashes or payloads in tickets, chat, dashboards, or ordinary logs.
+
+## Deployment and rollback
+
+Before enabling the worker:
+
+1. Apply the transactional `V3__add_etl_job_lease_fencing.sql` migration and then the nonblocking
+ `V4__add_etl_job_claim_eligibility_index.sql` migration.
+2. Confirm V4 uses `CREATE INDEX CONCURRENTLY`, its companion configuration contains
+ `executeInTransaction=false`, Flyway PostgreSQL transactional locking is disabled, and
+ `etl_job_claim_eligibility_index` is valid in the PostgreSQL catalog.
+3. Confirm the application principal has only the required table and advisory-lock permissions.
+4. Run migration, claim-contention, stale-lease rollback, response-replay, and target compatibility
+ tests against a production-equivalent PostgreSQL environment.
+5. Deploy with the worker disabled, inspect health and schema evidence, then enable a canary replica.
+6. Verify target, response-ledger, and terminal state atomicity before widening rollout.
+
+The dedicated rollout, invalid-index recovery, and concurrent rollback procedure is
+`docs/operations/durable-job-claim-index-rollout.md`.
+
+Rollback order is fail-closed:
+
+1. Disable intake when new accepted work must stop.
+2. Disable all workers and wait for active transactions to complete or roll back.
+3. Confirm no `RUNNING` rows remain; allow leases to expire if necessary.
+4. Decide whether `PENDING` payloads will be drained by the current version or retained under an
+ approved data-retention exception.
+5. Roll back application binaries before any schema compensation.
+6. Never edit or delete an applied Flyway versioned migration. Use a new forward compensating
+ migration only after every deployed binary no longer reads the lease columns.
+
+## Standards and primary documentation
+
+Fielding, R., Nottingham, M., & Reschke, J. (2022). *HTTP semantics* (RFC 9110). RFC Editor.
+https://www.rfc-editor.org/rfc/rfc9110
+
+OpenTelemetry Authors. (2026). *OpenTelemetry semantic conventions 1.43.0: Semantic conventions for
+SQL databases client operations*. Cloud Native Computing Foundation.
+https://opentelemetry.io/docs/specs/semconv/db/sql/
+
+PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation: SELECT*.
+https://www.postgresql.org/docs/18/sql-select.html
+
+PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation: CREATE INDEX*.
+https://www.postgresql.org/docs/18/sql-createindex.html
+
+Redgate Software. (2026). *Flyway PostgreSQL transactional lock setting*.
+https://documentation.red-gate.com/fd/flyway-postgresql-transactional-lock-setting-277579114.html
+
+Redgate Software. (2026). *Flyway script configuration*.
+https://documentation.red-gate.com/flyway/reference/script-configuration
+
+Spring Authors. (2026). *Task execution and scheduling*. Broadcom.
+https://docs.spring.io/spring-framework/reference/integration/scheduling.html
diff --git a/docs/superpowers/plans/2026-08-05-durable-job-lease-worker.md b/docs/superpowers/plans/2026-08-05-durable-job-lease-worker.md
new file mode 100644
index 00000000..cde8d2b0
--- /dev/null
+++ b/docs/superpowers/plans/2026-08-05-durable-job-lease-worker.md
@@ -0,0 +1,108 @@
+# Durable ETL Job Lease Worker Implementation Plan
+
+> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
+
+**Goal:** Execute accepted asynchronous ETL jobs safely across replicas with PostgreSQL claim locking, exact lease fencing, bounded retries, atomic success, terminal payload clearing, and operator-safe evidence.
+
+**Architecture:** A transaction-scoped repository owns claim and state transitions; a separate transactional execution service couples existing ETL target writes with an exact-live-lease success update; a fixed-delay coordinator classifies failures and performs retry or terminal transitions in a new transaction. PostgreSQL row state is the distribution and fencing authority, while scheduling only supplies repeated polling.
+
+**Tech Stack:** Java 25, Spring Boot, Spring JDBC transactions, Spring scheduling, PostgreSQL 18 SQL, Flyway, Micrometer, JUnit 5, Mockito, H2 compatibility tests, Maven/Jacoco.
+
+## Global Constraints
+
+- Preserve standalone operation and modular MSA compatibility with ContextualWisdomLab/.github, naruon, and other CWL services.
+- Database objects contain at least two descriptive words and use `snake_case`.
+- Worker activation is fail-closed and disabled by default.
+- Every public production type and method has beginner-readable Javadoc.
+- Added durable-job production code must have zero missed instruction, line, method, or branch coverage.
+- Payloads, principals, idempotency keys, job identifiers, lease identifiers, SQL, and exception messages never enter metrics or logs.
+- A stale or expired lease cannot commit target effects or state transitions.
+- Versioned Flyway migrations are immutable after publication; rollback uses a forward compensating migration.
+
+---
+
+### Task 1: Lock the schema and configuration contracts
+
+**Files:**
+- Create: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseMigrationTest.java`
+- Create: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobWorkerPropertiesTest.java`
+- Create: `etl-service/src/main/resources/db/migration/V3__add_etl_job_lease_fencing.sql`
+- Create: `etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobWorkerProperties.java`
+- Modify: `etl-service/src/main/java/com/xtrmetl/etl/EtlApplication.java`
+- Modify: `etl-service/src/main/resources/application.yml`
+
+**Interfaces:**
+- Produces: `EtlJobWorkerProperties` with `enabled`, `fixedDelayMilliseconds`, `initialDelayMilliseconds`, `leaseDurationSeconds`, `maxAttempts`, and `leaseOwnerId`.
+
+- [ ] Write migration and property tests first. Require the three lease columns, lifecycle constraints, claim index, fail-closed defaults, safe owner profile, and all numeric boundaries.
+- [ ] Run `./mvnw -B -pl etl-service -Dtest=EtlJobLeaseMigrationTest,EtlJobWorkerPropertiesTest test` and record the expected missing-file/type failure.
+- [ ] Add the migration, properties, application registration, and environment-backed defaults.
+- [ ] Re-run the focused tests and commit.
+
+### Task 2: Add exclusive claim and exact transition persistence
+
+**Files:**
+- Create: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobLeaseRepositoryIntegrationTest.java`
+- Create: `etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobLease.java`
+- Create: `etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobLeaseRepository.java`
+- Create: `etl-service/src/main/java/com/xtrmetl/etl/job/StaleEtlJobLeaseException.java`
+
+**Interfaces:**
+- Produces: `Optional claimNext(String leaseOwnerId, Duration leaseDuration, int maxAttempts)`.
+- Produces: `markSucceeded`, `releaseForRetry`, and `markFailed`, each returning only after an exact, unexpired lease transition or throwing `StaleEtlJobLeaseException`.
+
+- [ ] Write H2 integration tests for deterministic order, simultaneous single claim, expired reclaim, exhausted terminalization, success, retry, failure, and stale update refusal.
+- [ ] Run the focused test and record the missing-type failure.
+- [ ] Implement the two-statement lock-and-update claim transaction using `FOR UPDATE SKIP LOCKED` and database `CURRENT_TIMESTAMP`.
+- [ ] Implement exact-live-lease transition predicates and stable failure validation.
+- [ ] Re-run the focused test and commit.
+
+### Task 3: Couple ETL target effects to terminal success
+
+**Files:**
+- Create: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobExecutionServiceIntegrationTest.java`
+- Create: `etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobExecutionService.java`
+
+**Interfaces:**
+- Consumes: `EtlJobLease`, `EtlService.processData`, `EtlJobLeaseRepository.markSucceeded`.
+- Produces: `void execute(EtlJobLease lease)` in one Spring transaction.
+
+- [ ] Write integration tests proving target rows and `SUCCEEDED` commit together.
+- [ ] Add a stale-lease test that changes the claim before execution and asserts both the exception and zero committed target rows.
+- [ ] Run the focused test and record the missing-type failure.
+- [ ] Implement the minimal transactional service and re-run the tests.
+- [ ] Commit.
+
+### Task 4: Add bounded fixed-delay coordination and evidence
+
+**Files:**
+- Create: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobWorkerTest.java`
+- Create: `etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobWorker.java`
+
+**Interfaces:**
+- Consumes: repository claim/transitions, execution service, worker properties, `MeterRegistry`.
+- Produces: one `pollOnce()` invocation that claims at most one job and records finite outcomes.
+
+- [ ] Write tests for no work, success, transient retry, exhausted transient failure, deterministic request failure, non-transient target failure, unexpected failure, and stale evidence.
+- [ ] Run the focused test and record the missing-type failure.
+- [ ] Implement the conditional worker bean, fixed-delay method, failure classification, retry bound, duration timer, and finite-cardinality outcome counter.
+- [ ] Re-run the tests and commit.
+
+### Task 5: Complete operations, privacy, compatibility, and release evidence
+
+**Files:**
+- Modify: `docs/etl/durable-job-intake.md`
+- Create: `docs/operations/durable-job-worker.md`
+- Modify: `CHANGELOG.md`
+- Modify: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobMigrationDocumentationTest.java`
+- Modify: `etl-service/src/test/java/com/xtrmetl/etl/job/EtlJobCoveragePolicyTest.java`
+
+**Interfaces:**
+- Produces: authoritative activation, SLO, metrics, failure-code, recovery, retention, and rollback guidance.
+
+- [ ] Add documentation-first tests requiring activation pairs, privacy boundaries, exact failure codes, rollback ordering, and standards references.
+- [ ] Update the authoritative docs and changelog.
+- [ ] Run `./mvnw -B -pl etl-service test`.
+- [ ] Run `./mvnw -B test` across the full reactor.
+- [ ] Inspect Jacoco for zero missed durable-job instructions, lines, methods, and branches.
+- [ ] Open a stacked draft PR against `ci/hourly-opencode-nvidia-nim`, inspect every review and exact-head check, and mark ready only after all gates pass.
diff --git a/docs/superpowers/specs/2026-08-05-durable-job-lease-worker-design.md b/docs/superpowers/specs/2026-08-05-durable-job-lease-worker-design.md
new file mode 100644
index 00000000..ee347f50
--- /dev/null
+++ b/docs/superpowers/specs/2026-08-05-durable-job-lease-worker-design.md
@@ -0,0 +1,184 @@
+# Durable ETL Job Lease Worker Design
+
+## Status
+
+Accepted implementation design for issue #120. This is a bounded follow-on to the durable
+asynchronous intake merged in PR #119 and is stacked on PR #121 until that workflow-security
+prerequisite reaches `develop`.
+
+## Product outcome
+
+Accepted asynchronous ETL jobs must progress from `PENDING` to a terminal state without depending on
+the client connection or one service replica. The worker must distribute work across replicas through
+PostgreSQL row locking, fence stale owners, bound retry attempts, atomically couple response-ledger
+and target effects with terminal success, clear retained payloads at terminal state, and expose only
+stable non-sensitive status metadata through the existing owner-scoped API.
+
+## Scope
+
+This slice adds:
+
+- a PostgreSQL-owned claim operation using deterministic ordering and `FOR UPDATE SKIP LOCKED`;
+- process-lifetime `lease_owner_id` and per-claim `lease_claim_id` fencing;
+- lease expiry and reclaim;
+- bounded attempts with deterministic terminal failure codes;
+- fixed-delay polling that is disabled by default;
+- hashed durable execution identity copied from the accepted job;
+- response-ledger replay or creation, target writes, and conditional `SUCCEEDED` in one transaction;
+- retry and failure transitions that require the exact live lease;
+- finite-cardinality execution metrics;
+- migration, rollback, privacy, operations, and failure-recovery documentation.
+
+Cancellation, priorities, recurring schedules, manual replay, result-body exposure, and a dead-letter
+user interface remain out of scope.
+
+## Data model
+
+Flyway migration `V3__add_etl_job_lease_fencing.sql` adds the following descriptive `snake_case`
+columns to `etl_job_records`:
+
+- `lease_claim_id UUID` — unique token generated for every claim or reclaim;
+- `lease_owner_id VARCHAR(128)` — stable non-sensitive identifier for one worker process;
+- `lease_expires_at TIMESTAMPTZ` — database-time expiry boundary.
+
+A lifecycle constraint requires all three lease columns for `RUNNING` rows and requires all three to
+be null for every other state. A failure lifecycle constraint requires `failure_code` only for
+`FAILED` rows. The existing terminal-payload constraint remains authoritative. An eligibility index
+covers `job_status`, `lease_expires_at`, `created_at`, and `job_record_id`.
+
+The claim also carries the accepted row's `principal_scope_hash`, `submission_key_hash`, and
+`request_digest`. These independent lowercase SHA-256 values support durable execution without
+persisting or reconstructing raw authenticated principals or raw client idempotency keys.
+
+## Claim protocol
+
+`EtlJobLeaseRepository.claimNext` runs in one transaction:
+
+1. Terminalize eligible rows whose `attempt_count` has reached the configured maximum. Clear
+ `request_payload` and all lease columns and assign `etl_worker_attempts_exhausted`.
+2. Select one `PENDING` row or one expired `RUNNING` row with `attempt_count < max_attempts`, ordered
+ by `created_at, job_record_id`, using `FETCH FIRST 1 ROW ONLY FOR UPDATE SKIP LOCKED`.
+3. Read `CURRENT_TIMESTAMP` from the database in the same statement and derive the next expiry from
+ that database time.
+4. Generate a new `lease_claim_id`, increment `attempt_count`, set `RUNNING`, set the owner and
+ expiry, clear any prior failure code, and commit.
+
+The scheduler does not provide uniqueness. The database row lock and state predicate are the
+cross-replica authority. PostgreSQL documents `SKIP LOCKED` as suitable for avoiding contention among
+multiple consumers of a queue-like table while warning that it is not a general-purpose consistent
+view. That limitation is appropriate because each worker needs one exclusive claim rather than a
+complete snapshot.
+
+## Durable idempotent execution and fencing
+
+`EtlJobExecutionService.execute` starts one transaction and delegates to
+`EtlJobIdempotencyService` before attempting terminal success.
+
+The idempotency service:
+
+1. requires an actual Spring transaction;
+2. recomputes the SHA-256 digest of `request_payload` and compares it with the stored
+ `request_digest` before lock or table access;
+3. domain-separates and hashes `principal_scope_hash` plus `submission_key_hash` into a response
+ ledger key without recovering raw identity values;
+4. acquires the existing transaction-lifetime `EtlRequestLock` for that key;
+5. replays a matching `etl_idempotency_records` response or calls the existing validated
+ `EtlService.processData` target writer and inserts the response ledger row.
+
+The execution service then conditionally transitions the job to `SUCCEEDED` only when all of the
+following still match:
+
+- `job_record_id`;
+- `job_status = 'RUNNING'`;
+- exact `lease_claim_id`;
+- exact `lease_owner_id`;
+- `lease_expires_at > CURRENT_TIMESTAMP`.
+
+If the conditional update affects no row, `StaleEtlJobLeaseException` is thrown. The exception rolls
+back the same transaction, including target and response-ledger writes. An expired or superseded
+worker therefore cannot commit duplicate target effects, create a misleading response ledger, or
+terminalize a newer owner's job.
+
+## Failure policy
+
+The polling coordinator catches execution failures after the execution transaction rolls back and
+performs a separate exact-lease transition:
+
+- `TransientDataAccessException`: return to `PENDING` when attempts remain; otherwise terminal
+ `FAILED` with `etl_target_unavailable`;
+- `EtlJobIntegrityException`: terminal `FAILED` with `etl_job_integrity_failure`;
+- `EtlRequestException`: terminal `FAILED` with the existing stable request `errorCode`;
+- other `DataAccessException`: terminal `FAILED` with `etl_target_failure`;
+- other `RuntimeException`: terminal `FAILED` with `etl_internal_error`;
+- `StaleEtlJobLeaseException`: make no state change because another owner or expiry boundary is
+ authoritative.
+
+Every retry or failure update repeats the exact-live-lease predicate. A zero-row update is treated as
+stale evidence, not as success.
+
+## Scheduling and activation
+
+Spring fixed-delay scheduling is used because the next delay is measured after completion of the
+previous invocation. `mightyetl.etl.jobs.worker.enabled` and its supported `xtrmetl.*` alias default
+to `false`. Configurable values are bounded and validated:
+
+- `fixed-delay-milliseconds` > 0;
+- `initial-delay-milliseconds` >= 0;
+- `lease-duration-seconds` > 0;
+- `max-attempts` between 1 and 100;
+- `lease-owner-id` is 8–128 safe ASCII characters and defaults to a process-lifetime generated
+ identifier.
+
+One polling invocation claims at most one job. Horizontal throughput is achieved by replicas and
+repeated fixed-delay invocations rather than unbounded in-process fan-out.
+
+## Observability and privacy
+
+The worker emits a duration timer and a finite outcome counter for `idle`, `claimed`, `succeeded`,
+`retried`, `failed`, and `stale`. Metric tags never include payloads, principals, idempotency keys,
+hashes, job identifiers, SQL, lease identifiers, exception classes, or exception messages. Logs
+follow the same rule. Database client instrumentation should retain stable OpenTelemetry
+SQL/PostgreSQL semantic conventions and avoid opting raw query text or parameters into telemetry
+unless the deployment has separately assessed that exposure.
+
+## Testing strategy
+
+- Migration tests enforce descriptive names, lifecycle constraints, index shape, and rollback
+ instructions.
+- Repository integration tests use H2's supported `FOR UPDATE SKIP LOCKED` syntax to prove one live
+ claim, deterministic ordering, expiry reclaim, attempt increment, execution identity, and
+ exhaustion terminalization.
+- Idempotency integration tests prove first execution, response replay without duplicate target
+ writes, payload digest rejection before locking, ledger conflict rejection, transient lock
+ contention, and fail-closed transaction requirements.
+- Execution integration tests prove target rows, response ledger, and `SUCCEEDED` commit together and
+ prove a stale claim rolls all three effects back.
+- Coordinator tests cover every failure classification, retry bound, zero-work poll, metrics outcome,
+ and stale transition.
+- Property tests cover every validation boundary and generated owner identifier.
+- Documentation and coverage policy tests require complete public Javadoc and zero missed
+ instruction, line, method, and branch coverage for the durable-job package.
+
+## Rollback
+
+Before application rollback, stop all workers and disable intake. Allow active leases to expire,
+confirm no `RUNNING` rows remain, and decide whether pending payloads will be drained or retained
+under an approved exception. Roll back the application first. The three lease columns and eligibility
+index may be removed only after all rows are non-running and no deployed binary reads them. Flyway
+versioned migrations are not edited or deleted after publication; a forward compensating migration
+must perform any production schema reversal.
+
+## Standards and primary documentation
+
+Fielding, R., Nottingham, M., & Reschke, J. (2022). *HTTP semantics* (RFC 9110). RFC Editor.
+https://www.rfc-editor.org/rfc/rfc9110.html
+
+OpenTelemetry Authors. (2026). *OpenTelemetry semantic conventions 1.43.0: Semantic conventions for
+SQL databases client operations*. Cloud Native Computing Foundation.
+https://opentelemetry.io/docs/specs/semconv/db/sql/
+
+PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation: SELECT*.
+https://www.postgresql.org/docs/18/sql-select.html
+
+Spring Authors. (2026). *Task execution and scheduling*. Broadcom.
+https://docs.spring.io/spring-framework/reference/integration/scheduling.html
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/EtlApplication.java b/etl-service/src/main/java/com/xtrmetl/etl/EtlApplication.java
index c6796c2b..961b3d7d 100644
--- a/etl-service/src/main/java/com/xtrmetl/etl/EtlApplication.java
+++ b/etl-service/src/main/java/com/xtrmetl/etl/EtlApplication.java
@@ -1,6 +1,7 @@
package com.xtrmetl.etl;
import com.xtrmetl.etl.connector.ConnectorProperties;
+import com.xtrmetl.etl.job.EtlJobWorkerProperties;
import com.xtrmetl.etl.service.EtlBatchProperties;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@@ -8,6 +9,7 @@
import org.springframework.cloud.client.discovery.EnableDiscoveryClient;
import org.springframework.context.annotation.EnableAspectJAutoProxy;
import org.springframework.retry.annotation.EnableRetry;
+import org.springframework.scheduling.annotation.EnableScheduling;
/**
* Bootstraps the mightyETL transformation and loading service.
@@ -16,7 +18,12 @@
@EnableDiscoveryClient
@EnableAspectJAutoProxy(proxyTargetClass = true)
@EnableRetry
-@EnableConfigurationProperties({ConnectorProperties.class, EtlBatchProperties.class})
+@EnableScheduling
+@EnableConfigurationProperties({
+ ConnectorProperties.class,
+ EtlBatchProperties.class,
+ EtlJobWorkerProperties.class
+})
public class EtlApplication {
/**
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"
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);
+ }
+}
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");
+ }
+ }
+}
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.
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"
+ );
+ }
+ }
+}
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
+ ) {
+ }
+}
diff --git a/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobWorker.java b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobWorker.java
new file mode 100644
index 00000000..ef3737c2
--- /dev/null
+++ b/etl-service/src/main/java/com/xtrmetl/etl/job/EtlJobWorker.java
@@ -0,0 +1,212 @@
+package com.xtrmetl.etl.job;
+
+import com.xtrmetl.etl.service.EtlRequestException;
+import io.micrometer.core.instrument.Counter;
+import io.micrometer.core.instrument.MeterRegistry;
+import io.micrometer.core.instrument.Timer;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnBooleanProperty;
+import org.springframework.dao.DataAccessException;
+import org.springframework.dao.TransientDataAccessException;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import java.time.Duration;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+
+/**
+ * Polls and executes at most one durable ETL job per fixed-delay invocation.
+ *
+ *
The database claim repository, not the scheduler, distributes work across replicas. The
+ * worker classifies failures into stable non-sensitive codes, retries only transient database
+ * failures while attempts remain, and treats every failed exact-lease transition as stale evidence.
+ * Metrics use a fixed terminal-outcome vocabulary and never tag payloads, principals, keys, job
+ * identifiers, lease identifiers, SQL, exception classes, or exception messages. Every completed
+ * poll records exactly one terminal outcome counter and one matching duration sample, including
+ * idle polls and database failures while persisting retry or terminal transitions.
+ */
+@Component
+@ConditionalOnBooleanProperty(
+ prefix = "xtrmetl.etl.jobs.worker",
+ name = "enabled",
+ havingValue = true,
+ matchIfMissing = false
+)
+public class EtlJobWorker {
+
+ /** Stable target-unavailability code used after transient attempts are exhausted. */
+ public static final String TARGET_UNAVAILABLE_FAILURE_CODE = "etl_target_unavailable";
+
+ /** Stable non-transient database failure code. */
+ public static final String TARGET_FAILURE_CODE = "etl_target_failure";
+
+ /** Stable unexpected implementation failure code. */
+ public static final String INTERNAL_FAILURE_CODE = "etl_internal_error";
+
+ private static final String METRIC_OUTCOMES = "etl.jobs.worker.outcomes";
+ private static final String METRIC_DURATION = "etl.jobs.execution.duration";
+ private static final String IDLE_OUTCOME = "idle";
+ private static final String SUCCEEDED_OUTCOME = "succeeded";
+ private static final String RETRIED_OUTCOME = "retried";
+ private static final String FAILED_OUTCOME = "failed";
+ private static final String STALE_OUTCOME = "stale";
+ private static final List FINITE_OUTCOMES = List.of(
+ IDLE_OUTCOME,
+ SUCCEEDED_OUTCOME,
+ RETRIED_OUTCOME,
+ FAILED_OUTCOME,
+ STALE_OUTCOME
+ );
+
+ private final EtlJobLeaseRepository leaseRepository;
+ private final EtlJobExecutionService executionService;
+ private final EtlJobWorkerProperties properties;
+ private final MeterRegistry meterRegistry;
+ private final Map outcomeCounters;
+ private final Map outcomeTimers;
+
+ /**
+ * Creates one fail-closed worker and pre-registers its finite metric vocabulary.
+ *
+ * @param leaseRepository database claim and transition authority
+ * @param executionService atomic target-write and success boundary
+ * @param properties bounded worker configuration
+ * @param meterRegistry metrics registry for finite-cardinality evidence
+ */
+ public EtlJobWorker(
+ EtlJobLeaseRepository leaseRepository,
+ EtlJobExecutionService executionService,
+ EtlJobWorkerProperties properties,
+ MeterRegistry meterRegistry
+ ) {
+ this.leaseRepository = Objects.requireNonNull(
+ leaseRepository,
+ "leaseRepository must not be null"
+ );
+ this.executionService = Objects.requireNonNull(
+ executionService,
+ "executionService must not be null"
+ );
+ this.properties = Objects.requireNonNull(properties, "properties must not be null");
+ this.meterRegistry = Objects.requireNonNull(
+ meterRegistry,
+ "meterRegistry must not be null"
+ );
+
+ Map counters = new LinkedHashMap<>();
+ Map timers = new LinkedHashMap<>();
+ for (String outcome : FINITE_OUTCOMES) {
+ counters.put(
+ outcome,
+ Counter.builder(METRIC_OUTCOMES)
+ .description("Durable ETL worker terminal outcomes")
+ .tag("outcome", outcome)
+ .register(this.meterRegistry)
+ );
+ timers.put(
+ outcome,
+ Timer.builder(METRIC_DURATION)
+ .description("Duration of one durable ETL worker poll")
+ .tag("outcome", outcome)
+ .register(this.meterRegistry)
+ );
+ }
+ this.outcomeCounters = Map.copyOf(counters);
+ this.outcomeTimers = Map.copyOf(timers);
+ }
+
+ /**
+ * Claims and handles at most one eligible durable job.
+ *
+ *
Fixed delay is measured after this invocation completes. A database outage during claim or
+ * a database failure while persisting retry or terminal state is converted into a finite failed
+ * outcome without copying diagnostic text into application logs or telemetry. Spring invokes the
+ * method only when worker activation is explicitly enabled.
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;
+ }
+}
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.
+ */
+public class StaleEtlJobLeaseException extends RuntimeException {
+
+ /**
+ * Creates the stable non-sensitive stale-lease signal.
+ */
+ public StaleEtlJobLeaseException() {
+ super("The durable ETL job lease is stale or no longer owned");
+ }
+}
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);
}
}
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
diff --git a/etl-service/src/main/resources/application.yml b/etl-service/src/main/resources/application.yml
index dbb4ebcb..f5b566d5 100644
--- a/etl-service/src/main/resources/application.yml
+++ b/etl-service/src/main/resources/application.yml
@@ -27,9 +27,17 @@ xtrmetl:
max-payload-bytes: ${ETL_MAX_PAYLOAD_BYTES:1048576}
max-batch-records: ${ETL_MAX_BATCH_RECORDS:1000}
jobs:
- # Intake persists validated payloads but does not execute them in this bounded slice.
- # Keep disabled until an operator explicitly accepts the temporary retention boundary.
+ # Intake persists validated payloads for durable worker execution.
intake-enabled: ${ETL_JOB_INTAKE_ENABLED:false}
+ worker:
+ # Execution remains fail-closed until an operator enables the worker explicitly.
+ enabled: ${ETL_JOB_WORKER_ENABLED:false}
+ # Scheduler delays are bounded to one day (86,400,000 milliseconds).
+ fixed-delay-milliseconds: ${ETL_JOB_WORKER_FIXED_DELAY_MILLISECONDS:5000}
+ initial-delay-milliseconds: ${ETL_JOB_WORKER_INITIAL_DELAY_MILLISECONDS:5000}
+ # Leases are bounded to one day (86,400 seconds); longer jobs require lease renewal.
+ lease-duration-seconds: ${ETL_JOB_WORKER_LEASE_DURATION_SECONDS:300}
+ max-attempts: ${ETL_JOB_WORKER_MAX_ATTEMPTS:3}
# Warehouse/BI targets: SPI + config binding + validation + catalog; writes remain SCAFFOLD.
connectors:
databricks:
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
+ )
+ );
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
+ );
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
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();
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.
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.