Skip to content

Add SQL leased read-write locks - #1722

Merged
malatewang merged 15 commits into
mainfrom
feature/new-feature
Oct 6, 2026
Merged

malatewang merged 15 commits into
mainfrom
feature/new-feature

Conversation

@malatewang

@malatewang malatewang commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Purpose of the change

Add SQL-backed leased read/write locks for coordinating work across server instances.

Description

  • Add atomic SQLite and PostgreSQL lease storage with fencing tokens, expiration, renewal, and release.
  • Expose read/write lease acquisition and context managers that renew leases while work runs.
  • Add get_sql_lock_service() to the resource manager, using the configured session manager database.
  • The lock can be used as a distributed lock for synchronization to support scale out.

Type of change

  • New feature (non-breaking change which adds functionality)

How Has This Been Tested?

  • Unit Test
  • Integration Test

Test Results: uv run pytest packages/server/server_tests/memmachine_server/common/resource_manager/test_resource_manager.py packages/server/server_tests/memmachine_server/common/sql_lease_lock/test_service.py -q — 30 passed. uv run ruff check on the changed server and test modules — passed. PostgreSQL-specific tests are included in test_postgres.py but were not run locally for this PR preparation.

Checklist

  • My code follows the style guidelines of this project
  • I have performed a self-review of my own code
  • I have added tests for the new behavior
  • New and existing focused tests pass locally

Further comments

None.

@edwinyyyu

edwinyyyu commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Writer preference

A waiting writer has no representation in the database: it is a client polling every 25-50 ms, and each poll grants only at an instant with zero live holders. A read acquire checks only whether a writer holds the key, so while read leases overlap a writer never gets in.

Reproduced on the SQLite path: one acquire_write with a 3 s wait timeout behind two readers taking 100 ms read leases offset by 60 ms. The writer timed out while the readers took 48 leases.

Whether that matters depends on the intended consumer. The segment store's registry pattern (writers pin the row FOR SHARE, delete takes FOR UPDATE) gets writer preference from PostgreSQL's lock queue, since new FOR SHARE requests wait behind a queued FOR UPDATE. A consumer porting that shape onto this lock, readers serving a tenant and a writer deleting it, would find the tenant undeletable under continuous traffic.

If writer preference is wanted, it looks solvable with a small change: a leased pending-writer marker on the resource row (writer_waiting_until_ms). A write acquire that finds live holders stamps it with a short TTL and returns None; each poll refreshes it; a read acquire returns None while the stamp is live. A dead waiter unblocks readers when the stamp expires, the fencing token is unaffected, and the same SQL works on SQLite. Two trade-offs: a task already holding a read lease that acquires the same key again would wait behind a pending writer (today it succeeds), so same-key nesting must be forbidden for reads as it already effectively is for writes; and new readers stall until the current readers' bodies finish, which is the point. Full FIFO fairness (ticket rows) is the heavier alternative.

Is reader preference the intended semantics, or should this change? Either way the docstring should say which.

repro
import asyncio, tempfile, os, time
from datetime import timedelta
from sqlalchemy.ext.asyncio import create_async_engine
from memmachine_server.common.sql_lease_lock import SqlLeaseRWLockService, LockAcquireTimeout

async def main() -> None:
    path = os.path.join(tempfile.mkdtemp(), "locks.db")
    engine = create_async_engine(f"sqlite+aiosqlite:///{path}")
    service = SqlLeaseRWLockService(engine)
    await service.startup()
    stop = asyncio.Event()
    holds = 0

    async def reader(offset: float) -> None:
        nonlocal holds
        await asyncio.sleep(offset)
        while not stop.is_set():
            async with service.read_lock("k", lease_duration=timedelta(seconds=1)):
                holds += 1
                await asyncio.sleep(0.10)
            await asyncio.sleep(0.02)

    readers = [asyncio.create_task(reader(0.0)), asyncio.create_task(reader(0.06))]
    await asyncio.sleep(0.05)
    started = time.monotonic()
    try:
        lease = await service.acquire_write("k", lease_duration=timedelta(seconds=1), wait_timeout=timedelta(seconds=3))
        print(f"writer acquired after {time.monotonic() - started:.2f}s")
        await lease.release()
    except LockAcquireTimeout:
        print(f"writer starved: timed out after {time.monotonic() - started:.2f}s while readers took {holds} read leases")
    stop.set()
    await asyncio.gather(*readers)
    await engine.dispose()

asyncio.run(main())

🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

edwinyyyu commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Intended consumer, and leases versus SQL row locks

