Repository navigation
Add SQL leased read-write locks - #1722
Conversation
|
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 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 ( Is reader preference the intended semantics, or should this change? Either way the docstring should say which. reproimport 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. |
|
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 Under that assumption:
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. |
|
Nit: acronym casing in class names The server package capitalizes acronyms in CapWords names, as PEP 8 asks (
🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
Renewal failure during exit cancels a finished body and skips the release
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 Repro (SQLite, renew stubbed to block and then fail so the body exits while it is in flight): scriptimport 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: Fix: in 🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
Replacing the loop's cancellation with When renewal fails, Consequence: structured-concurrency primitives read that counter. Repro on Python 3.14: scriptimport 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: Fix: call 🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
Lock tables grow without bound
Measured on SQLite:
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: 🤖 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. |
|
Two more nits
🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
The shared counter row serializes every acquire in the deployment on PostgreSQL The rework takes fencing tokens from one row in Measured on postgres:16: with a transaction holding that same Fix: on PostgreSQL take the token from a 🤖 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. |
|
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 Measured on SQLite: one holder with a 5-minute lease and three waiters in Fix: back the delay off while waiting, for example doubling up to some fraction of 🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
|
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 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. |
|
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 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. |
|
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 🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
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:
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. |
|
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.
🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
Simplification and minor items
🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
21ab046 to
4c2fa04
Compare
|
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 🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu. |
|
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. |
4b9ff9d to
8a18019
Compare
Purpose of the change
Add SQL-backed leased read/write locks for coordinating work across server instances.
Description
get_sql_lock_service()to the resource manager, using the configured session manager database.Type of change
How Has This Been Tested?
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 checkon the changed server and test modules — passed. PostgreSQL-specific tests are included intest_postgres.pybut were not run locally for this PR preparation.Checklist
Further comments
None.