From 50144889fb456fb00e475c336a1ce08bb3c9ad00 Mon Sep 17 00:00:00 2001 From: T3ST3ST3R0N Date: Wed, 7 Oct 2026 04:12:33 +0330 Subject: [PATCH 1/2] fix(controller): give maintenance calls a timeout longer than node-serviced's update_node, update_core, update_geofiles and hard_reset are synchronous calls to node-serviced, which runs pg-node commands (docker pulls, downloads) with a 5 minute deadline (60 s for hard_reset). They used the node's default_timeout (10 s by default), so the caller got NodeAPIError(-5) while the update kept running, and the lifecycle lease was released early, letting a Start overlap the update. Each maintenance operation now has its own timeout (330 s, 90 s for hard_reset), a little above node-serviced's deadline. --- PasarGuardNodeBridge/controller.py | 13 +++++- tests/test_maintenance.py | 72 ++++++++++++++++++++++++++++++ 2 files changed, 84 insertions(+), 1 deletion(-) create mode 100644 tests/test_maintenance.py diff --git a/PasarGuardNodeBridge/controller.py b/PasarGuardNodeBridge/controller.py index 601045b..8d175ac 100644 --- a/PasarGuardNodeBridge/controller.py +++ b/PasarGuardNodeBridge/controller.py @@ -35,6 +35,15 @@ # Default timeout configuration (module-level constants) DEFAULT_API_TIMEOUT = 10 # Default timeout for public API methods DEFAULT_INTERNAL_TIMEOUT = 15 # Default timeout for internal gRPC/HTTP operations +# Maintenance calls block until node-serviced finishes the pg-node command (docker pulls, downloads). +# node-serviced allows 5 minutes for update/core_update/geofiles and 60 s for hard_reset, so wait a +# bit longer than that: the caller gets the real result and the lifecycle lease is held meanwhile. +MAINTENANCE_TIMEOUTS = { + LifecycleOperation.UPDATE_NODE: 330, + LifecycleOperation.UPDATE_CORE: 330, + LifecycleOperation.UPDATE_GEOFILES: 330, + LifecycleOperation.HARD_RESET: 90, +} class NodeAPIError(Exception): @@ -796,7 +805,9 @@ async def _run_coordinated_update( lease = await self._acquire_lifecycle_lease(operation) try: - return await self._make_json_request(method="POST", endpoint=endpoint, json=json) + return await self._make_json_request( + method="POST", endpoint=endpoint, json=json, timeout=MAINTENANCE_TIMEOUTS[operation] + ) finally: await self._release_lifecycle_lease(lease) diff --git a/tests/test_maintenance.py b/tests/test_maintenance.py new file mode 100644 index 0000000..274b7d5 --- /dev/null +++ b/tests/test_maintenance.py @@ -0,0 +1,72 @@ +"""Maintenance calls go to node-serviced, which runs pg-node commands that take minutes.""" + +import asyncio +import logging +import unittest +from unittest.mock import patch + +from aiohttp import web + +from PasarGuardNodeBridge.aiohttp_compat import LazyClientSession, make_timeout +from PasarGuardNodeBridge.rest import Node as RestNode +from PasarGuardNodeBridge.storage import InMemoryUserSyncStore + +SLOW_COMMAND_SECONDS = 1.5 + + +class MaintenanceTimeoutTests(unittest.IsolatedAsyncioTestCase): + async def asyncSetUp(self): + async def root(request): + return web.json_response({"status": "ok"}) + + async def slow_command(request): + await asyncio.sleep(SLOW_COMMAND_SECONDS) + return web.json_response({"status": "ok"}) + + app = web.Application() + app.router.add_get("/", root) + for path in ("/node/update", "/node/core_update", "/node/geofiles", "/node/hard_reset"): + app.router.add_post(path, slow_command) + self.runner = web.AppRunner(app) + await self.runner.setup() + site = web.TCPSite(self.runner, "127.0.0.1", 0) + await site.start() + port = site._server.sockets[0].getsockname()[1] + + with patch("ssl.SSLContext.load_verify_locations"): + self.node = RestNode( + address="127.0.0.1", + port=1, + api_port=port, + server_ca="test", + api_key="00000000-0000-0000-0000-000000000000", + node_id="maintenance", + user_sync_store=InMemoryUserSyncStore(), + default_timeout=1, + logger=logging.getLogger("maintenance-test"), + ) + # Talk plain HTTP to the local stand-in for node-serviced. + await self.node._json_client.close() + self.node._json_client = LazyClientSession( + ssl_context=None, headers={}, base_url=f"http://127.0.0.1:{port}", timeout=make_timeout(30) + ) + + async def asyncTearDown(self): + await self.node._json_client.close() + await self.runner.cleanup() + + async def test_maintenance_calls_outlive_default_timeout(self): + calls = { + "update_node": lambda: self.node.update_node(), + "update_core": lambda: self.node.update_core({"core_version": "v25.8.31"}), + "update_geofiles": lambda: self.node.update_geofiles({"region": "iran"}), + "hard_reset": lambda: self.node.hard_reset(), + } + for name, call in calls.items(): + with self.subTest(name): + response = await call() + self.assertEqual(response.status_code, 200) + + +if __name__ == "__main__": + unittest.main() From 65b54b9b08d89ca79d842c59e66749ef624bf8a2 Mon Sep 17 00:00:00 2001 From: T3ST3ST3R0N Date: Wed, 7 Oct 2026 05:36:20 +0330 Subject: [PATCH 2/2] fix(controller): keep renewing the lifecycle lease through coordinator errors The lease heartbeat only caught CancelledError. Any error from the coordinator (e.g. a NATS KV hiccup) killed the heartbeat task silently, the 60 s lease expired while a maintenance call could still be running for minutes, and releasing the lease then re-raised that error, so the caller saw a failure even if the operation itself succeeded. The heartbeat now logs coordinator errors and keeps renewing, and stops only when a coordinator reports the lease as lost (returns False). Tests: the heartbeat survives a failing coordinator during a long maintenance call, and the per-operation maintenance timeout is still enforced (NodeAPIError -5 when exceeded). --- PasarGuardNodeBridge/controller.py | 13 +++++- tests/test_maintenance.py | 75 +++++++++++++++++++++++------- 2 files changed, 69 insertions(+), 19 deletions(-) diff --git a/PasarGuardNodeBridge/controller.py b/PasarGuardNodeBridge/controller.py index 8d175ac..6e84bf0 100644 --- a/PasarGuardNodeBridge/controller.py +++ b/PasarGuardNodeBridge/controller.py @@ -353,7 +353,18 @@ async def _heartbeat_lifecycle_lease(self, lease: LifecycleLease) -> None: try: while True: await asyncio.sleep(interval) - await self._lifecycle_coordinator.heartbeat(lease) + try: + renewed = await self._lifecycle_coordinator.heartbeat(lease) + except Exception as e: + # Keep renewing: a transient store error must not let the lease lapse while a long + # operation (e.g. a node update that runs for minutes) is still in progress. + self.logger.warning( + f"[{self.name}] Lifecycle lease heartbeat failed | Error: {type(e).__name__} - {e!s}" + ) + continue + if renewed is False: # coordinators that report it: the lease is gone, stop renewing + self.logger.warning(f"[{self.name}] Lifecycle lease for {lease.operation} was lost") + return except asyncio.CancelledError: pass diff --git a/tests/test_maintenance.py b/tests/test_maintenance.py index 274b7d5..6aae6d8 100644 --- a/tests/test_maintenance.py +++ b/tests/test_maintenance.py @@ -7,13 +7,29 @@ from aiohttp import web +from PasarGuardNodeBridge import controller from PasarGuardNodeBridge.aiohttp_compat import LazyClientSession, make_timeout +from PasarGuardNodeBridge.controller import NodeAPIError from PasarGuardNodeBridge.rest import Node as RestNode -from PasarGuardNodeBridge.storage import InMemoryUserSyncStore +from PasarGuardNodeBridge.storage import InMemoryNodeLifecycleCoordinator, InMemoryUserSyncStore, LifecycleOperation SLOW_COMMAND_SECONDS = 1.5 +class FlakyHeartbeatCoordinator(InMemoryNodeLifecycleCoordinator): + """Fails the first lease heartbeat, then renews normally.""" + + def __init__(self): + super().__init__() + self.heartbeats = 0 + + async def heartbeat(self, lease): + self.heartbeats += 1 + if self.heartbeats == 1: + raise RuntimeError("store temporarily unavailable") + return await super().heartbeat(lease) + + class MaintenanceTimeoutTests(unittest.IsolatedAsyncioTestCase): async def asyncSetUp(self): async def root(request): @@ -29,44 +45,67 @@ async def slow_command(request): app.router.add_post(path, slow_command) self.runner = web.AppRunner(app) await self.runner.setup() - site = web.TCPSite(self.runner, "127.0.0.1", 0) - await site.start() - port = site._server.sockets[0].getsockname()[1] + await web.TCPSite(self.runner, "127.0.0.1", 0).start() + self.api_port = self.runner.addresses[0][1] + self.nodes = [] + + async def asyncTearDown(self): + for node in self.nodes: + await node._json_client.close() + await self.runner.cleanup() + def make_node(self, **kwargs): with patch("ssl.SSLContext.load_verify_locations"): - self.node = RestNode( + node = RestNode( address="127.0.0.1", port=1, - api_port=port, + api_port=self.api_port, server_ca="test", api_key="00000000-0000-0000-0000-000000000000", - node_id="maintenance", + node_id=f"maintenance-{len(self.nodes)}", user_sync_store=InMemoryUserSyncStore(), default_timeout=1, logger=logging.getLogger("maintenance-test"), + **kwargs, ) # Talk plain HTTP to the local stand-in for node-serviced. - await self.node._json_client.close() - self.node._json_client = LazyClientSession( - ssl_context=None, headers={}, base_url=f"http://127.0.0.1:{port}", timeout=make_timeout(30) + node._json_client = LazyClientSession( + ssl_context=None, headers={}, base_url=f"http://127.0.0.1:{self.api_port}", timeout=make_timeout(30) ) - - async def asyncTearDown(self): - await self.node._json_client.close() - await self.runner.cleanup() + self.nodes.append(node) + return node async def test_maintenance_calls_outlive_default_timeout(self): + node = self.make_node() calls = { - "update_node": lambda: self.node.update_node(), - "update_core": lambda: self.node.update_core({"core_version": "v25.8.31"}), - "update_geofiles": lambda: self.node.update_geofiles({"region": "iran"}), - "hard_reset": lambda: self.node.hard_reset(), + "update_node": lambda: node.update_node(), + "update_core": lambda: node.update_core({"core_version": "v25.8.31"}), + "update_geofiles": lambda: node.update_geofiles({"region": "iran"}), + "hard_reset": lambda: node.hard_reset(), } for name, call in calls.items(): with self.subTest(name): response = await call() self.assertEqual(response.status_code, 200) + async def test_maintenance_timeout_is_still_enforced(self): + node = self.make_node() + with ( + patch.dict(controller.MAINTENANCE_TIMEOUTS, {LifecycleOperation.UPDATE_NODE: 1}), + self.assertRaises(NodeAPIError) as raised, + ): + await node.update_node() + self.assertEqual(raised.exception.code, -5) + + async def test_lease_heartbeat_survives_a_coordinator_error(self): + coordinator = FlakyHeartbeatCoordinator() + node = self.make_node(lifecycle_coordinator=coordinator, lifecycle_lease_seconds=0.3) + + response = await node.update_node() + + self.assertEqual(response.status_code, 200) + self.assertGreaterEqual(coordinator.heartbeats, 3) + if __name__ == "__main__": unittest.main()