What is the intended consumer of this lock? The PR does not name one, nothing else references the service, and nothing in the diff uses it beyond the accessor. For the review I assumed the closest match: replacing the per-process _session_locks in EpisodicMemoryManager (#1655, finding 1), since the API mirrors AsyncRWLock's read_lock() / write_lock() and the service is built on the session database. If the consumer is something else, the two points below may not apply.

Under that assumption:

  1. Behavior change. AsyncRWLock guarantees that a waiting writer blocks new readers ("writers cannot be starved by a steady stream of readers", rw_locks.py). This lock does not (see the comment above). That consumer has many readers (every open) and few writers (delete, close), so a session under steady traffic could never be deleted or closed. That may or may not be acceptable, but it is a change in behavior and should be called out.

  2. Why leases rather than a row lock? The sections that need cross-process arbitration are short. Open's write branch already does SELECT ... FROM sessions WHERE session_key = :k before constructing the instance, and delete reads the same row. FOR SHARE on that read, held across the construction, and FOR UPDATE in delete give mutual exclusion with no extra round trips, no polling, no renewal task, no fencing token and no clocks, and PostgreSQL's lock queue supplies writer preference and FIFO, matching the in-process lock's guarantee. Release on disconnect is the database's job. The segment store already uses this shape across a remote write. SQLite has no row locks, but the segment store's SQLite path covers it: a self-checking UPDATE on the row as the transaction's first statement, so the driver opens a write transaction and SQLite's database-wide write lock serializes every such transaction, across processes as well, with the match count as the staleness check. Coarser than PostgreSQL's shared/exclusive split, but the same design works on both dialects. The lease instead costs two extra transactions per open, a renewal task per open, a lease_lock_resource row per session key that is never deleted, and polling waits.

Leases fit a different consumer: a claim that outlives a transaction, for example each process holding a read lease for as long as it caches the instance, so that delete waits for eviction everywhere. If that is the intent, it would help to say so, because then writer preference is mandatory rather than optional, and the fencing token needs a consumer that can actually check it.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Nit: acronym casing in class names

The server package capitalizes acronyms in CapWords names, as PEP 8 asks (HTTPServerError, not HttpServerError): SQLAlchemySegmentStore, SQLiteVectorStore, SQLConfigurationError, AsyncRWLock, LiteLLMLanguageModel, HTTPException. Across the package that is over 200 uses of the all-caps form against 13 of the mixed form, all of them the one pre-existing outlier SqlAlchemyConf.

SqlLeaseRWLockService and SqlLeaseStore use the mixed form, and the first one mixes both styles in a single name (Sql next to RW). Suggest SQLLeaseRWLockService and SQLLeaseStore. The snake_case names (get_sql_lock_service, sql_lease_lock, _sql_lock_service) are fine as they are.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Renewal failure during exit cancels a finished body and skips the release

_renew_loop checks stop_renewal before each renew but not after, and on any renewal exception calls owner.cancel() unconditionally (service.py:208-213). If the body finishes while a renew round trip is in flight, _lock_context is already in its finally: it has set stop_renewal and is awaiting renewal_task (service.py:251-252). If that in-flight renew then fails, the cancel lands on that await, inside the finally. The rest of the block is skipped: lease.release() never runs and the recorded renewal error is never raised. The caller gets a bare CancelledError for a block that completed, and if the failure was transient the lease stays live until it expires.

The window is one database round trip per renewal interval, so it is rare, but the failing renew can be the ordinary case the loop exists for: a lease that expired or was replaced during the last interval, discovered while exiting. The symptom is then CancelledError instead of LeaseLostError, which most asyncio code reads as shutdown.

Repro (SQLite, renew stubbed to block and then fail so the body exits while it is in flight):

script
import asyncio, tempfile, os
from datetime import timedelta
from sqlalchemy.ext.asyncio import create_async_engine
from memmachine_server.common.sql_lease_lock import SqlLeaseRWLockService

async def main() -> None:
    engine = create_async_engine(f"sqlite+aiosqlite:///{os.path.join(tempfile.mkdtemp(), 'locks.db')}")
    service = SqlLeaseRWLockService(engine)
    await service.startup()
    renew_started, let_renew_finish = asyncio.Event(), asyncio.Event()
    body_finished = False

    async def work() -> None:
        nonlocal body_finished
        async with service.write_lock("k", lease_duration=timedelta(milliseconds=120)) as lease:
            async def slow_failing_renew() -> None:
                renew_started.set()
                await let_renew_finish.wait()
                raise RuntimeError("transient db error during renew")
            lease.renew = slow_failing_renew
            await renew_started.wait()
            body_finished = True  # body exits normally with a renewal in flight

    task = asyncio.create_task(work())
    await renew_started.wait()
    await asyncio.sleep(0.02)  # owner is now in `finally`, awaiting renewal_task
    let_renew_finish.set()
    try:
        await task
        result = "context exited normally"
    except asyncio.CancelledError:
        result = "context raised CancelledError"
    except Exception as err:
        result = f"context raised {type(err).__name__}"
    probe = await service.try_acquire_write("k", lease_duration=timedelta(seconds=1))
    print(f"body_finished={body_finished}; {result}; lease released on exit={probe is not None}")
    await engine.dispose()

asyncio.run(main())

Output: body_finished=True; context raised CancelledError; lease released on exit=False

Fix: in _renew_loop, after a failed renew, return without cancelling when stop_renewal.is_set(). That check is race-free: the finally sets the event and awaits the task with no yield in between, and the loop's except block runs without yielding, so a clear flag means the owner is still in the body.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Replacing the loop's cancellation with LeaseLostError leaves a pending cancel request on the task

When renewal fails, _renew_loop calls owner.cancel(), and _lock_context catches the resulting CancelledError and raises the renewal error in its place (service.py:241-244, and the finally at 258-259 when the body swallowed the cancel). Since Python 3.11, Task.cancel() increments a counter that only Task.uncancel() decrements; delivering the CancelledError does not. Nothing here calls uncancel(), so after the caller catches LeaseLostError the task still reports cancelling() == 1. The Task.uncancel docs: code that suppresses a cancellation by catching CancelledError needs to call this method to remove the cancellation state.

Consequence: structured-concurrency primitives read that counter. TaskGroup.__aexit__ calls uncancel() after a child failure and, when the count does not reach zero, treats an outside cancellation as still pending; on 3.13+ it re-issues cancel() to keep the count stable, so the parent gets a fresh CancelledError at its next await. A worker that catches LeaseLostError, backs off and retries is the natural consumer, and it is the one this poisons.

Repro on Python 3.14:

script
import asyncio, tempfile, os
from datetime import timedelta
from sqlalchemy import text
from sqlalchemy.ext.asyncio import create_async_engine
from memmachine_server.common.sql_lease_lock import SqlLeaseRWLockService, LeaseLostError

async def main() -> None:
    engine = create_async_engine(f"sqlite+aiosqlite:///{os.path.join(tempfile.mkdtemp(), 'locks.db')}")
    service = SqlLeaseRWLockService(engine)
    await service.startup()
    me = asyncio.current_task()
    try:
        async with service.write_lock("k", lease_duration=timedelta(milliseconds=120)) as lease:
            async with engine.begin() as conn:
                await conn.execute(text("DELETE FROM lease_lock_holder WHERE lease_id = :id"), {"id": lease.lease_id})
            await asyncio.sleep(0.5)
    except LeaseLostError:
        pass
    print("cancelling() after LeaseLostError:", me.cancelling())

    async def child() -> None:
        await asyncio.sleep(0.01)
        raise ValueError("child failed")
    try:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(child())
    except BaseException as err:
        print("TaskGroup exit raised:", type(err).__name__)
    try:
        await asyncio.sleep(0.01)
        print("next await: fine")
    except asyncio.CancelledError:
        print("next await: CancelledError (stale cancel request re-delivered)")
    await engine.dispose()

asyncio.run(main())

Output:

cancelling() after LeaseLostError: 1
TaskGroup exit raised: ExceptionGroup
next await: CancelledError (stale cancel request re-delivered)

Fix: call owner.uncancel() wherever the loop's cancellation is replaced by the renewal error, in the except asyncio.CancelledError branch when renewal_errors is set and in the finally before raise renewal_errors[0]. One uncancel() per cancel() the loop issued; an outside cancellation that arrived at the same time keeps its own count.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Lock tables grow without bound

lease_lock_resource rows are never deleted, and expired lease_lock_holder rows are deleted only for the key being acquired (_store.py, try_acquire: DELETE ... WHERE key = :key AND expires_at_ms <= now). A key that is never touched again keeps its resource row and its last expired holders forever.

Measured on SQLite:

step lease_lock_resource lease_lock_holder
1000 acquire+release on distinct keys 1000 0
then 1000 leases left to expire on other distinct keys, then 100 acquires elsewhere 2100 1000 expired rows still present

For a handful of coordination keys this is nothing. For per-session or per-tenant keys, which a per-key lock invites, it is one row per key ever locked, forever, plus the expired holders. The resource row cannot simply be dropped when its last holder goes: generation has to survive to keep fencing tokens monotonic per key. Two ways out: take tokens from a global sequence (a PostgreSQL SEQUENCE; on SQLite a single counter row, already serialized by the write lock), so a resource row carries nothing worth keeping and can be deleted when no live holder remains, with expired holders swept in bounded batches rather than per key; or state in the docstring that keys are expected to be few and long-lived.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@malatewang

Copy link
Copy Markdown
Contributor Author

Intended consumer, and leases versus SQL row locks

What is the intended consumer of this lock? The PR does not name one, nothing else references the service, and nothing in the diff uses it beyond the accessor. For the review I assumed the closest match: replacing the per-process _session_locks in EpisodicMemoryManager (#1655, finding 1), since the API mirrors AsyncRWLock's read_lock() / write_lock() and the service is built on the session database. If the consumer is something else, the two points below may not apply.

Under that assumption:

  1. Behavior change. AsyncRWLock guarantees that a waiting writer blocks new readers ("writers cannot be starved by a steady stream of readers", rw_locks.py). This lock does not (see the comment above). That consumer has many readers (every open) and few writers (delete, close), so a session under steady traffic could never be deleted or closed. That may or may not be acceptable, but it is a change in behavior and should be called out.
  2. Why leases rather than a row lock? The sections that need cross-process arbitration are short. Open's write branch already does SELECT ... FROM sessions WHERE session_key = :k before constructing the instance, and delete reads the same row. FOR SHARE on that read, held across the construction, and FOR UPDATE in delete give mutual exclusion with no extra round trips, no polling, no renewal task, no fencing token and no clocks, and PostgreSQL's lock queue supplies writer preference and FIFO, matching the in-process lock's guarantee. Release on disconnect is the database's job. The segment store already uses this shape across a remote write. SQLite has no row locks, but the segment store's SQLite path covers it: a self-checking UPDATE on the row as the transaction's first statement, so the driver opens a write transaction and SQLite's database-wide write lock serializes every such transaction, across processes as well, with the match count as the staleness check. Coarser than PostgreSQL's shared/exclusive split, but the same design works on both dialects. The lease instead costs two extra transactions per open, a renewal task per open, a lease_lock_resource row per session key that is never deleted, and polling waits.

Leases fit a different consumer: a claim that outlives a transaction, for example each process holding a read lease for as long as it caches the instance, so that delete waits for eviction everywhere. If that is the intent, it would help to say so, because then writer preference is mandatory rather than optional, and the fencing token needs a consumer that can actually check it.

🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

This PR just implemented required distributed lock service. Other tasks that require distributed synchronization can use the service.

@edwinyyyu

edwinyyyu commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Two more nits

  1. StoredLease field order versus the table. lease_lock_holder declares lease_id, key, mode, fencing_token, expires_at_ms (_store.py:39-43); StoredLease declares key, mode, lease_id, expires_at_ms, fencing_token (_store.py:52-56); and try_acquire builds it positionally as StoredLease(key, mode, lease_id, expires_at_ms, token) (_store.py:146). The two integer fields sit next to each other in the opposite order from the columns, so a transposition there would type-check and pass a stale expiry off as a fencing token. Suggest one order for the columns, the dataclass and the insert, and a keyword construction.

  2. Return annotation of @asynccontextmanager functions. _transaction (_store.py:73), _lock_context (service.py:222) and the two test helpers are annotated -> AsyncIterator[...]. typeshed deprecated that form in Fix the signatures of @(async)contextmanager python/typeshed#12087 (merged 2026-01-23): the Iterator / AsyncIterator overloads of contextmanager / asynccontextmanager now carry @deprecated, and the accepted form is Generator[T, None, None] / AsyncGenerator[T, None], which is what rw_locks.py in this package already uses.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

edwinyyyu commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

The shared counter row serializes every acquire in the deployment on PostgreSQL

The rework takes fencing tokens from one row in lease_lock_counter via UPDATE ... RETURNING inside the acquire transaction (_store.py, _next_token), after the key row is locked. On PostgreSQL that row lock is held until the transaction commits, so an acquire on key A blocks every acquire on every other key until A's transaction is done. The previous version had per-key generations and no cross-key coupling.

Measured on postgres:16: with a transaction holding that same UPDATE uncommitted for 1.0 s, try_acquire on an unrelated key waited 0.93 s. Each acquire transaction is short, but this makes the lock service a single serialization point for the whole deployment, which is the opposite of what a scale-out primitive should be. SQLite is unaffected, since it already serializes all writers.

Fix: on PostgreSQL take the token from a SEQUENCE (nextval is non-transactional and takes no row lock; gaps are fine since only monotonicity matters). Keep the counter row for SQLite, where it costs nothing extra. Lock ordering stays key row first, so no deadlock either way.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@malatewang

Copy link
Copy Markdown
Contributor Author

The shared counter row serializes every acquire in the deployment on PostgreSQL

The rework takes fencing tokens from one row in lease_lock_counter via UPDATE ... RETURNING inside the acquire transaction (_store.py, _next_token), after the key row is locked. On PostgreSQL that row lock is held until the transaction commits, so an acquire on key A blocks every acquire on every other key until A's transaction is done. The previous version had per-key generations and no cross-key coupling.

Measured on postgres:16: with a transaction holding that same UPDATE uncommitted for 1.0 s, try_acquire on an unrelated key waited 0.93 s. Each acquire transaction is short, but this makes the lock service a single serialization point for the whole deployment, which is the opposite of what a scale-out primitive should be. SQLite is unaffected, since it already serializes all writers.

Fix: on PostgreSQL take the token from a SEQUENCE (nextval is non-transactional and takes no row lock; gaps are fine since only monotonicity matters). Keep the counter row for SQLite, where it costs nothing extra. Lock ordering stays key row first, so no deadlock either way.

🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

Such performance is not a major concern here. The distributed lock here is supposed to be used by long running background jobs. If the caller is sensitive to performance and short-lived task, other sql lock mechanism should be used.

@edwinyyyu

edwinyyyu commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

With long-running jobs as the consumer, the wait loop's fixed poll rate becomes the cost that matters

Taking the stated consumer as the contract, long-running background jobs with one owner per key, the acquires themselves are rare, but the waits are long: every standby replica sits in acquire or lock for as long as the current owner runs. _acquire polls with a fixed 25-50 ms jitter and no backoff (service.py, _acquire), and each poll is a full store transaction: key row locked (FOR UPDATE on PostgreSQL, BEGIN IMMEDIATE on SQLite), clock read, commit.

Measured on SQLite: one holder with a 5-minute lease and three waiters in acquire issued 378 acquire transactions in 5 s, 25 per waiter per second, and would keep doing so for the whole 5 minutes. For N standby replicas that is 25N write transactions per second on the session database, continuously, for the lifetime of the deployment. On SQLite each one takes the database-wide write lock the session manager also needs.

Fix: back the delay off while waiting, for example doubling up to some fraction of lease_duration (the holder cannot release earlier than it renews, so a long lease justifies a long poll), or take a poll_interval option so a consumer that expects to wait for minutes can say so. Either keeps the fast path for short waits.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@malatewang

Copy link
Copy Markdown
Contributor Author

With long-running jobs as the consumer, the wait loop's fixed poll rate becomes the cost that matters

Taking the stated consumer as the contract, long-running background jobs with one owner per key, the acquires themselves are rare, but the waits are long: every standby replica sits in acquire or lock for as long as the current owner runs. _acquire polls with a fixed 25-50 ms jitter and no backoff (service.py, _acquire), and each poll is a full store transaction: key row locked (FOR UPDATE on PostgreSQL, BEGIN IMMEDIATE on SQLite), clock read, commit.

Measured on SQLite: one holder with a 5-minute lease and three waiters in acquire issued 378 acquire transactions in 5 s, 25 per waiter per second, and would keep doing so for the whole 5 minutes. For N standby replicas that is 25N write transactions per second on the session database, continuously, for the lifetime of the deployment. On SQLite each one takes the database-wide write lock the session manager also needs.

Fix: back the delay off while waiting, for example doubling up to some fraction of lease_duration (the holder cannot release earlier than it renews, so a long lease justifies a long poll), or take a poll_interval option so a consumer that expects to wait for minutes can say so. Either keeps the fast path for short waits.

🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.
This makes sense. Changed the behavior

@edwinyyyu

edwinyyyu commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Which job depends on this?

"Long-running background jobs" names a class, not a job. Could you name the specific job, or the specific kind of job, you intend to add on top of this lock?

The reason it matters: the background work the server has today is item-shaped. The semantic ingestion loop and the deletion queue both consist of discrete units of work, and for those the standard mechanism is a claim on the next unit with SELECT ... FOR UPDATE SKIP LOCKED, which the segment store's purge path already uses. That runs in every replica in parallel with no leader and no duplicated work. A lease over the whole loop would instead turn N-fold duplication into a single-process bottleneck with a failover gap of one lease duration plus the backoff. Periodic sweeps fit the same mechanism by making the scheduled tick the claimable row.

A lease is the right tool for a continuous loop, a stream consumer, a subscription, a watcher, or anything that holds one exclusive external resource or in-process state for its lifetime. I do not know of such a loop in the server now. If one is planned, that is the consumer this PR needs, and naming it would settle the review. If the intended consumer is ingestion or deletion, the claim pattern fits better and this lock would not be the right primitive for it.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@malatewang

Copy link
Copy Markdown
Contributor Author

Which job depends on this?

"Long-running background jobs" names a class, not a job. Could you name the specific job, or the specific kind of job, you intend to add on top of this lock?

The reason it matters: the background work the server has today is item-shaped. The semantic ingestion loop and the deletion queue both consist of discrete units of work, and for those the standard mechanism is a claim on the next unit with SELECT ... FOR UPDATE SKIP LOCKED, which the segment store's purge path already uses. That runs in every replica in parallel with no leader and no duplicated work. A lease over the whole loop would instead turn N-fold duplication into a single-process bottleneck with a failover gap of one lease duration plus the backoff. Periodic sweeps fit the same mechanism by making the scheduled tick the claimable row.

A lease is the right tool for a continuous loop, a stream consumer, a subscription, a watcher, or anything that holds one exclusive external resource or in-process state for its lifetime. I do not know of such a loop in the server now. If one is planned, that is the consumer this PR needs, and naming it would settle the review. If the intended consumer is ingestion or deletion, the claim pattern fits better and this lock would not be the right primitive for it.

🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

If the lock is required by just update the database without lease, who is going to do clean up if the process crashed? This just provide a general tools.

@edwinyyyu

Copy link
Copy Markdown
Contributor

On cleanup after a crash

The two mechanisms, side by side, since the crash case is the one that distinguishes them.

A claim held in a transaction (the SKIP LOCKED shape) is cleaned up by the database. When the process dies, the client kernel closes its sockets, the server backend sees the connection drop, aborts the transaction and releases its locks. No application code is involved. The case the database cannot see is an unclean death where no close ever arrives: host failure, network partition, a frozen VM. Then the backend keeps the transaction and its locks until it notices, via TCP keepalive probes (PostgreSQL's tcp_keepalives_idle / _interval / _count, which fall back to the OS default, about two hours on Linux) or idle_in_transaction_session_timeout, or on PostgreSQL 14+ client_connection_check_interval. Those are server-side settings. They cost nothing in application code, and they bound the same exposure for every transaction the application holds, not only lock claims: a request handler's transaction orphaned by the same host failure pins the vacuum horizon and holds its row locks for exactly the same window. The engine configuration currently passes only asyncpg's command_timeout and connect timeout, none of these, so that gap exists independently of this PR.

A lease commits the holder's row and keeps no transaction open while the work runs. A dead holder therefore costs the database nothing, and the key frees itself after one lease duration, with no server setting involved. Its costs are the ones already discussed: polling while waiting, the fencing token that consumers must check, and the pause window after expiry.

So the crash question has an answer in both designs. What differs is the bound on an unclean death, a server timeout against a lease duration, and what else that bound covers, every transaction against this one lock.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

edwinyyyu commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

Where the crash cost does and does not apply

For the segment store purge (#1661) and the collection registry purge (#1631), both from @edwinyyyu's own recent work, a dead session should not matter to overall database performance provided the server timeouts are set. Their purge claims are held for one bounded batch per transaction: measured at 200 to 400 ms for a 10k-row segment batch on PostgreSQL, 50 to 100 ms per Milvus round and 1.2 to 1.6 s per one-shot Qdrant delete of a million points in the registry. With an idle-in-transaction timeout of a few minutes, an orphaned claim is cleared long before the vacuum horizon it pins costs anything, the lost work is one batch, and another purger picks the entry up. That is a configuration item for the engine rather than a property of either lock design.

The premise does hold for longer-running tasks that cannot be split into bounded batches, where a transactional claim would have to stay open for the whole task. A lease is the fitting shape for those, so the premise is not rejected; the question of which such task this lock is for remains the one that decides the review.

One detail on the timeouts, so they are not over-trusted: the idle-in-transaction timeout resets on every statement and TCP keepalives are answered by the kernel, so neither bounds a process that is alive and keeps issuing statements inside one transaction. Before PostgreSQL 17 nothing bounds total transaction age; 17 adds transaction_timeout. A lease has the mirror-image limit: a live holder that keeps renewing holds the key for as long as it likes. Both designs bound a dead or silent holder, not a live misbehaving one.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Context: where generic lease services sit in practice

A shared lease service is a standard piece of infrastructure, and this PR's shape has direct precedent. Google's Chubby paper made the case for a central lock service over a consensus library: existing systems can adopt it late and cheaply, and it doubles as a small name service. Kubernetes runs every controller's leader election through the same Lease API and one client library. In application code without such infrastructure, a database-table lease with expiry, renewal and a fencing token is what ShedLock, Quartz's clustered scheduler and the DynamoDB lock client are. All of these are used for coarse-grained, long-held ownership of a named resource: leader election, "run this job on one node", an exclusive external handle. Chubby's paper says explicitly that fine-grained locking was left to applications.

Work distribution, selecting and holding items from a backlog, is done differently everywhere: the lease is a property of the item inside the queue, and claiming is one operation that selects and leases together. SQS visibility timeouts, Kafka's partition assignment, Temporal's task tokens, and the PostgreSQL queue libraries (Oban, pg-boss, GoodJob, River) all work this way; the segment store purge here is the same pattern with the transaction as the hold. The reasons a separate lock service is not used for that:

  • A lock answers "may I have key K", not "which item next". Workers must enumerate candidates and probe keys one by one, and with a shared ordering they all probe the same head item.
  • Claim and item state live apart, so completion must be fenced by hand and the two can diverge; in a queue the claim and the state are one row and one statement.
  • Items carry lifecycle the lock has no place for: retry counts, backoff, not-before times, dead-lettering. Ephemeral keys also leave lock rows behind, as this service documents.
  • A dead worker's item is still in every other worker's candidate list, so each probes it and fails until the lease expires; a queue hides it by the same predicate that selects.
  • No ordering or fairness, which this service states, while queues provide it.

So the primitive is legitimate and its precedent is clear about its scope. The open question for this PR remains which singleton, long-running job it is for, since the server's current background work is item-shaped.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Renewal and cancellation: four remaining defects

The first two are reproduced on the current head (SQLite, store methods stubbed to control timing); the other two follow from the code.

  1. A transient renewal error kills a long-running job with most of its lease left. _renew_loop cancels the owner on the first exception of any kind (service.py, the except Exception around lease.renew()). With a 60 s lease renewed every 20 s, one pool timeout, dropped connection, or SQLite "database is locked" on the shared session file aborts the body with about 40 s of valid lease remaining. The lease is lost only when its expiry passes. The loop should retry renewal until then and cancel the owner only when renew() returns None, or when the time since the last successful renewal, measured locally, reaches the lease duration, which is the conservative bound.

  2. Cleanup cannot be bounded. _await_cleanup absorbs every cancellation while the cleanup task runs. With store.release stubbed to block, the owner task cancelled five times reports task.done()=False, cancelling()=5, an enclosing asyncio.timeout cannot end the block, and the process was still running 15 s later. A stuck database connection therefore pins the task for the process lifetime; asyncpg's command_timeout is optional and aiosqlite has none. The lease expiry is the designed fallback for a release that never lands, so cleanup can be bounded: absorb one cancellation to let a prompt release finish, then stop waiting and let the row expire.

  3. A database error on the cancelled acquire path replaces the cancellation. In _try_acquire, if the in-flight store operation raises while the cancellation is being drained, that exception propagates out of _await_cleanup and the raise of the CancelledError is never reached. Reproduced: a task cancelled while acquire() is in flight, whose store call then raises OperationalError, finishes with OperationalError, task.cancelled()=False, and asyncio logs "exception in shielded future". A caller that retries on database errors keeps running after it was told to stop. _finish_lock_context already guards this case; this path does not.

  4. An outside cancellation that coincides with a renewal failure is lost. The except asyncio.CancelledError branch raises the renewal error and the finally uncancels once. If an outside cancel() was pending at the same time, the count stays at one but asyncio never re-delivers it. The PR's own test_renewal_failure_restores_only_its_cancellation_state[True] ends in exactly that state: cancelling() == 1 and no CancelledError raised. Fix: after the uncancel, if owner.cancelling() is still positive, re-raise CancelledError instead of the renewal error.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Simplification and minor items

  • One conditional statement per store operation. Each operation is a four-to-five-statement transaction with two dialect-specific locking mechanisms (explicit BEGIN IMMEDIATE, SELECT ... FOR UPDATE), nullable columns that exist only for the placeholder insert, and a separate round trip for the database time. The same guarantee comes from a single statement: INSERT ... ON CONFLICT (key) DO UPDATE SET lease_id = ..., expires_at_ms = <now-expr> + :duration, fencing_token = ... WHERE lease_lock.lease_id IS NULL OR lease_lock.expires_at_ms <= <now-expr>, with the row count as the grant signal, and renew and release as one conditional UPDATE or DELETE with the now-expression inline. On PostgreSQL the token can be nextval of a sequence inside that statement, which also removes the counter-row serialization. That deletes _transaction, _lock_resource, _now_ms, the nullable columns and their is not None checks.

  • Backoff cap unrelated to the lease. The retry delay caps at a constant 30 s, so after a holder with a 2 s lease dies, a waiter already at the top rung sleeps up to 30 s before noticing. Capping at min(30 s, lease_duration) bounds takeover latency to the lease the caller chose.

  • Silent cleanup failures. _finish_lock_context discards cleanup exceptions, and when a cancellation was absorbed it drops the renewal error as well. A release that fails because the database is briefly unavailable leaves the key held until expiry with nothing logged. Log, or chain onto the propagating error.

  • No floor on the lease duration. _duration_ms rounds up to 1 ms and the renewal interval is a third of the duration, so a sub-millisecond lease renews in a busy loop against the database and loses the lease within its first renewal. Reject durations below a few renewal round trips.

  • fencing_token docstring contradicts the class docstring. The property says "a resource-local token that grows with each grant"; tokens now come from the global counter, as the class docstring states, so two grants of one key can carry 7 and 42 with no holder in between. The property is the one consumers are told to rely on.

  • Undocumented exits. Calling release() on the handle inside a lock() block gets the body cancelled on the next renewal, and a stall longer than two thirds of the lease makes a completed body raise LeaseLostError at exit. Both should be stated on lock() or the handle should not expose release() inside the block.

  • Table creation races across instances. startup() runs create_all from each instance at first use, and two instances starting together on a fresh PostgreSQL database both pass the existence check; the loser fails with a duplicate-table error and its first acquire raises. This is the schema-provisioning gap tracked in Schema provisioning runs from every process's boot path: create_all races on cold boot and never evolves a table, Alembic runs destructive migrations on first use #1570 and not specific to this PR, noted here because multi-instance startup is this PR's normal path.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@malatewang
malatewang force-pushed the feature/new-feature branch 2 times, most recently from 21ab046 to 4c2fa04 Compare October 2, 2026 18:56
@edwinyyyu

Copy link
Copy Markdown
Contributor

Verified the renewal and cleanup fixes; four new suppressions

Re-ran the probes against 99e4568: a transient renewal error no longer cancels the body (renewal retried, body completed); a cancelled acquire whose store call then raises now ends cancelled; a hanging release is abandoned after the 5 s cleanup deadline with the cancellation delivered; an outside cancellation coinciding with a renewal failure is honored; the exit-race and uncancel fixes still hold. Tests, ty and ruff pass.

The fix adds four # noqa lines on top of the existing N818 one. Two are SLF001 for reading lease._safe_until from the renewal loop, which is the signal that Lease should expose that as a method or property (for example a trusted_for() returning the remaining seconds, or the loop asking the handle to renew-with-deadline itself), since the loop and the handle are the same package and the value is part of the handle's contract. The two TRY301 ones go away by moving the raise out of the try block: compute remaining before the try, and raise LeaseLostError or the cleanup TimeoutError outside it. None of the five is needed to interface with a third-party library, so they should all go.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@edwinyyyu

Copy link
Copy Markdown
Contributor

Decision context: #1756 (background work across replicas) lists the consumers that exist in the server (session provision and delete jobs, ingestion claims per set, the two purges) and the claim primitive they share. It tracks this PR as "fold into the claim primitive, or close without a consumer", so the answer to "which job depends on this" decides it.


🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

@malatewang

Copy link
Copy Markdown
Contributor Author

Decision context: #1756 (background work across replicas) lists the consumers that exist in the server (session provision and delete jobs, ingestion claims per set, the two purges) and the claim primitive they share. It tracks this PR as "fold into the claim primitive, or close without a consumer", so the answer to "which job depends on this" decides it.

🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.

The first consumer of the lock will be the background session deletion job.

@malatewang
malatewang force-pushed the feature/new-feature branch from 4b9ff9d to 8a18019 Compare October 5, 2026 23:09
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants