From 9c1f387bfe6ae54e93a0abc96d454992c8d7b5e9 Mon Sep 17 00:00:00 2001 From: slawwan Date: Tue, 6 Oct 2026 19:36:37 +0500 Subject: [PATCH 1/4] limit lock acquisition attempt time --- CHANGELOG.md | 5 +++++ asyncpg_lock/guard.py | 9 +++++++-- tests/test_guard.py | 23 +++++++++++++++++++++++ 3 files changed, 35 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d482ae2..8f834d4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,8 @@ +## unreleased + +* Limit lock acquisition attempt time + + ## v0.0.3 (2026-10-06) * Use PyPI trusted publishing diff --git a/asyncpg_lock/guard.py b/asyncpg_lock/guard.py index 4e4ad34..5da7b76 100644 --- a/asyncpg_lock/guard.py +++ b/asyncpg_lock/guard.py @@ -102,12 +102,17 @@ async def __ensure_lock_acquired( logger.info("Lock %s might be lost", key) async def __acquire_lock(self, connection: asyncpg.Connection, key: int | tuple[int, int]) -> None: + per_attempt_acquire_budget = self.__after_acquire_delay / 3 while True: try: if isinstance(key, int): - acquired = await connection.fetchval("SELECT pg_try_advisory_lock($1)", key) + acquired = await connection.fetchval( + "SELECT pg_try_advisory_lock($1)", key, timeout=per_attempt_acquire_budget + ) else: - acquired = await connection.fetchval("SELECT pg_try_advisory_lock($1, $2)", key[0], key[1]) + acquired = await connection.fetchval( + "SELECT pg_try_advisory_lock($1, $2)", key[0], key[1], timeout=per_attempt_acquire_budget + ) except Exception: raise Exception(f"Lock {key} not acquired") if acquired: diff --git a/tests/test_guard.py b/tests/test_guard.py index b933778..c339e82 100644 --- a/tests/test_guard.py +++ b/tests/test_guard.py @@ -234,6 +234,29 @@ async def test_reacquire_lock_after_silent_disruption( assert not tracker.overlaps +async def test_acquire_lock_after_silent_disruption_while_waiting( + guard: asyncpg_lock.AdvisoryLockGuard, connector: PgConnector, proxy: TcpProxy, pg_14: pytest_pg.PG +) -> None: + holder = await asyncpg.connect(host=pg_14.host, port=pg_14.port, user=pg_14.user, database=pg_14.database) + await holder.execute("SELECT pg_advisory_lock($1)", LOCK_KEY) + + tracker = ExecutionTracker() + task = asyncio.create_task(guard.run(LOCK_KEY, tracker)) + try: + await asyncio.sleep(LOCK_ACQUIRE_RETRY_INTERVAL) + await proxy.freeze_connections() + await asyncio.sleep(LOCK_ACQUIRE_RETRY_INTERVAL * 2) + await holder.close() + async with asyncio.timeout((LOCK_ACQUIRE_RETRY_INTERVAL + LOCK_ACQUIRE_GRACE_PERIOD + PER_ATTEMPT_DELAY) * 2): + await tracker.min_completed_executions_event.wait() + finally: + await cancel_and_wait(task) + await holder.close() + + assert connector.total_open_connections == 2 + assert not tracker.overlaps + + async def test_no_overlapping_execution_for_same_keys( guard: asyncpg_lock.AdvisoryLockGuard, connector: PgConnector ) -> None: From 658841c907c24b337002be8d3102c61abec6918b Mon Sep 17 00:00:00 2001 From: slawwan Date: Tue, 6 Oct 2026 22:09:20 +0500 Subject: [PATCH 2/4] add acquire_timeout parameter instead of deriving it from after_acquire_delay --- asyncpg_lock/guard.py | 10 +++++++--- tests/test_guard.py | 6 +++++- 2 files changed, 12 insertions(+), 4 deletions(-) diff --git a/asyncpg_lock/guard.py b/asyncpg_lock/guard.py index 5da7b76..9a4a6d2 100644 --- a/asyncpg_lock/guard.py +++ b/asyncpg_lock/guard.py @@ -20,6 +20,7 @@ async def _connect() -> asyncpg.Connection: class AdvisoryLockGuard: __slots__ = ( + "__acquire_timeout", "__after_acquire_delay", "__connect", "__reacquire_delay", @@ -34,6 +35,7 @@ def __init__( reconnect_delay: float = 5, reacquire_delay: float = 5, after_acquire_delay: float = 5, + acquire_timeout: float = 5, ) -> None: if reconnect_delay < 0: raise ValueError("reconnect_delay must be non-negative") @@ -41,10 +43,13 @@ def __init__( raise ValueError("reacquire_delay must be positive") if after_acquire_delay <= 0: raise ValueError("after_acquire_delay must be positive") + if acquire_timeout <= 0: + raise ValueError("acquire_timeout must be positive") self.__reconnect_delay = reconnect_delay self.__after_acquire_delay = after_acquire_delay self.__reacquire_delay = reacquire_delay + self.__acquire_timeout = acquire_timeout self.__connect = connect self.__tasks: set[asyncio.Task[None]] = set() @@ -102,16 +107,15 @@ async def __ensure_lock_acquired( logger.info("Lock %s might be lost", key) async def __acquire_lock(self, connection: asyncpg.Connection, key: int | tuple[int, int]) -> None: - per_attempt_acquire_budget = self.__after_acquire_delay / 3 while True: try: if isinstance(key, int): acquired = await connection.fetchval( - "SELECT pg_try_advisory_lock($1)", key, timeout=per_attempt_acquire_budget + "SELECT pg_try_advisory_lock($1)", key, timeout=self.__acquire_timeout ) else: acquired = await connection.fetchval( - "SELECT pg_try_advisory_lock($1, $2)", key[0], key[1], timeout=per_attempt_acquire_budget + "SELECT pg_try_advisory_lock($1, $2)", key[0], key[1], timeout=self.__acquire_timeout ) except Exception: raise Exception(f"Lock {key} not acquired") diff --git a/tests/test_guard.py b/tests/test_guard.py index c339e82..e98b1ab 100644 --- a/tests/test_guard.py +++ b/tests/test_guard.py @@ -18,6 +18,7 @@ LOCK_ACQUIRE_GRACE_PERIOD = 0.75 LOCK_ACQUIRE_RETRY_INTERVAL = 0.5 PER_ATTEMPT_DELAY = 0.1 +LOCK_ACQUIRE_TIMEOUT = 0.25 LOCK_KEY = random.randint(0, 2**63 - 1) NON_CONFLICTING_LOCK_KEY = random.randint(-(2**63), -1) @@ -153,6 +154,7 @@ def guard(connector: PgConnector) -> asyncpg_lock.AdvisoryLockGuard: reconnect_delay=RECONNECT_DELAY, after_acquire_delay=LOCK_ACQUIRE_GRACE_PERIOD, reacquire_delay=LOCK_ACQUIRE_RETRY_INTERVAL, + acquire_timeout=LOCK_ACQUIRE_TIMEOUT, ) @@ -247,7 +249,9 @@ async def test_acquire_lock_after_silent_disruption_while_waiting( await proxy.freeze_connections() await asyncio.sleep(LOCK_ACQUIRE_RETRY_INTERVAL * 2) await holder.close() - async with asyncio.timeout((LOCK_ACQUIRE_RETRY_INTERVAL + LOCK_ACQUIRE_GRACE_PERIOD + PER_ATTEMPT_DELAY) * 2): + async with asyncio.timeout( + (LOCK_ACQUIRE_RETRY_INTERVAL + LOCK_ACQUIRE_TIMEOUT + LOCK_ACQUIRE_GRACE_PERIOD + PER_ATTEMPT_DELAY) * 2 + ): await tracker.min_completed_executions_event.wait() finally: await cancel_and_wait(task) From 4fcbf00a9fe4e7edaf4a3e36c70de1c91fff035a Mon Sep 17 00:00:00 2001 From: slawwan Date: Wed, 7 Oct 2026 09:16:46 +0500 Subject: [PATCH 3/4] assert guard keeps running after lock acquisition timeout --- tests/test_guard.py | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/test_guard.py b/tests/test_guard.py index e98b1ab..a1c235c 100644 --- a/tests/test_guard.py +++ b/tests/test_guard.py @@ -253,6 +253,7 @@ async def test_acquire_lock_after_silent_disruption_while_waiting( (LOCK_ACQUIRE_RETRY_INTERVAL + LOCK_ACQUIRE_TIMEOUT + LOCK_ACQUIRE_GRACE_PERIOD + PER_ATTEMPT_DELAY) * 2 ): await tracker.min_completed_executions_event.wait() + assert not task.done() finally: await cancel_and_wait(task) await holder.close() From c637783fc6bec68b82dd14e10552738b876271c2 Mon Sep 17 00:00:00 2001 From: slawwan Date: Wed, 7 Oct 2026 09:19:00 +0500 Subject: [PATCH 4/4] log failed lock acquisition attempts --- asyncpg_lock/guard.py | 1 + tests/test_guard.py | 24 ++++++++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/asyncpg_lock/guard.py b/asyncpg_lock/guard.py index 9a4a6d2..0227ff4 100644 --- a/asyncpg_lock/guard.py +++ b/asyncpg_lock/guard.py @@ -118,6 +118,7 @@ async def __acquire_lock(self, connection: asyncpg.Connection, key: int | tuple[ "SELECT pg_try_advisory_lock($1, $2)", key[0], key[1], timeout=self.__acquire_timeout ) except Exception: + logger.warning("Failed to acquire lock %s", key, exc_info=True) raise Exception(f"Lock {key} not acquired") if acquired: return diff --git a/tests/test_guard.py b/tests/test_guard.py index a1c235c..0e7bd2a 100644 --- a/tests/test_guard.py +++ b/tests/test_guard.py @@ -262,6 +262,30 @@ async def test_acquire_lock_after_silent_disruption_while_waiting( assert not tracker.overlaps +async def test_log_failed_lock_acquisition_attempt( + guard: asyncpg_lock.AdvisoryLockGuard, proxy: TcpProxy, pg_14: pytest_pg.PG, caplog: pytest.LogCaptureFixture +) -> None: + holder = await asyncpg.connect(host=pg_14.host, port=pg_14.port, user=pg_14.user, database=pg_14.database) + await holder.execute("SELECT pg_advisory_lock($1)", LOCK_KEY) + + task = asyncio.create_task(guard.run(LOCK_KEY, ExecutionTracker())) + try: + await asyncio.sleep(LOCK_ACQUIRE_RETRY_INTERVAL) + await proxy.freeze_connections() + await asyncio.sleep((LOCK_ACQUIRE_RETRY_INTERVAL + LOCK_ACQUIRE_TIMEOUT) * 2) + finally: + await cancel_and_wait(task) + await holder.close() + + assert any( + x.name == "asyncpg_lock" + and x.levelname == "WARNING" + and x.exc_info is not None + and x.exc_info[0] is TimeoutError + for x in caplog.records + ) + + async def test_no_overlapping_execution_for_same_keys( guard: asyncpg_lock.AdvisoryLockGuard, connector: PgConnector ) -> None: