Published on

Kill the Outbox Worker: Recovering PostgreSQL Claims with Leases and Generations

Authors

The worker has committed INPROGRESS. It has not recorded SENT. Then the process disappears.

The database is healthy. The next polling query runs successfully. Yet the event stays stuck: a query that selects only READY rows will never select it again. Restarting the worker does not change committed application state.

The earlier outbox publisher article left this as an explicit extension. This follow-up implements that extension in a small PostgreSQL lab and tests the database boundary with a real process kill. It is a new experiment, not a claim that this exact design was deployed in the earlier project.

The useful result is two separate protections: an expired lease makes work claimable again, and a generation prevents a previous owner from recording an outcome for the new claim. Neither prevents duplicate publication to Kafka.

The experiment and its limits

The companion example contains the schema, three SQL statements, six integration tests, and the intentionally unsafe baseline. It uses Python to coordinate PostgreSQL connections; no application framework is needed to observe the transaction boundary.

The recorded run used PostgreSQL 18.4, Python 3.12.14, and pg8000 1.31.5 on macOS x86_64. PostgreSQL ran in a separate temporary local cluster. The binary package was @embedded-postgres/darwin-x64@18.4.0-beta.17; the server itself reported 18.4. This identifies the tested environment, not a required or recommended upgrade.

The kill test starts a child process, lets it claim one event in autocommit mode, and waits for the child to print the returned claim. Only then does the parent send SIGKILL. That handshake matters: killing a process before commit would test rollback instead of abandoned committed work.

After verifying the persisted state, the test moves lease_until into the past. This is an explicit test fixture, not a measured thirty-second wait. There is no Kafka broker in this lab, so its results establish database recovery and ownership checks only.

First, reproduce the stranded row

The reduced baseline selects work this way:

SELECT id
FROM outbox_lab
WHERE status = 'READY'
  AND next_attempt_at <= statement_timestamp()
ORDER BY id
LIMIT 1
FOR UPDATE SKIP LOCKED;

The claim statement also changes the selected row to INPROGRESS and commits. The worker is then killed before recording any outcome. A second connection reads:

status       generation
INPROGRESS   1

Even after the lease fixture expires, the baseline returns no candidate. The problem is its eligibility predicate. PostgreSQL released the transaction's locks at commit; it did not promise to reverse the committed status when the application process died.

The lab also tests a different interruption: a claim transaction that remains open while another connection polls. The second connection selects a different row, and rolling back the first transaction makes its original row available again. PostgreSQL documents SKIP LOCKED as suitable for queue-like access, while warning that it does not provide a consistent general-purpose view. See the SELECT locking reference.

Those are different recovery cases. An uncommitted claim rolls back; a committed claim needs an application policy.

Give every claim an expiration and an identity

The lab's minimal table is:

CREATE TABLE outbox_lab (
    id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    event_id text NOT NULL UNIQUE,
    status text NOT NULL DEFAULT 'READY'
        CHECK (status IN ('READY', 'INPROGRESS', 'SENT')),
    generation bigint NOT NULL DEFAULT 0,
    lease_until timestamptz,
    next_attempt_at timestamptz NOT NULL DEFAULT statement_timestamp(),
    CHECK ((status = 'INPROGRESS') = (lease_until IS NOT NULL))
);

event_id identifies the logical message and stays the same across retries. generation identifies one attempt to own its publication. These identifiers have different lifetimes: using the generation as a new event ID would hide redelivery from downstream deduplication.

The check constraint prevents an INPROGRESS row without a lease. It also requires terminal and waiting states to clear the lease. Payloads, priorities, and consumer receipts are omitted to keep this experiment about publisher ownership.

The replacement claim statement accepts either due READY work or expired INPROGRESS work:

WITH candidate AS (
    SELECT id FROM outbox_lab
    WHERE (status = 'READY' AND next_attempt_at <= statement_timestamp())
       OR (status = 'INPROGRESS' AND lease_until <= statement_timestamp())
    ORDER BY id
    LIMIT 1
    FOR UPDATE SKIP LOCKED
)
UPDATE outbox_lab AS o
SET status = 'INPROGRESS', generation = o.generation + 1,
    lease_until = statement_timestamp() + interval '30 seconds'
FROM candidate AS c
WHERE o.id = c.id
RETURNING o.id, o.event_id, o.generation;

Selection, ownership change, and generation increment happen in one statement. The worker retains the returned generation. In the kill test, the first claim returns [1, "event-001", 1]; after expiration, the replacement returns [1, "event-001", 2].

Thirty seconds is an illustrative lease, not a tuned value. The example claims one row at a time. Increasing the batch size also means accounting for time spent waiting behind other messages in the same batch.

Reject the old worker's outcome

A killed process cannot resume. A paused process can. Imagine worker A stalls, its lease expires, and worker B reclaims the row. When A resumes, an update using only the row ID could mark B's in-flight work complete or send it back to the retry queue.

The completion statement therefore checks both ownership and state:

UPDATE outbox_lab SET status = 'SENT', lease_until = NULL
WHERE id = :id
  AND status = 'INPROGRESS'
  AND generation = :generation
  AND lease_until > clock_timestamp()
RETURNING id;

:id and :generation are driver-bound parameters. The lab chooses a strict policy: an expired owner cannot finalize even before another worker reclaims the row. An empty returned result means the ownership condition was not satisfied; the caller must not follow it with an unconditional update.

Retry needs exactly the same protection:

UPDATE outbox_lab
SET status = 'READY', lease_until = NULL,
    next_attempt_at = statement_timestamp() + interval '1 minute'
WHERE id = :id
  AND status = 'INPROGRESS'
  AND generation = :generation
  AND lease_until > clock_timestamp()
RETURNING id;

A late failure callback is as capable of corrupting a newer claim as a late success callback. The tests verify both paths. They also check that a retry cannot move a SENT row back to READY.

The claim uses one statement timestamp to determine eligibility and its illustrative deadline. Completion checks the database's wall clock when evaluating the predicate. PostgreSQL distinguishes transaction, statement, and wall-clock time in its date/time documentation. This is not a guarantee that a transaction commits before the lease deadline. Keep these statements short, commit each claim before external I/O, and do not hold a database transaction open around a broker send.

What the tests actually returned

Running the same six tests against the unsafe baseline produced four failures and two passes. With expired-claim selection and guarded outcome updates, all six passed.

ExperimentUnsafe baselineLease and generation version
Kill after claim commit, then expire the leaseNo replacement claimSame event, generation 2
Expired owner attempts completion before reclaimIncorrectly accepts completionRejects completion and retry
Generation 1 reports an outcome after generation 2 owns the rowIncorrectly accepts completionRejects old success and retry; accepts current success
Current owner schedules a delayed retryPassesPasses
Another connection polls while the first holds its row lockClaims the other rowClaims the other row
Late retry arrives after SENTIncorrectly returns row to READYRejects retry

The tests assert returned rows and persisted state. The child-process test checks the actual SIGKILL exit status. The stale-owner test separately constructs a newer live generation so it can test write protection independently of the reclaim query.

To reproduce, use a disposable PostgreSQL database and a role allowed to create schemas. From the example directory:

python3 -m venv .venv
.venv/bin/python -m pip install -r requirements.txt
export PGHOST=127.0.0.1 PGPORT=55439 PGUSER=outbox_lab PGDATABASE=postgres
# Set PGPASSWORD if required by your disposable server.
.venv/bin/python test_recovery.py
OUTBOX_LAB_BASELINE=1 .venv/bin/python test_recovery.py

The last command intentionally exits with status 1. Each test creates and drops its own UUID-named schema. The process-kill test requires macOS or Linux. The example README records the environment and the captured output; it does not include a throughput claim.

Recovery still allows duplicate publication

The database can reject an old owner's update. It cannot retract a request the old owner has already sent to Kafka.

One unresolved sequence is:

A claims event-001, generation 1
A sends event-001; broker accepts it
A stops before recording SENT
Lease expires
B claims event-001, generation 2
B sends event-001 again

Another is a paused A resuming and sending after B has reclaimed. A pre-send ownership check narrows the window but cannot atomically coordinate PostgreSQL ownership with a later network operation. Kafka's delivery-semantics documentation explains the limits of guarantees when processing involves external systems.

The consumer still needs stable event identity and a deliberate policy for repeated effects. The notification consumer article discusses the separate boundary between consumer state, offsets, and external delivery.

This lab does not test those broker or consumer behaviors. It demonstrates that an abandoned database claim can be recovered without allowing a previous generation to overwrite the replacement's state.

Decisions left for the application

Before using this pattern, choose a lease duration from observed send latency and pause behavior. A lease that is too short creates unnecessary reclamation; a long lease delays recovery. Long-running work may need renewal, which must also be guarded by current ownership and must have a policy for renewal failure.

The lab intentionally omits retry exhaustion. Count abandoned claims as well as reported send failures when designing an attempt budget, or repeated process deaths can bypass a catch-block-based counter. Exhausted work needs a visible failure state and a controlled repair path, not silent deletion.

For a larger table, inspect the plan for the combined eligibility predicate, choose indexes for the actual workload, and test backlog behavior. This tiny fixture says nothing about large-table performance or per-aggregate ordering. Useful operational signals include oldest eligible event age, expired claim count, reclamation rate, and rejected stale updates.

The next broker-level test should stop a publisher after acknowledgment but before recording SENT, then observe redelivery using the unchanged event ID. Until that test exists, keep the claim narrow: database claim recovery and stale-owner rejection are verified here; end-to-end exactly-once delivery is not.